diff --git a/.changeset/chat-stop-successor-boundary.md b/.changeset/chat-stop-successor-boundary.md new file mode 100644 index 00000000000..f356ac582b0 --- /dev/null +++ b/.changeset/chat-stop-successor-boundary.md @@ -0,0 +1,8 @@ +--- +"@trigger.dev/sdk": patch +--- + +Keep new chat responses intact after Stop, including slow Stop acknowledgments and page reloads. +Sequence-free replies after Stop require a transcript reload before further messages. +Loading a fresh transcript through `useLoadTranscript` restores blocked sessions only after its saved input cursor covers the stopped turn. +Transcript recovery reports missing cursor evidence and empty output polls. An empty recovery poll keeps the accepted message available for reconnect. diff --git a/knip.json b/knip.json index fb75544f4d9..b0c5339fe77 100644 --- a/knip.json +++ b/knip.json @@ -88,7 +88,7 @@ }, "packages/trigger-sdk": { "ignoreFiles": ["src/**/*-cjs.cts", "src/v3/index-browser.mts"], - "ignoreDependencies": ["ai-v7", "react"] + "ignoreDependencies": ["ai-v7"] }, "docs": { "ignoreFiles": ["style.css"] diff --git a/packages/trigger-sdk/package.json b/packages/trigger-sdk/package.json index 4bdb5a0c080..c3d2355c978 100644 --- a/packages/trigger-sdk/package.json +++ b/packages/trigger-sdk/package.json @@ -87,8 +87,12 @@ "@ai-sdk/provider": "3.0.8", "@arethetypeswrong/cli": "^0.18.5", "@types/react": "^19.2.14", + "@types/react-dom": "19.2.3", "ai": "^6.0.116", "ai-v7": "npm:ai@7.0.0-canary.159", + "jsdom": "30.0.1", + "react": "18.3.1", + "react-dom": "18.3.1", "rimraf": "^6.0.1", "tshy": "^4.1.3", "tsx": "4.17.0", diff --git a/packages/trigger-sdk/src/v3/chat-react.ts b/packages/trigger-sdk/src/v3/chat-react.ts index f9fdaa47a19..a40d473c2fd 100644 --- a/packages/trigger-sdk/src/v3/chat-react.ts +++ b/packages/trigger-sdk/src/v3/chat-react.ts @@ -32,6 +32,7 @@ import { type InferChatUIMessage, } from "./ai-shared.js"; import type { UIMessage, ChatRequestOptions } from "ai"; +import type { TranscriptCursors } from "./transcriptStorage.js"; /** * Options for `useTriggerChatTransport`, with a type-safe `task` field. @@ -55,7 +56,7 @@ export type { ChatTransportEvent, ChatTransportSendSource } from "./chat.js"; /** What a `chat.createLoadTranscriptAction` action returns, as `useLoadTranscript` reads it. */ export type LoadTranscriptResult = { messages: TUIMessage[]; - cursors?: { lastOutEventId?: string; lastInEventId?: string }; + cursors?: TranscriptCursors; nextCursor?: string; }; @@ -65,6 +66,8 @@ export type UseLoadTranscriptOptions = { * the loaded transcript, so the live subscription opens just past the * persisted history instead of replaying it. Only applies once the * transport knows the session (from `sessions` or after `start`). + * A blocked session requires a newer output cursor and an input cursor that covers the stopped input. + * A stale result leaves sends blocked. */ transport?: TriggerChatTransport; /** Page size passed to the action. */ @@ -77,12 +80,13 @@ export type UseLoadTranscriptOptions = { * history. Applied to the session now if it exists, otherwise held by the * transport until the session is created, so a load that resolves before the * session exists still moves the cursor. A no-op when the transcript carries - * no cursor. Returns whether a cursor was provided. + * no cursor. + * Returns whether the cursor was accepted. */ export function seedTranscriptCursor( transport: Pick, chatId: string, - cursors: { lastOutEventId?: string } | undefined + cursors: TranscriptCursors | undefined ): boolean { const lastEventId = cursors?.lastOutEventId; if (!lastEventId) return false; @@ -130,8 +134,7 @@ export function useLoadTranscript( const loadRef = useRef(load); loadRef.current = load; - const transportRef = useRef(options?.transport); - transportRef.current = options?.transport; + const transport = options?.transport; const limit = options?.limit; useEffect(() => { @@ -140,13 +143,18 @@ export function useLoadTranscript( return; } let cancelled = false; + const completeRecovery = transport?.prepareTranscriptRecovery(chatId); setState({ chatId, messages: [], isLoading: true, error: undefined, nextCursor: undefined }); loadRef .current({ chatId, ...(limit !== undefined ? { limit } : {}) }) .then((result) => { if (cancelled) return; - if (transportRef.current) { - seedTranscriptCursor(transportRef.current, chatId, result.cursors); + if (completeRecovery) { + if (!completeRecovery(result.cursors)) { + throw new Error("The loaded transcript is not current. Reload the chat again."); + } + } else if (transport) { + seedTranscriptCursor(transport, chatId, result.cursors); } setState({ chatId, @@ -164,7 +172,7 @@ export function useLoadTranscript( return () => { cancelled = true; }; - }, [chatId, limit]); + }, [chatId, limit, transport]); return { messages: state.chatId === chatId ? state.messages : [], diff --git a/packages/trigger-sdk/src/v3/chat-stop.test.ts b/packages/trigger-sdk/src/v3/chat-stop.test.ts new file mode 100644 index 00000000000..53d6dcee8c8 --- /dev/null +++ b/packages/trigger-sdk/src/v3/chat-stop.test.ts @@ -0,0 +1,931 @@ +import { createServer, type Server, type ServerResponse } from "node:http"; +import { readUIMessageStream, type UIMessageChunk } from "ai"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { + TriggerChatTransport, + type ChatSessionPersistedState, + type ChatTransportEvent, + type TriggerChatTransportOptions, +} from "./chat.js"; + +type OutputRecord = { + seq_num: number; + timestamp: number; + body: string; + headers: string[][]; +}; + +function chunk(seq: number, data: UIMessageChunk): OutputRecord { + return { + seq_num: seq, + timestamp: seq, + body: JSON.stringify({ id: `part-${seq}`, data }), + headers: [], + }; +} + +function complete(seq: number, input: number): OutputRecord { + return { + seq_num: seq, + timestamp: seq, + body: "", + headers: [ + ["trigger-control", "turn-complete"], + ["session-in-event-id", String(input)], + ], + }; +} + +function reply(start: number): OutputRecord[] { + return [ + chunk(start, { type: "start", messageId: "new" }), + chunk(start + 1, { type: "text-start", id: "text" }), + chunk(start + 2, { type: "text-delta", id: "text", delta: "New response" }), + chunk(start + 3, { type: "text-end", id: "text" }), + chunk(start + 4, { type: "finish" }), + ]; +} + +async function readText(stream: ReadableStream): Promise { + let text = ""; + for await (const message of readUIMessageStream({ stream, terminateOnError: true })) { + text = message.parts + .filter((part) => part.type === "text") + .map((part) => part.text) + .join(""); + } + return text; +} + +function readWatchedTurn(stream: ReadableStream): Promise { + return readText( + stream.pipeThrough( + new TransformStream({ + transform(value, controller) { + controller.enqueue(value); + if (value.type === "finish") controller.terminate(); + }, + }) + ) + ); +} + +describe("Stop with a successor response", () => { + let server: Server; + let baseURL: string; + let transport: TriggerChatTransport; + let outputs: ServerResponse[]; + let outputHeaders: { peek: boolean; timeout: number }[]; + let inputSeq: number; + let holdStop: boolean; + let stopStatus: number; + let settled: boolean; + let resumeAfterStoppedCheckpoint: boolean; + let emptyRecoveredOutput: boolean; + let includeSequence: boolean; + let pendingStop: { response: ServerResponse; seq: number } | undefined; + let saved: ChatSessionPersistedState | null; + + function createTransport( + session: ChatSessionPersistedState, + options: Partial = {} + ) { + return new TriggerChatTransport({ + task: "test-chat", + baseURL, + accessToken: () => "test-token", + sessions: { chat: session }, + onSessionChange: (_chatId, session) => { + saved = session; + }, + ...options, + }); + } + + function appendResponse(response: ServerResponse, seq: number, status = 200) { + response + .writeHead(status, { "Content-Type": "application/json" }) + .end(JSON.stringify(includeSequence ? { seq } : {})); + } + + beforeEach(async () => { + outputs = []; + outputHeaders = []; + inputSeq = 10; + holdStop = false; + stopStatus = 200; + settled = false; + resumeAfterStoppedCheckpoint = false; + emptyRecoveredOutput = false; + includeSequence = true; + pendingStop = undefined; + saved = null; + server = createServer(async (request, response) => { + if (request.method === "POST") { + let body = ""; + for await (const data of request) body += data; + const input: unknown = JSON.parse(body); + const isStop = + typeof input === "object" && input !== null && "kind" in input && input.kind === "stop"; + const seq = inputSeq++; + if (isStop && holdStop) { + pendingStop = { response, seq }; + } else { + appendResponse(response, seq, isStop ? stopStatus : 200); + } + return; + } + const stoppedCheckpointPeek = + resumeAfterStoppedCheckpoint && request.headers["x-peek-settled"] !== undefined; + outputHeaders.push({ + peek: request.headers["x-peek-settled"] !== undefined, + timeout: Number(request.headers["timeout-seconds"]), + }); + response.writeHead(200, { + "Content-Type": "text/event-stream", + "X-Stream-Version": "v2", + "X-Session-Settled": String(settled || stoppedCheckpointPeek), + }); + response.flushHeaders(); + outputs.push(response); + if (stoppedCheckpointPeek || emptyRecoveredOutput) { + response.end(); + } else if (resumeAfterStoppedCheckpoint) { + response.write( + `event: batch\ndata: ${JSON.stringify({ records: [...reply(12), complete(17, 12)] })}\n\n` + ); + } + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const address = server.address(); + if (!address || typeof address === "string") throw new Error("Expected a TCP address"); + baseURL = `http://127.0.0.1:${address.port}`; + transport = createTransport({ publicAccessToken: "test-token" }); + }); + + afterEach(async () => { + transport.dispose(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + }); + + async function send(abortSignal?: AbortSignal) { + const before = outputs.length; + const stream = await transport.sendMessages({ + chatId: "chat", + trigger: "submit-message", + messageId: "user", + messages: [{ id: "user", role: "user", parts: [{ type: "text", text: "Continue" }] }], + abortSignal, + }); + await vi.waitFor(() => expect(outputs.length).toBeGreaterThan(before)); + return stream; + } + + function emit(records: OutputRecord[]) { + const response = outputs.at(-1); + if (!response || response.destroyed) throw new Error("The output subscription is closed"); + response.write(`event: batch\ndata: ${JSON.stringify({ records })}\n\n`); + } + + function oldTailAndReply(oldInput = 11, newInput = 12): OutputRecord[] { + return [ + chunk(4, { type: "tool-output-available", toolCallId: "old-tool", output: "Late output" }), + complete(5, oldInput), + ...reply(6), + complete(11, newInput), + ]; + } + + async function hydrateBlockedSession(hydrate: "constructor" | "setSession") { + const first = await send(); + const reader = first.getReader(); + emit([chunk(1, { type: "start", messageId: "old" })]); + await reader.read(); + await transport.stopGeneration("chat"); + includeSequence = false; + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + const session = transport.getSession("chat"); + if (!session) throw new Error("Expected persisted state"); + expect(session).toMatchObject({ requiresTranscriptReload: true, lastEventId: "1" }); + transport.dispose(); + transport = createTransport( + hydrate === "constructor" ? session : { publicAccessToken: "test-token" } + ); + if (hydrate === "setSession") transport.setSession("chat", session); + includeSequence = true; + } + + it.each([false, true])( + "keeps a successor before the Stop acknowledgment (resumed: %s)", + async (resumed) => { + const abort = new AbortController(); + let first: ReadableStream; + if (resumed) { + transport.setSession("chat", { publicAccessToken: "test-token", lastEventId: "1" }); + inputSeq = 11; + const stream = await transport.reconnectToStream({ + chatId: "chat", + abortSignal: abort.signal, + }); + if (!stream) throw new Error("Expected a resumed stream"); + first = stream; + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + } else { + first = await send(abort.signal); + } + const reader = first.getReader(); + emit([ + chunk(2, { type: "start", messageId: "old" }), + chunk(3, { + type: "tool-input-available", + toolCallId: "old-tool", + toolName: "bash", + input: {}, + }), + ]); + await reader.read(); + await reader.read(); + // Resumed streams do not send Stop on abort. This matches useChat.stop(). + if (resumed) { + abort.abort(); + await reader.read(); + } + holdStop = true; + const stopped = transport.stopGeneration("chat"); + await vi.waitFor(() => expect(pendingStop).toBeDefined()); + const next = await send(); + const pending = pendingStop!; + appendResponse(pending.response, pending.seq); + expect(await stopped).toBe(true); + expect(transport.getSession("chat")?.isStreaming).toBe(true); + emit(oldTailAndReply()); + await expect(readText(next)).resolves.toBe("New response"); + } + ); + + it.each([false, true])( + "discards stopped output before the first resumed record (abort first: %s)", + async (abortFirst) => { + transport.setSession("chat", { publicAccessToken: "test-token", lastEventId: "1" }); + const abort = new AbortController(); + const resumed = await transport.reconnectToStream({ + chatId: "chat", + abortSignal: abort.signal, + }); + if (!resumed) throw new Error("Expected a resumed stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + const reader = resumed.getReader(); + if (abortFirst) { + abort.abort(); + expect((await reader.read()).done).toBe(true); + } + await transport.stopGeneration("chat"); + const next = await send(); + emit(oldTailAndReply(10, 11)); + await expect(readText(next)).resolves.toBe("New response"); + expect(transport.getSession("chat")?.skipToTurnComplete).toBe(false); + } + ); + + it("does not gate a response after an empty settled resume", async () => { + settled = true; + transport.setSession("chat", { publicAccessToken: "test-token", lastEventId: "1" }); + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + outputs[0]!.end(); + await expect(readText(resumed)).resolves.toBe(""); + await transport.stopGeneration("chat"); + settled = false; + const next = await send(); + emit([...reply(2), complete(7, 11)]); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("does not gate a response after Stop on a known idle watch", async () => { + transport.dispose(); + transport = createTransport( + { publicAccessToken: "test-token", lastEventId: "1", isStreaming: false }, + { watch: true } + ); + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a watch stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + await transport.stopGeneration("chat"); + const next = await send(); + emit([...reply(2), complete(7, 11)]); + await expect(readWatchedTurn(next)).resolves.toBe("New response"); + }); + + it("does not gate a response after a passive watch abort with unknown turn state", async () => { + transport.dispose(); + transport = createTransport( + { publicAccessToken: "test-token", lastEventId: "1" }, + { watch: true } + ); + const abort = new AbortController(); + const resumed = await transport.reconnectToStream({ + chatId: "chat", + abortSignal: abort.signal, + }); + if (!resumed) throw new Error("Expected a watch stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + abort.abort(); + await expect(readText(resumed)).resolves.toBe(""); + const next = await send(); + emit([...reply(2), complete(7, 10)]); + await expect(readWatchedTurn(next)).resolves.toBe("New response"); + expect(inputSeq).toBe(11); + }); + + it.each(["constructor", "setSession"] as const)( + "retains the stopped boundary through %s hydration", + async (hydrate) => { + await send(); + await transport.stopGeneration("chat"); + expect(saved).toMatchObject({ + skipToTurnComplete: true, + supersededInputSeq: 10, + isStreaming: false, + }); + const session = transport.getSession("chat"); + if (!session) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport( + hydrate === "constructor" ? session : { publicAccessToken: "test-token" } + ); + if (hydrate === "setSession") transport.setSession("chat", session); + const next = await send(); + emit(oldTailAndReply()); + await expect(readText(next)).resolves.toBe("New response"); + expect(saved).toMatchObject({ skipToTurnComplete: false, supersededInputSeq: undefined }); + } + ); + + it("retains unread output after a failed Stop request", async () => { + await send(); + stopStatus = 400; + expect(await transport.stopGeneration("chat")).toBe(false); + expect(transport.getSession("chat")?.isStreaming).toBe(false); + await vi.waitFor(() => expect(outputs[0]?.destroyed).toBe(true)); + const next = await send(); + emit(oldTailAndReply(10, 12)); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("does not resume a stopped owning consumer after hydration", async () => { + const abort = new AbortController(); + await send(abort.signal); + abort.abort(); + expect(saved).toMatchObject({ + skipToTurnComplete: true, + supersededInputSeq: 10, + isStreaming: false, + }); + if (!saved) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport(saved); + expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull(); + }); + + it("closes the SSE connection when the stopped boundary is missing", async () => { + await send(); + await transport.stopGeneration("chat"); + const next = await send(); + emit([...reply(1), complete(6, 12)]); + await expect(readText(next)).rejects.toThrow("The previous turn's output was lost"); + await vi.waitFor(() => expect(outputs.at(-1)?.destroyed).toBe(true)); + expect(saved).toMatchObject({ + skipToTurnComplete: false, + isStreaming: false, + activeInputSeq: undefined, + }); + }); + + it("does not retain an old stopped input after session recreation", async () => { + transport.dispose(); + transport = createTransport( + { publicAccessToken: "test-token" }, + { + startSession: async () => { + stopStatus = 200; + return { publicAccessToken: "replacement-token" }; + }, + } + ); + await send(); + stopStatus = 404; + expect(await transport.stopGeneration("chat")).toBe(true); + expect(saved).toMatchObject({ + skipToTurnComplete: false, + supersededInputSeq: undefined, + activeInputSeq: undefined, + }); + await transport.stopGeneration("chat"); + const next = await send(); + emit([...reply(1), complete(6, 14)]); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("does not gate a response after a completed turn", async () => { + const first = await send(); + emit([...reply(1), complete(6, 10)]); + await expect(readText(first)).resolves.toBe("New response"); + await transport.stopGeneration("chat"); + const next = await send(); + emit([...reply(7), complete(12, 12)]); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("retains the first stopped boundary after repeated Stop calls", async () => { + await send(); + await transport.stopGeneration("chat"); + await transport.stopGeneration("chat"); + const next = await send(); + emit(oldTailAndReply(12, 13)); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("does not stop a successor when an old consumer aborts after settled EOF", async () => { + const abort = new AbortController(); + settled = true; + const first = await send(abort.signal); + emit(reply(1)); + outputs[0]!.end(); + await expect(readText(first)).resolves.toBe("New response"); + settled = false; + const next = await send(); + abort.abort(); + expect(transport.getSession("chat")?.isStreaming).toBe(true); + emit([...reply(6), complete(11, 11)]); + await expect(readText(next)).resolves.toBe("New response"); + expect(inputSeq).toBe(12); + }); + + it.each(["message", "action"] as const)( + "requires transcript reload when a stopped successor %s has no sequence", + async (kind) => { + await send(); + await transport.stopGeneration("chat"); + includeSequence = false; + const reloadError = + "Stopped chat response cannot be matched. Reload the chat before sending another message."; + const next = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" }); + await expect(next).rejects.toThrow(reloadError); + expect(saved).toMatchObject({ requiresTranscriptReload: true, isStreaming: false }); + expect(inputSeq).toBe(13); + await expect(send()).rejects.toThrow(reloadError); + await expect(transport.sendAction("chat", { type: "undo" })).rejects.toThrow(reloadError); + expect(inputSeq).toBe(13); + + if (!saved) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport(saved, { watch: true }); + await expect(send()).rejects.toThrow(reloadError); + expect(inputSeq).toBe(13); + expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull(); + expect(outputs).toHaveLength(1); + + // A fresh transcript supplies a cursor beyond the accepted response. + transport.dispose(); + transport = createTransport(saved); + transport.setSession("chat", { + publicAccessToken: "test-token", + lastEventId: "11", + isStreaming: false, + }); + includeSequence = true; + const afterReload = await send(); + emit([...reply(12), complete(17, 13)]); + await expect(readText(afterReload)).resolves.toBe("New response"); + } + ); + + it("rejects a newer snapshot from before the stopped input", async () => { + await hydrateBlockedSession("constructor"); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(recover({ lastOutEventId: "5", lastInEventId: "9" })).toBe(false); + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + expect(inputSeq).toBe(13); + }); + + it.each([ + ["explicit", "constructor"], + ["explicit", "setSession"], + ["abort", "constructor"], + ["abort", "setSession"], + ] as const)( + "rejects an older snapshot after cursor-free %s Stop and %s hydration", + async (stopMode, hydrate) => { + const abort = new AbortController(); + const resumed = await transport.reconnectToStream({ + chatId: "chat", + abortSignal: abort.signal, + stopOnAbort: true, + }); + if (!resumed) throw new Error("Expected a resumed stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + if (stopMode === "explicit") await transport.stopGeneration("chat"); + else abort.abort(); + await vi.waitFor(() => + expect(transport.getSession("chat")?.transcriptRecoveryInputSeq).toBe(10) + ); + includeSequence = false; + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + const session = transport.getSession("chat"); + if (!session) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport( + hydrate === "constructor" ? session : { publicAccessToken: "test-token" } + ); + if (hydrate === "setSession") transport.setSession("chat", session); + const stale = transport.prepareTranscriptRecovery("chat"); + if (!stale) throw new Error("Expected transcript recovery"); + expect(stale({ lastOutEventId: "5", lastInEventId: "9" })).toBe(false); + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + const fresh = transport.prepareTranscriptRecovery("chat"); + if (!fresh) throw new Error("Expected transcript recovery"); + expect(fresh({ lastOutEventId: "11", lastInEventId: "10" })).toBe(true); + const next = await transport.reconnectToStream({ chatId: "chat" }); + if (!next) throw new Error("Expected a resumed stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(2)); + emit([...reply(12), complete(17, 11)]); + await expect(readText(next)).resolves.toBe("New response"); + expect(inputSeq).toBe(12); + } + ); + + it.each([false, true])( + "retains the first Stop input through repeated Stop (delayed first acknowledgment: %s)", + async (delayed) => { + holdStop = delayed; + const firstStop = transport.stopGeneration("chat"); + if (delayed) await vi.waitFor(() => expect(pendingStop).toBeDefined()); + else expect(await firstStop).toBe(true); + includeSequence = false; + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + includeSequence = true; + holdStop = false; + expect(await transport.stopGeneration("chat")).toBe(true); + if (delayed) { + const pending = pendingStop!; + appendResponse(pending.response, pending.seq); + expect(await firstStop).toBe(true); + } + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(recover({ lastOutEventId: "17", lastInEventId: "11" })).toBe(true); + const next = await send(); + emit([...reply(18), complete(23, 13)]); + await expect(readText(next)).resolves.toBe("New response"); + } + ); + + it("does not install an old Stop sequence into a replacement session", async () => { + holdStop = true; + const stopped = transport.stopGeneration("chat"); + await vi.waitFor(() => expect(pendingStop).toBeDefined()); + transport.setSession("chat", { + publicAccessToken: "replacement-token", + skipToTurnComplete: true, + requiresTranscriptReload: true, + isStreaming: false, + }); + const pending = pendingStop!; + appendResponse(pending.response, pending.seq); + expect(await stopped).toBe(true); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(() => recover({ lastOutEventId: "17", lastInEventId: "11" })).toThrow( + "Transcript recovery requires a stopped input sequence" + ); + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + holdStop = false; + expect(await transport.stopGeneration("chat")).toBe(true); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "17", lastInEventId: "11" }) + ).toBe(true); + }); + + it("accepts a sequence-free response without a stopped boundary", async () => { + includeSequence = false; + const stream = await send(); + emit([...reply(1), complete(6, 10)]); + await expect(readText(stream)).resolves.toBe("New response"); + }); + + it("rejects steering before an append when transcript reload is required", async () => { + await send(); + await transport.stopGeneration("chat"); + includeSequence = false; + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + const accepted = await transport.sendPendingMessage("chat", { + id: "steering-message", + role: "user", + parts: [{ type: "text", text: "Use the new instructions" }], + }); + expect({ accepted, inputSeq }).toEqual({ accepted: false, inputSeq: 13 }); + }); + + it.each([ + ["constructor", "reconnect"], + ["constructor", "send"], + ["setSession", "reconnect"], + ["setSession", "send"], + ] as const)( + "recovers a blocked session through %s hydration and %s after a fresh transcript", + async (hydrate, operation) => { + await hydrateBlockedSession(hydrate); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(recover({ lastOutEventId: "11", lastInEventId: "11" })).toBe(true); + transport.seedResumeCursor("chat", "11"); + expect(saved).toMatchObject({ + lastEventId: "11", + requiresTranscriptReload: false, + skipToTurnComplete: false, + supersededInputSeq: undefined, + activeInputSeq: undefined, + }); + const before = outputs.length; + const stream = + operation === "send" ? await send() : await transport.reconnectToStream({ chatId: "chat" }); + if (!stream) throw new Error("Expected a response stream"); + await vi.waitFor(() => expect(outputs.length).toBeGreaterThan(before)); + emit([...reply(12), complete(17, operation === "send" ? 13 : 12)]); + await expect(readText(stream)).resolves.toBe("New response"); + expect(inputSeq).toBe(operation === "send" ? 14 : 13); + } + ); + + it.each([undefined, "0", "1"])( + "retains the reload guard for a missing or stale transcript cursor (%s)", + async (cursor) => { + await hydrateBlockedSession("constructor"); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + if (cursor === undefined) { + expect(() => recover({ lastOutEventId: cursor, lastInEventId: "11" })).toThrow( + "Transcript recovery requires numeric input and output cursors" + ); + } else { + expect(recover({ lastOutEventId: cursor, lastInEventId: "11" })).toBe(false); + } + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + await expect(transport.sendAction("chat", { type: "undo" })).rejects.toThrow( + "Stopped chat response cannot be matched" + ); + expect( + await transport.sendPendingMessage("chat", { + id: "steering-message", + role: "user", + parts: [{ type: "text", text: "Use the new instructions" }], + }) + ).toBe(false); + expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull(); + expect(inputSeq).toBe(13); + expect(transport.getSession("chat")).toMatchObject({ + lastEventId: "1", + requiresTranscriptReload: true, + }); + } + ); + + it.each(["none", "constructor", "setSession"] as const)( + "resumes accepted output after a stopped checkpoint (recovery hydration: %s)", + async (hydrate) => { + await hydrateBlockedSession("constructor"); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(recover({ lastOutEventId: "11", lastInEventId: "11" })).toBe(true); + if (hydrate !== "none") { + const session = transport.getSession("chat"); + if (!session) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport( + hydrate === "constructor" ? session : { publicAccessToken: "test-token" } + ); + if (hydrate === "setSession") transport.setSession("chat", session); + } + resumeAfterStoppedCheckpoint = true; + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + await expect(readText(resumed)).resolves.toBe("New response"); + expect(inputSeq).toBe(13); + expect(transport.getSession("chat")).toMatchObject({ + skipSettledPeek: false, + isStreaming: false, + }); + } + ); + + it("reports an empty recovery poll and reconnects without another append", async () => { + await hydrateBlockedSession("constructor"); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(recover({ lastOutEventId: "11", lastInEventId: "11" })).toBe(true); + const events: ChatTransportEvent[] = []; + transport.setOnEvent((event) => events.push(event)); + resumeAfterStoppedCheckpoint = true; + emptyRecoveredOutput = true; + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + await expect(readText(resumed)).rejects.toThrow( + "Chat recovery received no output before the poll ended. Reconnect to resume the accepted message." + ); + expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1); + const request = outputHeaders.at(-1); + if (!request) throw new Error("Expected a stream request"); + expect(request.peek).toBe(false); + expect(request.timeout).toBeGreaterThan(0); + expect(request.timeout).toBeLessThanOrEqual(30); + expect(outputs).toHaveLength(2); + expect(transport.getSession("chat")).toMatchObject({ + lastEventId: "11", + isStreaming: undefined, + skipSettledPeek: true, + }); + emptyRecoveredOutput = false; + const next = await transport.reconnectToStream({ chatId: "chat" }); + if (!next) throw new Error("Expected a resumed stream"); + await expect(readText(next)).resolves.toBe("New response"); + expect(inputSeq).toBe(13); + }); + + it("rejects captured recovery after an explicit Stop before its acknowledgment", async () => { + await hydrateBlockedSession("constructor"); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + holdStop = true; + const stopped = transport.stopGeneration("chat"); + await vi.waitFor(() => expect(pendingStop).toBeDefined()); + expect(recover({ lastOutEventId: "11", lastInEventId: "11" })).toBe(false); + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + const pending = pendingStop!; + appendResponse(pending.response, pending.seq); + expect(await stopped).toBe(true); + expect(inputSeq).toBe(14); + expect(transport.getSession("chat")).toMatchObject({ requiresTranscriptReload: true }); + }); + + it.each(["abort", "stop"] as const)("closes recovery quietly after %s", async (operation) => { + await hydrateBlockedSession("constructor"); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" }) + ).toBe(true); + const events: ChatTransportEvent[] = []; + transport.setOnEvent((event) => events.push(event)); + const abort = new AbortController(); + const resumed = await transport.reconnectToStream({ + chatId: "chat", + abortSignal: abort.signal, + }); + if (!resumed) throw new Error("Expected a resumed stream"); + const result = readText(resumed); + await vi.waitFor(() => expect(outputs).toHaveLength(2)); + if (operation === "abort") abort.abort(); + else await transport.stopGeneration("chat"); + await expect(result).resolves.toBe(""); + expect(events.filter((event) => event.type === "stream-error")).toEqual([]); + if (operation === "abort") { + expect(transport.getSession("chat")).toMatchObject({ + isStreaming: undefined, + skipSettledPeek: true, + }); + expect(inputSeq).toBe(13); + } else { + expect(transport.getSession("chat")).toMatchObject({ + isStreaming: false, + skipToTurnComplete: true, + }); + expect(inputSeq).toBe(14); + } + }); + + it("closes an empty settled recovery without an error", async () => { + await hydrateBlockedSession("constructor"); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" }) + ).toBe(true); + const events: ChatTransportEvent[] = []; + transport.setOnEvent((event) => events.push(event)); + settled = true; + emptyRecoveredOutput = true; + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + await expect(readText(resumed)).resolves.toBe(""); + expect(events.filter((event) => event.type === "stream-error")).toEqual([]); + expect(transport.getSession("chat")?.isStreaming).toBe(false); + expect(inputSeq).toBe(13); + }); + + it("reconnects a watch after an empty recovery response", async () => { + await hydrateBlockedSession("constructor"); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" }) + ).toBe(true); + const session = transport.getSession("chat")!; + transport.dispose(); + const events: ChatTransportEvent[] = []; + transport = createTransport(session, { watch: true, onEvent: (event) => events.push(event) }); + emptyRecoveredOutput = true; + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + const result = readWatchedTurn(resumed); + await vi.waitFor(() => expect(outputs.length).toBeGreaterThanOrEqual(2)); + emptyRecoveredOutput = false; + resumeAfterStoppedCheckpoint = true; + await expect(result).resolves.toBe("New response"); + expect(events.filter((event) => event.type === "stream-error")).toEqual([]); + expect(inputSeq).toBe(13); + }); + + it("uses the active-turn retry policy after recovery receives data", async () => { + await hydrateBlockedSession("constructor"); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" }) + ).toBe(true); + const events: ChatTransportEvent[] = []; + transport.setOnEvent((event) => events.push(event)); + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + const reader = resumed.getReader(); + await vi.waitFor(() => expect(outputs).toHaveLength(2)); + emit([chunk(12, { type: "start", messageId: "new" })]); + await expect(reader.read()).resolves.toMatchObject({ value: { type: "start" } }); + outputs.at(-1)!.end(); + await vi.waitFor(() => expect(outputs).toHaveLength(3)); + emit([...reply(12).slice(1), complete(17, 12)]); + while (!(await reader.read()).done) {} + expect(events.filter((event) => event.type === "stream-error")).toEqual([]); + expect(transport.getSession("chat")?.isStreaming).toBe(false); + expect(inputSeq).toBe(13); + }); + + it("does not let a replaced recovery stream settle the new response", async () => { + await hydrateBlockedSession("constructor"); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" }) + ).toBe(true); + const events: ChatTransportEvent[] = []; + transport.setOnEvent((event) => events.push(event)); + const resumed = await transport.reconnectToStream({ chatId: "chat" }); + if (!resumed) throw new Error("Expected a resumed stream"); + const oldResult = readText(resumed); + await vi.waitFor(() => expect(outputs).toHaveLength(2)); + const replacement = await send(); + await expect(oldResult).resolves.toBe(""); + expect(transport.getSession("chat")?.isStreaming).toBe(true); + emit([...reply(12), complete(17, 13)]); + await expect(readText(replacement)).resolves.toBe("New response"); + expect(events.filter((event) => event.type === "stream-error")).toEqual([]); + expect(inputSeq).toBe(14); + }); + + it.each(["constructor", "setSession"] as const)( + "retains the abandoned-turn marker through %s hydration in watch mode", + async (hydrate) => { + await send(); + transport.clearSupersedeGate("chat"); + const session = transport.getSession("chat"); + if (!session) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport( + hydrate === "constructor" ? session : { publicAccessToken: "test-token" }, + { watch: true } + ); + if (hydrate === "setSession") transport.setSession("chat", session); + const watched = await transport.reconnectToStream({ chatId: "chat" }); + if (!watched) throw new Error("Expected a watch stream"); + await vi.waitFor(() => expect(outputs).toHaveLength(2)); + const reader = watched.getReader(); + emit([chunk(1, { type: "start", messageId: "abandoned" })]); + await reader.read(); + await transport.stopGeneration("chat"); + const next = await send(); + emit([...reply(2), complete(7, 12)]); + await expect(readWatchedTurn(next)).resolves.toBe("New response"); + expect(session).toMatchObject({ outstandingTurnAbandoned: true }); + } + ); + + it("persists a cleared boundary without rearming the abandoned turn after hydration", async () => { + await send(); + transport.clearSupersedeGate("chat"); + expect(saved).toMatchObject({ + skipToTurnComplete: false, + supersededInputSeq: undefined, + activeInputSeq: undefined, + isStreaming: false, + }); + if (!saved) throw new Error("Expected persisted state"); + transport.dispose(); + transport = createTransport(saved); + await transport.stopGeneration("chat"); + const next = await send(); + emit([...reply(1), complete(6, 12)]); + await expect(readText(next)).resolves.toBe("New response"); + }); +}); diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index c3157a5bbff..afe31e84d6d 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1438,7 +1438,10 @@ describe("TriggerChatTransport", () => { it("does not gate a stop with no turn outstanding", async () => { mockFetch([() => defaultSseResponse()]); - const transport = await armedGate("chat-idle-stop", { publicAccessToken: "p" }); + const transport = await armedGate("chat-idle-stop", { + publicAccessToken: "p", + isStreaming: false, + }); const stream = await send(transport, "chat-idle-stop"); @@ -1577,13 +1580,13 @@ describe("TriggerChatTransport", () => { ]); }); - it("keeps the gate out of the persisted session", async () => { + it("persists the stopped boundary", async () => { mockFetch([() => defaultSseResponse()]); const sessions: Record = {}; const transport = await armedGate( "chat-persist", - { publicAccessToken: "p" }, + { publicAccessToken: "p", isStreaming: true, activeInputSeq: 5 }, { onSessionChange: (chatId, session) => { sessions[chatId] = session; @@ -1591,8 +1594,14 @@ describe("TriggerChatTransport", () => { } ); - expect(transport.getSession("chat-persist")).not.toHaveProperty("skipToTurnComplete"); - expect(sessions["chat-persist"]).not.toHaveProperty("skipToTurnComplete"); + expect(transport.getSession("chat-persist")).toMatchObject({ + skipToTurnComplete: true, + supersededInputSeq: 5, + }); + expect(sessions["chat-persist"]).toMatchObject({ + skipToTurnComplete: true, + supersededInputSeq: 5, + }); }); it("clears on the first turn-complete after two consecutive stops", async () => { diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index b3c10316df9..4a277cc7645 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -45,6 +45,7 @@ function byteLength(body: string): number { return new TextEncoder().encode(body).byteLength; } import { ChatTabCoordinator } from "./chat-tab-coordinator.js"; +import type { TranscriptCursors } from "./transcriptStorage.js"; import { MAX_EOF_RESUBSCRIBES, slimSubmitMessageForWire, @@ -496,6 +497,18 @@ export type ChatSessionPersistedState = { /** The `.in` append sequence of the last send this client owned; reused as `sinceInSeq` on reconnect. */ activeInputSeq?: number; isStreaming?: boolean; + /** Discard unread output from a stopped turn before the next response. */ + skipToTurnComplete?: boolean; + /** The stopped input sequence excludes older completion records from the boundary. */ + supersededInputSeq?: number; + /** A send lacks an input sequence while stopped output remains unread. Reload the transcript before another send. */ + requiresTranscriptReload?: boolean; + /** The acknowledged Stop sequence proves a transcript covers an unknown stopped input. */ + transcriptRecoveryInputSeq?: number; + /** An abandoned turn must not restore its discard boundary after a reload. */ + outstandingTurnAbandoned?: boolean; + /** A recovered transcript can precede accepted output. Wait for output instead of peeking at the previous completion. */ + skipSettledPeek?: boolean; /** Set once the session is closed. Persisted so a reload doesn't retry a dead session. */ closed?: boolean; /** The reason the session was closed, when one was given. */ @@ -698,29 +711,11 @@ export type TriggerChatTransportOptions = { * `end-and-continue`, etc. * @internal */ -type ChatSessionState = { - /** Session-scoped PAT — `read:sessions:{chatId} + write:sessions:{chatId}`. */ - publicAccessToken: string; - /** Last SSE event ID — used to resume the stream without replaying old events. */ - lastEventId?: string; - /** `.in` append sequence used to filter stale turn boundaries after reconnecting. */ - activeInputSeq?: number; - /** - * Set when the stream was aborted mid-turn (stop). Skip chunks until the - * stopped turn's trigger:turn-complete — survives a reconnect and a retry - * send, so the stopped turn's tail never renders into the new turn. - */ - skipToTurnComplete?: boolean; - /** `.in` seq of the turn the gate supersedes; only its boundary (or a later one) clears the gate. */ - supersededInputSeq?: number; - /** Whether the agent is currently streaming a response. Set on first chunk, cleared on turn-complete. */ - isStreaming?: boolean; - /** Set once the outstanding turn is declared dead: a later stop must not gate the next turn on it. */ - outstandingTurnAbandoned?: boolean; - /** Set once the session is closed. Terminal — sends and reconnects stop. */ - closed?: boolean; - /** The reason the session was closed, when one was given. */ - closedReason?: string; +type ChatSessionState = ChatSessionPersistedState & { + /** Identifies the stopped boundary for pending Stop acknowledgments. Never persisted. */ + stoppedBoundary?: symbol; + /** Identifies the latest transcript load for this blocked session. Never persisted. */ + transcriptRecovery?: symbol; }; /** @@ -802,14 +797,7 @@ export class TriggerChatTransport implements ChatTransport { if (options.sessions) { for (const [chatId, session] of Object.entries(options.sessions)) { - this.sessions.set(chatId, { - publicAccessToken: session.publicAccessToken, - lastEventId: session.lastEventId, - activeInputSeq: session.activeInputSeq, - isStreaming: session.isStreaming, - closed: session.closed, - closedReason: session.closedReason, - }); + this.sessions.set(chatId, this.toPersisted(session)); } } } @@ -951,6 +939,7 @@ export class TriggerChatTransport implements ChatTransport { // Generated outside the closure so auth-retries reuse the same part id // and the server-side dedupe sees one logical append. + this.assertTranscriptReady(chatId, state); const partId = crypto.randomUUID(); const serializedBody = this.serializeInputChunk({ kind: "message", payload: wirePayload }); const sendChatMessage = (token: string) => @@ -975,8 +964,10 @@ export class TriggerChatTransport implements ChatTransport { } state.activeInputSeq = inSeq; + this.requireStoppedTurnCorrelation(chatId, state, inSeq); state.isStreaming = true; state.outstandingTurnAbandoned = false; + state.skipSettledPeek = false; this.notifySessionChange(chatId, state); // Owning turn: aborting this live send stops the turn the user drives. @@ -1225,7 +1216,7 @@ export class TriggerChatTransport implements ChatTransport { metadata?: Record ): Promise => { const state = this.sessions.get(chatId); - if (!state) return false; + if (!state || state.requiresTranscriptReload) return false; const mergedMetadata = this.defaultMetadata || metadata @@ -1280,6 +1271,7 @@ export class TriggerChatTransport implements ChatTransport { if (!state) return null; // A closed session has no further turns to resume. if (state.closed) return null; + if (state.requiresTranscriptReload) return null; // Watch is a standing subscription: a settled session is exactly the // state it waits in, so a completed last turn must not block the resume. @@ -1303,7 +1295,7 @@ export class TriggerChatTransport implements ChatTransport { // can remain at the tail until the current turn writes its first chunk. // Watch mode must NOT peek: a settled peek between turns closes the // standing subscription, so the viewer never sees the next turn. - peekSettled: !this.watchMode && state.activeInputSeq === undefined, + peekSettled: !this.watchMode && !state.skipSettledPeek && state.activeInputSeq === undefined, }); }; @@ -1315,33 +1307,20 @@ export class TriggerChatTransport implements ChatTransport { stopGeneration = async (chatId: string): Promise => { const state = this.sessions.get(chatId); if (!state) return false; - - const partId = crypto.randomUUID(); - const serializedBody = this.serializeInputChunk({ kind: "stop" }); - const send = async (token: string) => { - await this.appendInputChunk(chatId, token, serializedBody, partId); - }; - - try { - await this.sendWithEvents( - chatId, - "stop", - { partId, bodyBytes: byteLength(serializedBody) }, - () => this.callWithAuthRetry(chatId, state, send) - ); - } catch { - return false; - } - + state.transcriptRecovery = undefined; + // Close the captured turn before the request awaits. A delayed acknowledgment + // must not change a successor's reader or stopped-output boundary. // Only gate when a sent turn is still outstanding. A stop at a boundary has // nothing to supersede, and gating it would swallow the next turn. if ( !state.outstandingTurnAbandoned && - (state.isStreaming || state.activeInputSeq !== undefined) + (state.isStreaming !== false || state.activeInputSeq !== undefined) ) { - state.skipToTurnComplete = true; - state.supersededInputSeq = state.activeInputSeq; + this.armStoppedBoundary(state); } + const stoppedBoundary = state.skipToTurnComplete + ? (state.stoppedBoundary ??= Symbol("stopped-boundary")) + : undefined; const activeStream = this.activeStreams.get(chatId); if (activeStream) { @@ -1360,7 +1339,23 @@ export class TriggerChatTransport implements ChatTransport { // explicitly stopped. state.isStreaming = false; this.notifySessionChange(chatId, state); - return true; + + const partId = crypto.randomUUID(); + const serializedBody = this.serializeInputChunk({ kind: "stop" }); + const send = (token: string) => this.appendInputChunk(chatId, token, serializedBody, partId); + try { + const inSeq = await this.sendWithEvents( + chatId, + "stop", + { partId, bodyBytes: byteLength(serializedBody) }, + () => this.callWithAuthRetry(chatId, state, send) + ); + this.recordStoppedInput(chatId, state, stoppedBoundary, inSeq); + return true; + } catch { + // The reader already closed. Retain its unread boundary for the next send. + return false; + } }; /** @@ -1371,9 +1366,12 @@ export class TriggerChatTransport implements ChatTransport { clearSupersedeGate = (chatId: string): void => { const state = this.sessions.get(chatId); if (!state) return; - state.skipToTurnComplete = false; - state.supersededInputSeq = undefined; + this.clearStoppedBoundary(state); + state.activeInputSeq = undefined; state.outstandingTurnAbandoned = true; + state.skipSettledPeek = false; + state.isStreaming = false; + this.notifySessionChange(chatId, state); }; /** @@ -1408,6 +1406,7 @@ export class TriggerChatTransport implements ChatTransport { : undefined, }; + this.assertTranscriptReady(chatId, state); const body = this.serializeInputChunk({ kind: "message", payload: wirePayload }); const partId = crypto.randomUUID(); const send = (token: string) => this.appendInputChunk(chatId, token, body, partId); @@ -1432,8 +1431,10 @@ export class TriggerChatTransport implements ChatTransport { // Mark streaming + persist so a reload mid-action resumes (reconnectToStream // no-ops when the persisted session says isStreaming: false). state.activeInputSeq = inSeq; + this.requireStoppedTurnCorrelation(chatId, state, inSeq); state.isStreaming = true; state.outstandingTurnAbandoned = false; + state.skipSettledPeek = false; this.notifySessionChange(chatId, state); // Owning action: aborting this send stops the turn the user drives. @@ -1454,16 +1455,12 @@ export class TriggerChatTransport implements ChatTransport { }; setSession(chatId: string, session: ChatSessionPersistedState): void { - this.sessions.set( - chatId, - this.applyPendingResumeCursor(chatId, { - publicAccessToken: session.publicAccessToken, - lastEventId: session.lastEventId, - activeInputSeq: session.activeInputSeq, - isStreaming: session.isStreaming, - }) - ); - this.notifySessionChange(chatId, this.toPersisted(this.sessions.get(chatId)!)); + const state = this.toPersisted(session); + // Explicit session replacement resets the closed state, unlike constructor hydration. + state.closed = undefined; + state.closedReason = undefined; + this.sessions.set(chatId, this.applyPendingResumeCursor(chatId, state)); + this.notifySessionChange(chatId, state); } /** @@ -1479,7 +1476,7 @@ export class TriggerChatTransport implements ChatTransport { if (existing?.publicAccessToken) { if (existing.lastEventId === undefined) { existing.lastEventId = lastEventId; - this.notifySessionChange(chatId, this.toPersisted(existing)); + this.notifySessionChange(chatId, existing); } this.pendingResumeCursors.delete(chatId); return; @@ -1487,6 +1484,86 @@ export class TriggerChatTransport implements ChatTransport { this.pendingResumeCursors.set(chatId, lastEventId); }; + /** + * Capture recovery before a fresh transcript load starts. The returned callback + * accepts a newer saved cursor once, with input evidence for the stopped boundary. + * Ordinary transcript loads do not reset session state. + * Stale recovery returns false. Missing or invalid recovery evidence throws an error. + */ + prepareTranscriptRecovery = ( + chatId: string + ): ((cursors: TranscriptCursors | undefined) => boolean) | undefined => { + const state = this.sessions.get(chatId); + if (!state?.requiresTranscriptReload || state.closed) return undefined; + + const token = Symbol("transcript-recovery"); + state.transcriptRecovery = token; + const { + lastEventId, + activeInputSeq, + skipToTurnComplete, + supersededInputSeq, + transcriptRecoveryInputSeq, + } = state; + const stoppedInputSeq = supersededInputSeq ?? transcriptRecoveryInputSeq; + return (cursors) => { + if (state.transcriptRecovery !== token) return false; + state.transcriptRecovery = undefined; + if ( + this.sessions.get(chatId) !== state || + !state.requiresTranscriptReload || + state.closed || + this.activeStreams.has(chatId) || + state.lastEventId !== lastEventId || + state.activeInputSeq !== activeInputSeq || + state.skipToTurnComplete !== skipToTurnComplete || + state.supersededInputSeq !== supersededInputSeq || + state.transcriptRecoveryInputSeq !== transcriptRecoveryInputSeq + ) { + return false; + } + if ( + stoppedInputSeq === undefined || + !Number.isSafeInteger(stoppedInputSeq) || + stoppedInputSeq < 0 + ) { + throw new Error( + "Transcript recovery requires a stopped input sequence. Stop the chat again, then reload its transcript." + ); + } + const loadedEventId = cursors?.lastOutEventId; + const loadedInEventId = cursors?.lastInEventId; + if ( + loadedInEventId === undefined || + !/^\d+$/.test(loadedInEventId) || + loadedEventId === undefined || + !/^\d+$/.test(loadedEventId) + ) { + throw new Error( + "Transcript recovery requires numeric input and output cursors. Return both cursors from the transcript loader." + ); + } + if ( + BigInt(loadedInEventId) < BigInt(stoppedInputSeq) || + (lastEventId !== undefined && + (!/^\d+$/.test(lastEventId) || BigInt(loadedEventId) <= BigInt(lastEventId))) + ) { + return false; + } + + state.lastEventId = loadedEventId; + state.requiresTranscriptReload = false; + this.clearStoppedBoundary(state); + state.activeInputSeq = undefined; + state.outstandingTurnAbandoned = false; + state.skipSettledPeek = true; + // The saved transcript can precede the accepted response. Resume from its checkpoint. + state.isStreaming = undefined; + this.notifySessionChange(chatId, state); + return true; + }; + }; + private applyPendingResumeCursor(chatId: string, state: ChatSessionState): ChatSessionState { const pending = this.pendingResumeCursors.get(chatId); if (pending !== undefined && state.lastEventId === undefined) { @@ -1633,6 +1710,10 @@ export class TriggerChatTransport implements ChatTransport { controller.abort(); } this.activeStreams.clear(); + for (const state of this.sessions.values()) { + state.transcriptRecovery = undefined; + state.stoppedBoundary = undefined; + } this.coordinator?.dispose(); this.coordinator = null; } @@ -1641,15 +1722,85 @@ export class TriggerChatTransport implements ChatTransport { // Internal helpers // ------------------------------------------------------------------------- + private armStoppedBoundary(state: ChatSessionState): symbol { + state.skipToTurnComplete = true; + state.supersededInputSeq = state.activeInputSeq; + state.transcriptRecoveryInputSeq = undefined; + state.transcriptRecovery = undefined; + return (state.stoppedBoundary = Symbol("stopped-boundary")); + } + + private clearStoppedBoundary(state: ChatSessionState): void { + state.skipToTurnComplete = false; + state.supersededInputSeq = undefined; + state.transcriptRecoveryInputSeq = undefined; + state.transcriptRecovery = undefined; + state.stoppedBoundary = undefined; + } + + private recordStoppedInput( + chatId: string, + state: ChatSessionState, + stoppedBoundary: symbol | undefined, + inSeq: number | undefined + ): void { + if ( + stoppedBoundary === undefined || + state.stoppedBoundary !== stoppedBoundary || + this.sessions.get(chatId) !== state || + !state.skipToTurnComplete || + state.closed || + state.supersededInputSeq !== undefined || + inSeq === undefined || + !Number.isSafeInteger(inSeq) || + inSeq < 0 || + (state.transcriptRecoveryInputSeq !== undefined && state.transcriptRecoveryInputSeq <= inSeq) + ) { + return; + } + // Keep the Stop sequence separate from the stopped input sequence. + state.transcriptRecoveryInputSeq = inSeq; + state.transcriptRecovery = undefined; + this.notifySessionChange(chatId, state); + } + private serializeInputChunk(chunk: ChatInputChunk): string { return JSON.stringify(chunk); } + private assertTranscriptReady(chatId: string, state: ChatSessionState): void { + if (!state.requiresTranscriptReload) return; + this.coordinator?.release(chatId); + throw new Error( + "Stopped chat response cannot be matched. Reload the chat before sending another message." + ); + } + + private requireStoppedTurnCorrelation( + chatId: string, + state: ChatSessionState, + inSeq: number | undefined + ): void { + if (!state.skipToTurnComplete || inSeq !== undefined) return; + // The server accepted the prompt. A retry can create a duplicate turn. + state.requiresTranscriptReload = true; + state.transcriptRecovery = undefined; + state.isStreaming = false; + this.notifySessionChange(chatId, state); + this.assertTranscriptReady(chatId, state); + } + private toPersisted = (state: ChatSessionState): ChatSessionPersistedState => ({ publicAccessToken: state.publicAccessToken, lastEventId: state.lastEventId, activeInputSeq: state.activeInputSeq, isStreaming: state.isStreaming, + skipToTurnComplete: state.skipToTurnComplete, + supersededInputSeq: state.supersededInputSeq, + requiresTranscriptReload: state.requiresTranscriptReload, + transcriptRecoveryInputSeq: state.transcriptRecoveryInputSeq, + outstandingTurnAbandoned: state.outstandingTurnAbandoned, + skipSettledPeek: state.skipSettledPeek, closed: state.closed, closedReason: state.closedReason, }); @@ -1667,11 +1818,11 @@ export class TriggerChatTransport implements ChatTransport { if (state.closed) return; state.closed = true; + state.skipSettledPeek = false; if (reason) state.closedReason = reason; state.isStreaming = false; state.activeInputSeq = undefined; - state.skipToTurnComplete = false; - state.supersededInputSeq = undefined; + this.clearStoppedBoundary(state); this.emitEvent({ type: "session-closed", @@ -1961,7 +2112,13 @@ export class TriggerChatTransport implements ChatTransport { } state.publicAccessToken = publicAccessToken; state.lastEventId = undefined; + state.activeInputSeq = undefined; state.isStreaming = false; + this.clearStoppedBoundary(state); + state.requiresTranscriptReload = false; + state.outstandingTurnAbandoned = false; + state.skipSettledPeek = false; + state.transcriptRecovery = undefined; this.sessions.set(chatId, state); this.notifySessionChange(chatId, state); } @@ -2001,15 +2158,23 @@ export class TriggerChatTransport implements ChatTransport { // has no turn to stop: don't gate the next one, don't write a stop. const outstanding = !state.outstandingTurnAbandoned && - (state.isStreaming || state.activeInputSeq !== undefined); - if (options?.sendStopOnAbort !== false && outstanding && !internalAbort.signal.aborted) { - state.skipToTurnComplete = true; - state.supersededInputSeq = state.activeInputSeq; + (state.isStreaming !== false || state.activeInputSeq !== undefined); + if ( + options?.sendStopOnAbort !== false && + outstanding && + !internalAbort.signal.aborted && + this.activeStreams.get(chatId) === internalAbort + ) { + const stoppedBoundary = this.armStoppedBoundary(state); + state.isStreaming = false; + this.notifySessionChange(chatId, state); this.appendInputChunk( chatId, state.publicAccessToken, this.serializeInputChunk({ kind: "stop" }) - ).catch(() => {}); + ) + .then((inSeq) => this.recordStoppedInput(chatId, state, stoppedBoundary, inSeq)) + .catch(() => {}); } internalAbort.abort(); }, @@ -2088,7 +2253,12 @@ export class TriggerChatTransport implements ChatTransport { ...(options?.peekSettled ? { "X-Peek-Settled": "1" } : {}), }, signal: combinedSignal, - timeoutInSeconds: this.streamTimeoutSeconds, + // A saved transcript can already include the accepted response. Bound the + // initial recovery poll below the stall deadline without claiming completion. + timeoutInSeconds: + state.skipSettledPeek && state.isStreaming === undefined + ? Math.min(this.streamTimeoutSeconds, 30) + : this.streamTimeoutSeconds, lastEventId: state.lastEventId, // Reconnect if no decoded record arrives for 60 seconds. stallTimeoutMs: 60_000, @@ -2129,6 +2299,9 @@ export class TriggerChatTransport implements ChatTransport { }; let eofResubscribes = 0; + const emptyRecoveryError = new Error( + "Chat recovery received no output before the poll ended. Reconnect to resume the accepted message." + ); const resumeAfterEof = async () => { // Watch mode is a standing subscription: it outlives turn-complete @@ -2146,7 +2319,7 @@ export class TriggerChatTransport implements ChatTransport { if (opened) return opened; } - // A settled session or an abort ends the turn cleanly. Exhausting the + // A settled session or an abort ends the subscription cleanly. Exhausting the // resubscribe budget while the turn is still streaming means it was cut // off — surface an error so the UI doesn't read a truncated reply as // complete. The caller's catch emits stream-error and errors the stream. @@ -2160,9 +2333,24 @@ export class TriggerChatTransport implements ChatTransport { ); } - // Settled close, or the turn is gone — tell the UI instead of - // leaving it spinning on a stream nobody will finish. - if (state.isStreaming && this.activeStreams.get(chatId) === internalAbort) { + if ( + !this.watchMode && + state.skipSettledPeek && + state.isStreaming === undefined && + !currentSubscription?.sessionSettled && + !combinedSignal.aborted + ) { + throw emptyRecoveryError; + } + + // A passive abort closes this view, not the remote turn. Only a + // settled subscription changes the turn's persisted streaming state. + if ( + state.isStreaming !== false && + currentSubscription?.sessionSettled && + !combinedSignal.aborted && + this.activeStreams.get(chatId) === internalAbort + ) { state.isStreaming = false; this.notifySessionChange(chatId, state); } @@ -2295,8 +2483,8 @@ export class TriggerChatTransport implements ChatTransport { ) { continue; } - state.skipToTurnComplete = false; - state.supersededInputSeq = undefined; + this.clearStoppedBoundary(state); + this.notifySessionChange(chatId, state); // This boundary is the new turn's own, so the gate swallowed its // output: fail the turn instead of completing an empty answer, and // leave nothing armed for the retry. @@ -2392,6 +2580,7 @@ export class TriggerChatTransport implements ChatTransport { }); state.activeInputSeq = undefined; sinceInSeq = undefined; + state.skipSettledPeek = false; state.isStreaming = false; this.notifySessionChange(chatId, state); this.coordinator?.release(chatId); @@ -2416,6 +2605,12 @@ export class TriggerChatTransport implements ChatTransport { // unwrapped from the S2 record envelope (the parser does the // JSON unwrap). Drop empty/malformed payloads defensively. if (value.chunk == null) continue; + // A resumed session can contain only a token and output cursor. + // Its first data record establishes an active turn for Stop. + if (!state.outstandingTurnAbandoned && state.isStreaming !== true) { + state.isStreaming = true; + this.notifySessionChange(chatId, state); + } if (!sawFirstChunk) { sawFirstChunk = true; this.emitEvent({ @@ -2430,6 +2625,7 @@ export class TriggerChatTransport implements ChatTransport { controller.enqueue(value.chunk as UIMessageChunk); } } catch (error) { + internalAbort.abort(); if (error instanceof Error && error.name === "AbortError") { try { controller.close(); @@ -2440,7 +2636,8 @@ export class TriggerChatTransport implements ChatTransport { } const errorStatus = (error as { status?: unknown }).status; // A superseded stream cannot settle the replacement stream. - if (this.activeStreams.get(chatId) === internalAbort) { + // An empty recovery poll retains unknown state so reconnect can resume the accepted message. + if (this.activeStreams.get(chatId) === internalAbort && error !== emptyRecoveryError) { state.isStreaming = false; this.notifySessionChange(chatId, state); } diff --git a/packages/trigger-sdk/test/chat-transport-events.test.ts b/packages/trigger-sdk/test/chat-transport-events.test.ts index 2c4d24b85a0..99c29d3bb68 100644 --- a/packages/trigger-sdk/test/chat-transport-events.test.ts +++ b/packages/trigger-sdk/test/chat-transport-events.test.ts @@ -175,68 +175,30 @@ describe("transport send events", () => { }); describe("stopped turn followed by a new turn", () => { - /** - * `.out` stub that honours the `Last-Event-ID` cursor like the server does, so - * a resubscribe cannot replay records the reader already consumed. Legacy v1 - * frames carry no `session-in-event-id`, so the stopped turn's boundary is - * indistinguishable from this turn's: the tail is dropped and the turn closes. - */ - function cursoredTwoTurnTransport() { - const frames = [ - { id: "1", data: `{"type":"text-delta","id":"t1","delta":"stale"}` }, - { id: "2", data: `{"type":"trigger:turn-complete"}` }, - ]; - - return makeTransport({ - sessions: { c1: { publicAccessToken: "tok_test", isStreaming: true } }, - fetch: async (_url, init, ctx) => { - if (ctx.endpoint === "in") return jsonOk(); - - const cursor = new Headers(init.headers).get("Last-Event-ID"); - const from = cursor ? frames.findIndex((f) => f.id === cursor) + 1 : 0; - const remaining = frames.slice(from); - const response = sseResponse( - remaining.map((f) => `id: ${f.id}\ndata: ${f.data}\n\n`).join("") - ); - // Nothing left to send: the session is settled, so the reader stops - // instead of resubscribing. - if (remaining.length === 0) response.headers.set("X-Session-Settled", "true"); - return response; - }, - }); - } - - it("drops the stopped turn's tail and closes the sendMessages turn", async () => { - const { transport, events } = cursoredTwoTurnTransport(); - - expect(await transport.stopGeneration("c1")).toBe(true); - events.length = 0; - - const stream = await transport.sendMessages({ - trigger: "submit-message", - chatId: "c1", - messageId: undefined, - messages: [user("after stop", "u-2")], - abortSignal: undefined, - }); - const chunks = await readAll(stream); - - expect(chunks).toEqual([]); - expect(events.some((e) => e.type === "turn-completed")).toBe(true); - }); - - it("drops the stopped turn's tail and closes the sendAction turn", async () => { - const { transport, events } = cursoredTwoTurnTransport(); - - expect(await transport.stopGeneration("c1")).toBe(true); - events.length = 0; - - const stream = await transport.sendAction("c1", { type: "undo" }); - const chunks = await readAll(stream); - - expect(chunks).toEqual([]); - expect(events.some((e) => e.type === "turn-completed")).toBe(true); - }); + it.each(["message", "action"] as const)( + "does not report completion for an uncorrelated %s", + async (kind) => { + const { transport, events } = makeTransport({ + sessions: { c1: { publicAccessToken: "tok_test", isStreaming: true } }, + }); + expect(await transport.stopGeneration("c1")).toBe(true); + events.length = 0; + const sent = + kind === "action" + ? transport.sendAction("c1", { type: "undo" }) + : transport.sendMessages({ + trigger: "submit-message", + chatId: "c1", + messageId: undefined, + messages: [user("after stop", "u-2")], + abortSignal: undefined, + }); + await expect(sent).rejects.toThrow("Reload the chat before sending another message"); + expect(events.some((e) => e.type === "message-sent")).toBe(true); + expect(events.some((e) => e.type === "stream-connected")).toBe(false); + expect(events.some((e) => e.type === "turn-completed")).toBe(false); + } + ); }); describe("transport stream events", () => { diff --git a/packages/trigger-sdk/test/chat-turn-correlation.test.ts b/packages/trigger-sdk/test/chat-turn-correlation.test.ts index f81ca755a27..f37d638931f 100644 --- a/packages/trigger-sdk/test/chat-turn-correlation.test.ts +++ b/packages/trigger-sdk/test/chat-turn-correlation.test.ts @@ -120,6 +120,8 @@ describe("transport turn correlation", () => { lastEventId: undefined, activeInputSeq: 5, isStreaming: true, + outstandingTurnAbandoned: false, + skipSettledPeek: false, }); expect(transport.getSession("c1")?.activeInputSeq).toBe(5); await readDeltas(stream); diff --git a/packages/trigger-sdk/test/use-load-transcript-react.test.ts b/packages/trigger-sdk/test/use-load-transcript-react.test.ts new file mode 100644 index 00000000000..520473c59c2 --- /dev/null +++ b/packages/trigger-sdk/test/use-load-transcript-react.test.ts @@ -0,0 +1,230 @@ +// @vitest-environment jsdom + +import { act, createElement, StrictMode } from "react"; +import { createRoot, type Root } from "react-dom/client"; +import { afterAll, afterEach, beforeAll, describe, expect, it } from "vitest"; +import { useLoadTranscript, type LoadTranscriptResult } from "../src/v3/chat-react.js"; +import { TriggerChatTransport, type ChatSessionPersistedState } from "../src/v3/chat.js"; + +type ProbeProps = { + transport: TriggerChatTransport; + load: (params: { chatId: string; limit?: number }) => Promise; +}; + +function Probe({ transport, load }: ProbeProps) { + const result = useLoadTranscript("chat", load, { transport }); + return createElement( + "output", + { + "data-status": result.isLoading ? "loading" : result.error ? "error" : "ready", + "data-next-cursor": result.nextCursor, + }, + result.error?.message ?? result.messages.map((message) => message.id).join(",") + ); +} + +function pendingLoad() { + let resolve!: (result: LoadTranscriptResult) => void; + const promise = new Promise((settle) => { + resolve = settle; + }); + return { promise, resolve }; +} + +function transcript(messageId = "saved-answer", lastOutEventId = "11"): LoadTranscriptResult { + return { + messages: [ + { id: messageId, role: "assistant", parts: [{ type: "text", text: "Saved answer" }] }, + ], + cursors: { lastOutEventId, lastInEventId: "10" }, + nextCursor: "older-messages", + }; +} + +describe("useLoadTranscript mounted recovery", () => { + const mounted: Array<{ root: Root; container: HTMLDivElement; unmounted: boolean }> = []; + const transports: TriggerChatTransport[] = []; + const actEnvironment = Object.getOwnPropertyDescriptor(globalThis, "IS_REACT_ACT_ENVIRONMENT"); + + beforeAll(() => { + Object.defineProperty(globalThis, "IS_REACT_ACT_ENVIRONMENT", { + configurable: true, + writable: true, + value: true, + }); + }); + + afterEach(async () => { + for (const view of mounted.splice(0)) { + if (!view.unmounted) await act(async () => view.root.unmount()); + view.container.remove(); + } + for (const transport of transports.splice(0)) transport.dispose(); + }); + + afterAll(() => { + if (actEnvironment) { + Object.defineProperty(globalThis, "IS_REACT_ACT_ENVIRONMENT", actEnvironment); + } else { + Reflect.deleteProperty(globalThis, "IS_REACT_ACT_ENVIRONMENT"); + } + }); + + function blockedTransport() { + const saved: ChatSessionPersistedState[] = []; + const transport = new TriggerChatTransport({ + task: "chat-task", + accessToken: () => "test-token", + sessions: { + chat: { + publicAccessToken: "test-token", + lastEventId: "1", + supersededInputSeq: 10, + skipToTurnComplete: true, + requiresTranscriptReload: true, + isStreaming: false, + }, + }, + onSessionChange: (_chatId, session) => { + if (session) saved.push(session); + }, + }); + transports.push(transport); + return { transport, saved }; + } + + async function mount(props: ProbeProps, strict = false) { + const container = document.createElement("div"); + document.body.append(container); + const view = { root: createRoot(container), container, unmounted: false }; + mounted.push(view); + const render = async (next: ProbeProps) => { + await act(async () => { + const probe = createElement(Probe, next); + view.root.render(strict ? createElement(StrictMode, null, probe) : probe); + }); + }; + await render(props); + return { + container, + render, + async unmount() { + await act(async () => view.root.unmount()); + view.unmounted = true; + }, + }; + } + + it("renders the transcript and persists recovery after a valid load", async () => { + const { transport, saved } = blockedTransport(); + const pending = pendingLoad(); + const view = await mount({ transport, load: () => pending.promise }); + + expect(view.container.querySelector("output")?.dataset.status).toBe("loading"); + expect(transport.getSession("chat")?.requiresTranscriptReload).toBe(true); + + await act(async () => pending.resolve(transcript())); + + expect(view.container.querySelector("output")?.dataset.status).toBe("ready"); + expect(view.container.querySelector("output")?.dataset.nextCursor).toBe("older-messages"); + expect(view.container.textContent).toBe("saved-answer"); + expect(transport.getSession("chat")).toMatchObject({ + lastEventId: "11", + requiresTranscriptReload: false, + skipToTurnComplete: false, + }); + expect(saved).toHaveLength(1); + }); + + it("ignores the first Strict Mode load and recovers from the current load", async () => { + const { transport, saved } = blockedTransport(); + const loads: ReturnType[] = []; + const view = await mount( + { + transport, + load: () => { + const pending = pendingLoad(); + loads.push(pending); + return pending.promise; + }, + }, + true + ); + expect(loads).toHaveLength(2); + + await act(async () => loads[0]!.resolve(transcript("cancelled-answer", "99"))); + expect(view.container.querySelector("output")?.dataset.status).toBe("loading"); + expect(transport.getSession("chat")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + + await act(async () => loads[1]!.resolve(transcript("current-answer"))); + expect(view.container.textContent).toBe("current-answer"); + expect(transport.getSession("chat")?.lastEventId).toBe("11"); + expect(saved).toHaveLength(1); + }); + + it("does not recover after the component unmounts", async () => { + const { transport, saved } = blockedTransport(); + const pending = pendingLoad(); + const view = await mount({ transport, load: () => pending.promise }); + await view.unmount(); + + await act(async () => pending.resolve(transcript())); + + expect(view.container.textContent).toBe(""); + expect(transport.getSession("chat")).toMatchObject({ + lastEventId: "1", + requiresTranscriptReload: true, + }); + expect(saved).toEqual([]); + }); + + it("loads again for a replacement transport and ignores the previous result", async () => { + const first = blockedTransport(); + const replacement = blockedTransport(); + const loads: ReturnType[] = []; + const load = () => { + const pending = pendingLoad(); + loads.push(pending); + return pending.promise; + }; + const view = await mount({ transport: first.transport, load }); + await view.render({ transport: replacement.transport, load }); + expect(loads).toHaveLength(2); + + await act(async () => loads[0]!.resolve(transcript("previous-answer", "99"))); + expect(view.container.querySelector("output")?.dataset.status).toBe("loading"); + expect(first.saved).toEqual([]); + expect(replacement.saved).toEqual([]); + + await act(async () => loads[1]!.resolve(transcript("replacement-answer"))); + expect(view.container.textContent).toBe("replacement-answer"); + expect(first.transport.getSession("chat")?.requiresTranscriptReload).toBe(true); + expect(replacement.transport.getSession("chat")?.requiresTranscriptReload).toBe(false); + expect(replacement.saved).toHaveLength(1); + }); + + it("reports a stale checkpoint and retains the send block", async () => { + const { transport, saved } = blockedTransport(); + const loaded = transcript(); + loaded.cursors = { lastOutEventId: "11", lastInEventId: "9" }; + const view = await mount({ transport, load: async () => loaded }); + + expect(view.container.querySelector("output")?.dataset.status).toBe("error"); + expect(view.container.textContent).toContain("not current"); + expect(transport.getSession("chat")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + }); + + it("reports missing input evidence and retains the send block", async () => { + const { transport, saved } = blockedTransport(); + const loaded = transcript(); + loaded.cursors = { lastOutEventId: "11" }; + const view = await mount({ transport, load: async () => loaded }); + + expect(view.container.querySelector("output")?.dataset.status).toBe("error"); + expect(view.container.textContent).toMatch(/input/i); + expect(transport.getSession("chat")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + }); +}); diff --git a/packages/trigger-sdk/test/use-load-transcript.test.ts b/packages/trigger-sdk/test/use-load-transcript.test.ts index 911f7c7ad10..a1c85a331b1 100644 --- a/packages/trigger-sdk/test/use-load-transcript.test.ts +++ b/packages/trigger-sdk/test/use-load-transcript.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from "vitest"; -import { TriggerChatTransport } from "../src/v3/chat.js"; +import { TriggerChatTransport, type ChatSessionPersistedState } from "../src/v3/chat.js"; import { seedTranscriptCursor } from "../src/v3/chat-react.js"; function transportWithStart() { @@ -56,3 +56,197 @@ describe("seedTranscriptCursor + TriggerChatTransport resume cursor", () => { expect(transport.getSession("chat-1")?.lastEventId).toBe("50"); }); }); + +describe("transcript recovery", () => { + function blockedTransport(overrides: Partial = {}) { + const saved: ChatSessionPersistedState[] = []; + const transport = new TriggerChatTransport({ + task: "my-chat", + accessToken: () => "pat", + sessions: { + "chat-1": { + publicAccessToken: "pat", + lastEventId: "9007199254740992", + isStreaming: false, + skipToTurnComplete: true, + supersededInputSeq: 4, + requiresTranscriptReload: true, + ...overrides, + }, + }, + onSessionChange: (_chatId, session) => { + if (session) saved.push(session); + }, + }); + return { transport, saved }; + } + + it("installs a newer checkpoint and persists recovery once", () => { + const { transport, saved } = blockedTransport(); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect(recover).toBeDefined(); + expect(recover?.({ lastOutEventId: "9007199254740993", lastInEventId: "4" })).toBe(true); + expect(transport.getSession("chat-1")).toMatchObject({ + lastEventId: "9007199254740993", + requiresTranscriptReload: false, + skipToTurnComplete: false, + supersededInputSeq: undefined, + activeInputSeq: undefined, + outstandingTurnAbandoned: false, + skipSettledPeek: true, + isStreaming: undefined, + }); + expect(saved).toHaveLength(1); + expect(recover?.({ lastOutEventId: "9007199254740994", lastInEventId: "4" })).toBe(false); + expect(saved).toHaveLength(1); + }); + + it.each([undefined, "", "NaN", "-1", "1e20", "9007199254740993x"])( + "reports invalid output evidence and keeps sends blocked: %s", + (lastOutEventId) => { + const { transport, saved } = blockedTransport(); + const before = transport.getSession("chat-1"); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect(() => recover?.({ lastOutEventId, lastInEventId: "4" })).toThrow( + "Transcript recovery requires numeric input and output cursors" + ); + expect(transport.getSession("chat-1")).toEqual(before); + expect(saved).toEqual([]); + } + ); + + it.each(["9007199254740992", "42"])("rejects a stale output checkpoint: %s", (lastOutEventId) => { + const { transport, saved } = blockedTransport(); + expect( + transport.prepareTranscriptRecovery("chat-1")?.({ lastOutEventId, lastInEventId: "4" }) + ).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + }); + + it("reports missing cursor evidence", () => { + const { transport } = blockedTransport(); + expect(() => transport.prepareTranscriptRecovery("chat-1")?.(undefined)).toThrow( + "Transcript recovery requires numeric input and output cursors" + ); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it("rejects a malformed persisted cursor", () => { + const { transport } = blockedTransport({ lastEventId: "invalid" }); + expect( + transport.prepareTranscriptRecovery("chat-1")?.({ + lastOutEventId: "9007199254740993", + lastInEventId: "4", + }) + ).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it("recovers a session with no previous cursor", () => { + const { transport } = blockedTransport({ lastEventId: undefined }); + expect( + transport.prepareTranscriptRecovery("chat-1")?.({ lastOutEventId: "42", lastInEventId: "4" }) + ).toBe(true); + expect(transport.getSession("chat-1")?.lastEventId).toBe("42"); + }); + + it("does not seed a cursor after a superseded recovery fails", () => { + const { transport } = blockedTransport({ lastEventId: undefined }); + const stale = transport.prepareTranscriptRecovery("chat-1"); + const current = transport.prepareTranscriptRecovery("chat-1"); + expect(stale?.({ lastOutEventId: "50", lastInEventId: "4" })).toBe(false); + expect(transport.getSession("chat-1")?.lastEventId).toBeUndefined(); + expect(current?.({ lastOutEventId: "42", lastInEventId: "4" })).toBe(true); + }); + + it("rejects a load for a replaced session", () => { + const { transport } = blockedTransport(); + const recover = transport.prepareTranscriptRecovery("chat-1"); + transport.setSession("chat-1", transport.getSession("chat-1")!); + expect(recover?.({ lastOutEventId: "9007199254740993", lastInEventId: "4" })).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it("rejects a load after the cursor changes", () => { + const { transport } = blockedTransport({ lastEventId: undefined }); + const recover = transport.prepareTranscriptRecovery("chat-1"); + transport.seedResumeCursor("chat-1", "42"); + expect(recover?.({ lastOutEventId: "43", lastInEventId: "4" })).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it.each(["abandon", "dispose"])("rejects a load after %s", (operation) => { + const { transport } = blockedTransport(); + const recover = transport.prepareTranscriptRecovery("chat-1"); + if (operation === "abandon") transport.clearSupersedeGate("chat-1"); + else transport.dispose(); + expect(recover?.({ lastOutEventId: "9007199254740993", lastInEventId: "4" })).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it.each([undefined, "", "NaN", "-1", "4x"])( + "reports invalid input evidence and keeps sends blocked: %s", + (lastInEventId) => { + const { transport, saved } = blockedTransport({ lastEventId: undefined }); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect(() => recover?.({ lastOutEventId: "42", lastInEventId })).toThrow( + "Transcript recovery requires numeric input and output cursors" + ); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + } + ); + + it("rejects a stale input checkpoint", () => { + const { transport, saved } = blockedTransport(); + expect( + transport.prepareTranscriptRecovery("chat-1")?.({ + lastOutEventId: "9007199254740993", + lastInEventId: "3", + }) + ).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + }); + + it("rejects an obsolete callback before it examines missing evidence", () => { + const { transport } = blockedTransport(); + const stale = transport.prepareTranscriptRecovery("chat-1"); + transport.prepareTranscriptRecovery("chat-1"); + expect(stale?.(undefined)).toBe(false); + }); + + it("keeps recovery blocked when the stopped input is unknown", () => { + const { transport } = blockedTransport({ + lastEventId: undefined, + supersededInputSeq: undefined, + }); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect(() => recover?.({ lastOutEventId: "42", lastInEventId: "100" })).toThrow( + "Transcript recovery requires a stopped input sequence" + ); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it("does not prepare recovery for an ordinary load or a closed session", () => { + const ordinary = blockedTransport({ requiresTranscriptReload: false }).transport; + expect(ordinary.prepareTranscriptRecovery("chat-1")).toBeUndefined(); + const closed = blockedTransport({ closed: true }).transport; + expect(closed.prepareTranscriptRecovery("chat-1")).toBeUndefined(); + expect(closed.getSession("chat-1")?.closed).toBe(true); + expect(closed.prepareTranscriptRecovery("unknown")).toBeUndefined(); + }); + + it("preserves closed state on hydration but resets it on explicit replacement", () => { + const { transport } = blockedTransport({ closed: true, closedReason: "finished" }); + const session = transport.getSession("chat-1")!; + expect(session).toMatchObject({ closed: true, closedReason: "finished" }); + transport.setSession("chat-1", session); + expect(transport.getSession("chat-1")).toMatchObject({ + closed: undefined, + closedReason: undefined, + }); + expect(transport.prepareTranscriptRecovery("chat-1")).toBeDefined(); + }); +}); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 69bcc2bf1a7..4a4f65051b2 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -2114,9 +2114,6 @@ importers: '@trigger.dev/core': specifier: workspace:4.6.3 version: link:../core - react: - specifier: 18.3.1 - version: 18.3.1 uncrypto: specifier: ^0.1.3 version: 0.1.3 @@ -2130,12 +2127,24 @@ importers: '@types/react': specifier: ^19.2.14 version: 19.2.14 + '@types/react-dom': + specifier: 19.2.3 + version: 19.2.3(@types/react@19.2.14) ai: specifier: 6.0.116 version: 6.0.116(zod@4.5.4) ai-v7: specifier: npm:ai@7.0.0-canary.159 version: ai@7.0.0-canary.159(zod@4.5.4) + jsdom: + specifier: 30.0.1 + version: 30.0.1 + react: + specifier: 18.3.1 + version: 18.3.1 + react-dom: + specifier: 18.3.1 + version: 18.3.1(react@18.3.1) rimraf: specifier: ^6.0.1 version: 6.0.1 @@ -8063,6 +8072,11 @@ packages: '@types/react-dom@18.2.7': resolution: {integrity: sha512-GRaAEriuT4zp9N4p1i8BDBYmEyfo+xQ3yHjJU4eiK5NDa1RmUZG+unZABUTK4/Ox/M+GaHwb6Ow8rUITrtjszA==} + '@types/react-dom@19.2.3': + resolution: {integrity: sha512-jp2L/eY6fn+KgVVQAOqYItbF0VY/YApe5Mz2F0aykSO8gx31bYCZyvSeYxCHKvzHG5eZjc+zyaS5BrBWya2+kQ==} + peerDependencies: + '@types/react': ^19.2.0 + '@types/react@18.2.69': resolution: {integrity: sha512-W1HOMUWY/1Yyw0ba5TkCV+oqynRjG7BnteBB+B7JmAK7iw3l2SW+VGOxL+akPweix6jk2NNJtyJKpn4TkpfK3Q==} @@ -22603,6 +22617,10 @@ snapshots: dependencies: '@types/react': 18.2.69 + '@types/react-dom@19.2.3(@types/react@19.2.14)': + dependencies: + '@types/react': 19.2.14 + '@types/react@18.2.69': dependencies: '@types/prop-types': 15.7.5