diff --git a/src/gateway/server-runtime-subscriptions.test.ts b/src/gateway/server-runtime-subscriptions.test.ts index ae0e3a4dd604..a899e21c6684 100644 --- a/src/gateway/server-runtime-subscriptions.test.ts +++ b/src/gateway/server-runtime-subscriptions.test.ts @@ -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(); + return { + ...actual, + readSessionMessageCountAsync: transcriptBroadcastMocks.readMessageCount, + }; +}); + +vi.mock("./server-session-events.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + createTranscriptUpdateBroadcastHandler: ( + ...args: Parameters + ) => { + 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[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()); diff --git a/src/gateway/server-session-events.test.ts b/src/gateway/server-session-events.test.ts index 8fe907b4f65e..9f11e3b68390 100644 --- a/src/gateway/server-session-events.test.ts +++ b/src/gateway/server-session-events.test.ts @@ -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", diff --git a/src/gateway/server-session-events.ts b/src/gateway/server-session-events.ts index 1aba29d36156..ec40745b07d7 100644 --- a/src/gateway/server-session-events.ts +++ b/src/gateway/server-session-events.ts @@ -164,7 +164,7 @@ export function createTranscriptUpdateBroadcastHandler(params: { chatAbortControllers: Map; }) { let broadcastQueue = Promise.resolve(); - return (update: InternalSessionTranscriptUpdate): void => { + return (update: InternalSessionTranscriptUpdate): Promise => { // 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; }; }