fix(gateway): report queued transcript broadcast failures

This commit is contained in:
Peter Steinberger
2026-08-01 10:09:17 -07:00
parent 3346507d84
commit ce47ea331f
3 changed files with 110 additions and 34 deletions

View File

@@ -54,6 +54,10 @@ const auditTestState = vi.hoisted(() => ({
const agentEventHandlerMocks = vi.hoisted(() => ({
create: vi.fn(),
}));
const transcriptBroadcastMocks = vi.hoisted(() => ({
useActualHandler: false,
readMessageCount: vi.fn(),
}));
vi.mock("../audit/audit-config.js", () => ({
isAuditLedgerEnabled: () => auditTestState.enabled,
@@ -84,14 +88,33 @@ vi.mock("./server-session-key.js", () => ({
resolveSessionKeyForRun: () => "agent:main:main",
}));
vi.mock("./server-session-events.js", () => ({
createTranscriptUpdateBroadcastHandler: () => () => {
throw new Error("transcript handler failure");
},
createLifecycleEventBroadcastHandler: () => () => {
throw new Error("lifecycle handler failure");
},
}));
vi.mock("./session-transcript-readers.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("./session-transcript-readers.js")>();
return {
...actual,
readSessionMessageCountAsync: transcriptBroadcastMocks.readMessageCount,
};
});
vi.mock("./server-session-events.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("./server-session-events.js")>();
return {
...actual,
createTranscriptUpdateBroadcastHandler: (
...args: Parameters<typeof actual.createTranscriptUpdateBroadcastHandler>
) => {
if (transcriptBroadcastMocks.useActualHandler) {
return actual.createTranscriptUpdateBroadcastHandler(...args);
}
return () => {
throw new Error("transcript handler failure");
};
},
createLifecycleEventBroadcastHandler: () => () => {
throw new Error("lifecycle handler failure");
},
};
});
const { startGatewayEventSubscriptions } = await import("./server-runtime-subscriptions.js");
type SubscriptionParams = Parameters<typeof startGatewayEventSubscriptions>[0];
@@ -122,6 +145,8 @@ describe("startGatewayEventSubscriptions", () => {
auditTestState.created = 0;
auditTestState.recorded = 0;
auditTestState.stopped = 0;
transcriptBroadcastMocks.useActualHandler = false;
transcriptBroadcastMocks.readMessageCount.mockReset();
agentEventHandlerMocks.create.mockReset().mockImplementation(() => {
throw new Error("server-chat lazy load failure");
});
@@ -219,6 +244,58 @@ describe("startGatewayEventSubscriptions", () => {
);
});
it("logs real asynchronous transcript failures and recovers the broadcast queue", async () => {
transcriptBroadcastMocks.useActualHandler = true;
const persistenceFailure = new Error("session transcript read failed");
transcriptBroadcastMocks.readMessageCount
.mockRejectedValueOnce(persistenceFailure)
.mockResolvedValueOnce(2);
const params = createParams();
params.sessionEventSubscribers.subscribe("conn-transcript");
unsubs = startGatewayEventSubscriptions(params);
const emitMessage = (messageId: string) =>
emitSessionTranscriptUpdate({
sessionFile: "/tmp/openclaw-transcript-dispatch.sqlite",
sessionKey: "agent:main:main",
message: { role: "assistant", content: [{ type: "text", text: "visible answer" }] },
messageId,
target: {
agentId: "main",
sessionId: "sess-transcript",
sessionKey: "agent:main:main",
storePath: "/tmp/openclaw-transcript-dispatch-sessions.json",
},
});
emitMessage("failed-message");
await waitForFast(() =>
expect(transcriptBroadcastMocks.readMessageCount).toHaveBeenCalledOnce(),
);
await waitForFast(() =>
expect(warn).toHaveBeenCalledWith("Transcript update dispatch failed", {
sessionKey: "agent:main:main",
error: persistenceFailure,
}),
);
expect(params.broadcastToConnIds).not.toHaveBeenCalled();
emitMessage("recovered-message");
await waitForFast(() => expect(params.broadcastToConnIds).toHaveBeenCalledOnce());
expect(params.broadcastToConnIds).toHaveBeenCalledWith(
"session.message",
expect.objectContaining({
sessionKey: "agent:main:main",
messageId: "recovered-message",
messageSeq: 2,
}),
new Set(["conn-transcript"]),
);
expect(transcriptBroadcastMocks.readMessageCount).toHaveBeenCalledTimes(2);
expect(warn).toHaveBeenCalledOnce();
});
it("logs lifecycle handler failures", async () => {
unsubs = startGatewayEventSubscriptions(createParams());

View File

@@ -85,14 +85,14 @@ async function emitAssistantTranscriptUpdate(
message: unknown = { role: "assistant", content: [{ type: "text", text: "Final answer" }] },
) {
const { broadcastToConnIds, handler } = createHandler(projectSessionActive);
handler({
await handler({
sessionFile: "/tmp/sess-main.jsonl",
sessionKey: "agent:main:main",
message,
messageId: "message-1",
messageSeq: 1,
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledTimes(1));
expect(broadcastToConnIds).toHaveBeenCalledTimes(1);
return broadcastToConnIds.mock.calls[0]?.[1];
}
@@ -109,7 +109,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
it("never silently drops an authoritative session message for a slow subscriber", async () => {
const { broadcastToConnIds, handler } = createHandler(false);
handler({
await handler({
sessionFile: "/tmp/sess-main.jsonl",
sessionKey: "agent:main:main",
message: { role: "user", content: [{ type: "text", text: "shared durable prompt" }] },
@@ -117,7 +117,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
messageSeq: 1,
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledTimes(1));
expect(broadcastToConnIds).toHaveBeenCalledTimes(1);
expect(broadcastToConnIds).toHaveBeenCalledWith(
"session.message",
expect.objectContaining({
@@ -148,7 +148,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
chatAbortControllers: new Map(),
});
handler({
await handler({
target: {
agentId: "main",
sessionId: "sess-main",
@@ -158,7 +158,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
lifecycleRevision: "committed-revision",
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledOnce());
expect(broadcastToConnIds).toHaveBeenCalledOnce();
expect(getSessionMessageSubscribers).toHaveBeenCalledWith("agent:main:main");
expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledWith({
agentId: "main",
@@ -188,7 +188,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
it("rejects an identity-only invalidation when its custom-store owner was deleted", async () => {
const { broadcastToConnIds, handler } = createHandler(false);
handler({
await handler({
target: {
agentId: "main",
sessionId: "sess-main",
@@ -198,7 +198,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
lifecycleRevision: "deleted-revision",
});
await vi.waitFor(() => expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledOnce());
expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledOnce();
expect(broadcastToConnIds).not.toHaveBeenCalled();
expect(readSessionMessageCountAsyncMock).not.toHaveBeenCalled();
});
@@ -268,7 +268,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
);
const { broadcastToConnIds, handler } = createHandler(false);
handler({
const pendingBroadcast = handler({
target: {
agentId: "main",
sessionId: "sess-main",
@@ -290,9 +290,8 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
}
resolveMessageCount?.(1);
await vi.waitFor(() =>
expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledTimes(expectedEntryReads),
);
await pendingBroadcast;
expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledTimes(expectedEntryReads);
expect(loadGatewaySessionEntryReadOnlyMock).not.toHaveBeenCalled();
expect(broadcastToConnIds).not.toHaveBeenCalled();
},
@@ -314,7 +313,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
});
const { broadcastToConnIds, handler } = createHandler(false);
handler({
await handler({
target: {
agentId: "main",
sessionId: "sess-main",
@@ -327,7 +326,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
messageSeq: 1,
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledOnce());
expect(broadcastToConnIds).toHaveBeenCalledOnce();
expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledWith({
agentId: "main",
sessionKey: "agent:main:main",
@@ -356,7 +355,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
});
const { broadcastToConnIds, handler } = createHandler(false);
handler({
await handler({
target: {
agentId: "main",
sessionId: "sess-main",
@@ -369,7 +368,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
messageSeq: 1,
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledOnce());
expect(broadcastToConnIds).toHaveBeenCalledOnce();
expect(broadcastToConnIds).toHaveBeenCalledWith(
"session.message",
expect.objectContaining({
@@ -385,7 +384,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
readSessionMessageCountAsyncMock.mockResolvedValue(3);
const { broadcastToConnIds, handler } = createHandler(false);
handler({
await handler({
target: {
agentId: "main",
sessionId: "sess-main",
@@ -396,7 +395,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
messageId: "legacy-lifecycle-message",
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledOnce());
expect(broadcastToConnIds).toHaveBeenCalledOnce();
expect(broadcastToConnIds).toHaveBeenCalledWith(
"session.message",
expect.objectContaining({
@@ -505,14 +504,14 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
projectChatDisplayMessageMock.mockReturnValueOnce(undefined).mockReturnValueOnce(undefined);
try {
handler({
await handler({
sessionFile: "/tmp/sess-main.jsonl",
sessionKey: "agent:main:main",
message: { role: "toolResult", content: [] },
messageId: "message-1",
messageSeq: 1,
});
await vi.waitFor(() => expect(received).toHaveBeenCalledOnce());
expect(received).toHaveBeenCalledOnce();
expect(received).toHaveBeenCalledWith({
sessionKey: "agent:main:main",
phase: "message",
@@ -530,7 +529,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
readSessionMessageCountAsyncMock.mockResolvedValue(7);
const { broadcastToConnIds, handler } = createHandler(false);
handler({
await handler({
agentId: "main",
message: { role: "assistant", content: [{ type: "text", text: "Final answer" }] },
messageId: "message-partial-target",
@@ -542,7 +541,7 @@ describe("createTranscriptUpdateBroadcastHandler", () => {
storePath: "/tmp/explicit-sessions.json",
},
});
await vi.waitFor(() => expect(broadcastToConnIds).toHaveBeenCalledTimes(1));
expect(broadcastToConnIds).toHaveBeenCalledTimes(1);
expect(loadAccessorSessionEntryReadOnlyMock).toHaveBeenCalledWith({
agentId: "main",

View File

@@ -164,7 +164,7 @@ export function createTranscriptUpdateBroadcastHandler(params: {
chatAbortControllers: Map<string, ChatAbortControllerEntry>;
}) {
let broadcastQueue = Promise.resolve();
return (update: InternalSessionTranscriptUpdate): void => {
return (update: InternalSessionTranscriptUpdate): Promise<void> => {
// Capture legacy ownership before the async queue can cross a same-id reset;
// committed producer ownership always wins over a later session-store read.
const lifecycleRevision =
@@ -175,9 +175,9 @@ export function createTranscriptUpdateBroadcastHandler(params: {
const queuedUpdate = lifecycleRevision ? { ...update, lifecycleRevision } : update;
// Preserve transcript update order even when counting messages requires an
// async read from the session file.
broadcastQueue = broadcastQueue
.then(() => handleTranscriptUpdateBroadcast(params, queuedUpdate))
.catch(() => undefined);
const task = broadcastQueue.then(() => handleTranscriptUpdateBroadcast(params, queuedUpdate));
broadcastQueue = task.catch(() => undefined);
return task;
};
}