From 4ca683dc74dcb5552cf27bb0d4f4f077802f7846 Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Fri, 18 Sep 2026 18:29:36 -0700 Subject: [PATCH 1/5] fix(sdk): preserve stopped chat boundaries across successor sends --- .changeset/chat-stop-successor-boundary.md | 6 + packages/trigger-sdk/src/v3/chat-stop.test.ts | 402 ++++++++++++++++++ packages/trigger-sdk/src/v3/chat.test.ts | 14 +- packages/trigger-sdk/src/v3/chat.ts | 117 +++-- .../test/chat-transport-events.test.ts | 86 ++-- 5 files changed, 535 insertions(+), 90 deletions(-) create mode 100644 .changeset/chat-stop-successor-boundary.md create mode 100644 packages/trigger-sdk/src/v3/chat-stop.test.ts diff --git a/.changeset/chat-stop-successor-boundary.md b/.changeset/chat-stop-successor-boundary.md new file mode 100644 index 00000000000..47fbc8e1e4f --- /dev/null +++ b/.changeset/chat-stop-successor-boundary.md @@ -0,0 +1,6 @@ +--- +"@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. 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..861de1f09f5 --- /dev/null +++ b/packages/trigger-sdk/src/v3/chat-stop.test.ts @@ -0,0 +1,402 @@ +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 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; +} + +describe("Stop with a successor response", () => { + let server: Server; + let baseURL: string; + let transport: TriggerChatTransport; + let outputs: ServerResponse[]; + let inputSeq: number; + let holdStop: boolean; + let stopStatus: number; + let settled: 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 = []; + inputSeq = 10; + holdStop = false; + stopStatus = 200; + settled = 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; + } + response.writeHead(200, { + "Content-Type": "text/event-stream", + "X-Stream-Version": "v2", + "X-Session-Settled": String(settled), + }); + response.flushHeaders(); + outputs.push(response); + }); + 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), + ]; + } + + 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(["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("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("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..d75d8d7a547 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1577,13 +1577,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 +1591,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..3c16665d530 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -496,6 +496,12 @@ 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; /** 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. */ @@ -713,6 +719,7 @@ type ChatSessionState = { skipToTurnComplete?: boolean; /** `.in` seq of the turn the gate supersedes; only its boundary (or a later one) clears the gate. */ supersededInputSeq?: number; + requiresTranscriptReload?: boolean; /** 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. */ @@ -807,6 +814,9 @@ export class TriggerChatTransport implements ChatTransport { lastEventId: session.lastEventId, activeInputSeq: session.activeInputSeq, isStreaming: session.isStreaming, + skipToTurnComplete: session.skipToTurnComplete, + supersededInputSeq: session.supersededInputSeq, + requiresTranscriptReload: session.requiresTranscriptReload, closed: session.closed, closedReason: session.closedReason, }); @@ -951,6 +961,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,6 +986,7 @@ export class TriggerChatTransport implements ChatTransport { } state.activeInputSeq = inSeq; + this.requireStoppedTurnCorrelation(chatId, state, inSeq); state.isStreaming = true; state.outstandingTurnAbandoned = false; this.notifySessionChange(chatId, state); @@ -1280,6 +1292,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. @@ -1315,24 +1328,8 @@ 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; - } - + // 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 ( @@ -1360,7 +1357,24 @@ 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 = 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) + ); + return true; + } catch { + // The reader already closed. Retain its unread boundary for the next send. + return false; + } }; /** @@ -1373,7 +1387,10 @@ export class TriggerChatTransport implements ChatTransport { if (!state) return; state.skipToTurnComplete = false; state.supersededInputSeq = undefined; + state.activeInputSeq = undefined; state.outstandingTurnAbandoned = true; + state.isStreaming = false; + this.notifySessionChange(chatId, state); }; /** @@ -1408,6 +1425,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,6 +1450,7 @@ 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; this.notifySessionChange(chatId, state); @@ -1461,6 +1480,9 @@ export class TriggerChatTransport implements ChatTransport { lastEventId: session.lastEventId, activeInputSeq: session.activeInputSeq, isStreaming: session.isStreaming, + skipToTurnComplete: session.skipToTurnComplete, + supersededInputSeq: session.supersededInputSeq, + requiresTranscriptReload: session.requiresTranscriptReload, }) ); this.notifySessionChange(chatId, this.toPersisted(this.sessions.get(chatId)!)); @@ -1645,11 +1667,35 @@ export class TriggerChatTransport implements ChatTransport { 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.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, closed: state.closed, closedReason: state.closedReason, }); @@ -1961,7 +2007,11 @@ export class TriggerChatTransport implements ChatTransport { } state.publicAccessToken = publicAccessToken; state.lastEventId = undefined; + state.activeInputSeq = undefined; state.isStreaming = false; + state.skipToTurnComplete = false; + state.supersededInputSeq = undefined; + state.requiresTranscriptReload = false; this.sessions.set(chatId, state); this.notifySessionChange(chatId, state); } @@ -2002,9 +2052,16 @@ export class TriggerChatTransport implements ChatTransport { const outstanding = !state.outstandingTurnAbandoned && (state.isStreaming || state.activeInputSeq !== undefined); - if (options?.sendStopOnAbort !== false && outstanding && !internalAbort.signal.aborted) { + if ( + options?.sendStopOnAbort !== false && + outstanding && + !internalAbort.signal.aborted && + this.activeStreams.get(chatId) === internalAbort + ) { state.skipToTurnComplete = true; state.supersededInputSeq = state.activeInputSeq; + state.isStreaming = false; + this.notifySessionChange(chatId, state); this.appendInputChunk( chatId, state.publicAccessToken, @@ -2146,7 +2203,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 +2217,13 @@ 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) { + // A passive abort closes this view, not the remote turn. Only a + // settled subscription changes the turn's persisted streaming state. + if ( + state.isStreaming && + !combinedSignal.aborted && + this.activeStreams.get(chatId) === internalAbort + ) { state.isStreaming = false; this.notifySessionChange(chatId, state); } @@ -2297,6 +2358,7 @@ export class TriggerChatTransport implements ChatTransport { } state.skipToTurnComplete = false; state.supersededInputSeq = undefined; + 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. @@ -2416,6 +2478,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 +2498,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(); 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", () => { From 92b6a916f0b2544a08e72d9c2de9f4a7ff66e57a Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Fri, 18 Sep 2026 18:58:23 -0700 Subject: [PATCH 2/5] fix(sdk): recover stopped sessions after transcript reload --- .changeset/chat-stop-successor-boundary.md | 1 + packages/trigger-sdk/src/v3/chat-react.ts | 22 +- packages/trigger-sdk/src/v3/chat-stop.test.ts | 298 +++++++++++++++++- packages/trigger-sdk/src/v3/chat.test.ts | 5 +- packages/trigger-sdk/src/v3/chat.ts | 93 +++++- .../test/chat-turn-correlation.test.ts | 2 + .../test/use-load-transcript.test.ts | 116 ++++++- 7 files changed, 521 insertions(+), 16 deletions(-) diff --git a/.changeset/chat-stop-successor-boundary.md b/.changeset/chat-stop-successor-boundary.md index 47fbc8e1e4f..e29c95ce622 100644 --- a/.changeset/chat-stop-successor-boundary.md +++ b/.changeset/chat-stop-successor-boundary.md @@ -4,3 +4,4 @@ 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 without replaying old output. diff --git a/packages/trigger-sdk/src/v3/chat-react.ts b/packages/trigger-sdk/src/v3/chat-react.ts index f9fdaa47a19..9de0eca1f97 100644 --- a/packages/trigger-sdk/src/v3/chat-react.ts +++ b/packages/trigger-sdk/src/v3/chat-react.ts @@ -65,6 +65,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`). + * For a blocked session, a fresh load with a newer saved cursor also clears + * the transcript reload requirement. A stale result leaves sends blocked. */ transport?: TriggerChatTransport; /** Page size passed to the action. */ @@ -77,14 +79,17 @@ 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. A captured recovery callback replaces ordinary cursor seeding. + * Returns whether the cursor was accepted. */ export function seedTranscriptCursor( transport: Pick, chatId: string, - cursors: { lastOutEventId?: string } | undefined + cursors: { lastOutEventId?: string } | undefined, + completeRecovery?: (lastEventId: string | undefined) => boolean ): boolean { const lastEventId = cursors?.lastOutEventId; + if (completeRecovery) return completeRecovery(lastEventId); if (!lastEventId) return false; transport.seedResumeCursor(chatId, lastEventId); return true; @@ -130,8 +135,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 +144,17 @@ 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 (transport) { + const seeded = seedTranscriptCursor(transport, chatId, result.cursors, completeRecovery); + if (completeRecovery && !seeded) { + throw new Error("The loaded transcript is not current. Reload the chat again."); + } } 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 index 861de1f09f5..0c810d98412 100644 --- a/packages/trigger-sdk/src/v3/chat-stop.test.ts +++ b/packages/trigger-sdk/src/v3/chat-stop.test.ts @@ -56,15 +56,31 @@ async function readText(stream: ReadableStream): Promise 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; @@ -93,10 +109,13 @@ describe("Stop with a successor response", () => { beforeEach(async () => { outputs = []; + outputHeaders = []; inputSeq = 10; holdStop = false; stopStatus = 200; settled = false; + resumeAfterStoppedCheckpoint = false; + emptyRecoveredOutput = false; includeSequence = true; pendingStop = undefined; saved = null; @@ -115,13 +134,26 @@ describe("Stop with a successor response", () => { } 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), + "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(); @@ -164,6 +196,25 @@ describe("Stop with a successor response", () => { ]; } + 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) => { @@ -212,6 +263,81 @@ describe("Stop with a successor 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) => { @@ -382,6 +508,176 @@ describe("Stop with a successor response", () => { 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("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"); + expect(recover(cursor)).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("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("retains unknown recovery state after a bounded empty response", async () => { + await hydrateBlockedSession("constructor"); + const recover = transport.prepareTranscriptRecovery("chat"); + if (!recover) throw new Error("Expected transcript recovery"); + expect(recover("11")).toBe(true); + resumeAfterStoppedCheckpoint = 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(""); + 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("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(["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"); diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index d75d8d7a547..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"); diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 3c16665d530..f052a60dd4a 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -502,6 +502,10 @@ export type ChatSessionPersistedState = { supersededInputSeq?: number; /** A send lacks an input sequence while stopped output remains unread. Reload the transcript before another send. */ requiresTranscriptReload?: boolean; + /** 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. */ @@ -724,6 +728,9 @@ type ChatSessionState = { isStreaming?: boolean; /** Set once the outstanding turn is declared dead: a later stop must not gate the next turn on it. */ outstandingTurnAbandoned?: boolean; + /** Identifies the latest transcript load for this blocked session. Never persisted. */ + transcriptRecovery?: symbol; + skipSettledPeek?: boolean; /** Set once the session is closed. Terminal — sends and reconnects stop. */ closed?: boolean; /** The reason the session was closed, when one was given. */ @@ -817,6 +824,8 @@ export class TriggerChatTransport implements ChatTransport { skipToTurnComplete: session.skipToTurnComplete, supersededInputSeq: session.supersededInputSeq, requiresTranscriptReload: session.requiresTranscriptReload, + outstandingTurnAbandoned: session.outstandingTurnAbandoned, + skipSettledPeek: session.skipSettledPeek, closed: session.closed, closedReason: session.closedReason, }); @@ -989,6 +998,7 @@ export class TriggerChatTransport implements ChatTransport { 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. @@ -1237,7 +1247,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 @@ -1316,7 +1326,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, }); }; @@ -1328,13 +1338,14 @@ export class TriggerChatTransport implements ChatTransport { stopGeneration = async (chatId: string): Promise => { const state = this.sessions.get(chatId); if (!state) 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; @@ -1385,10 +1396,12 @@ export class TriggerChatTransport implements ChatTransport { clearSupersedeGate = (chatId: string): void => { const state = this.sessions.get(chatId); if (!state) return; + state.transcriptRecovery = undefined; state.skipToTurnComplete = false; state.supersededInputSeq = undefined; state.activeInputSeq = undefined; state.outstandingTurnAbandoned = true; + state.skipSettledPeek = false; state.isStreaming = false; this.notifySessionChange(chatId, state); }; @@ -1453,6 +1466,7 @@ export class TriggerChatTransport implements ChatTransport { 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. @@ -1483,6 +1497,8 @@ export class TriggerChatTransport implements ChatTransport { skipToTurnComplete: session.skipToTurnComplete, supersededInputSeq: session.supersededInputSeq, requiresTranscriptReload: session.requiresTranscriptReload, + outstandingTurnAbandoned: session.outstandingTurnAbandoned, + skipSettledPeek: session.skipSettledPeek, }) ); this.notifySessionChange(chatId, this.toPersisted(this.sessions.get(chatId)!)); @@ -1509,6 +1525,54 @@ 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, while the blocked session remains unchanged. + * Ordinary transcript loads do not reset session state. + */ + prepareTranscriptRecovery = ( + chatId: string + ): ((lastEventId: string | 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 } = state; + return (loadedEventId) => { + 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 || + loadedEventId === undefined || + !/^\d+$/.test(loadedEventId) || + (lastEventId !== undefined && + (!/^\d+$/.test(lastEventId) || BigInt(loadedEventId) <= BigInt(lastEventId))) + ) { + return false; + } + + state.lastEventId = loadedEventId; + state.requiresTranscriptReload = false; + state.skipToTurnComplete = false; + state.supersededInputSeq = undefined; + 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) { @@ -1655,6 +1719,9 @@ export class TriggerChatTransport implements ChatTransport { controller.abort(); } this.activeStreams.clear(); + for (const state of this.sessions.values()) { + state.transcriptRecovery = undefined; + } this.coordinator?.dispose(); this.coordinator = null; } @@ -1683,6 +1750,7 @@ export class TriggerChatTransport implements ChatTransport { 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); @@ -1696,6 +1764,8 @@ export class TriggerChatTransport implements ChatTransport { skipToTurnComplete: state.skipToTurnComplete, supersededInputSeq: state.supersededInputSeq, requiresTranscriptReload: state.requiresTranscriptReload, + outstandingTurnAbandoned: state.outstandingTurnAbandoned, + skipSettledPeek: state.skipSettledPeek, closed: state.closed, closedReason: state.closedReason, }); @@ -1713,6 +1783,7 @@ 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; @@ -2012,6 +2083,9 @@ export class TriggerChatTransport implements ChatTransport { state.skipToTurnComplete = false; state.supersededInputSeq = undefined; state.requiresTranscriptReload = false; + state.outstandingTurnAbandoned = false; + state.skipSettledPeek = false; + state.transcriptRecovery = undefined; this.sessions.set(chatId, state); this.notifySessionChange(chatId, state); } @@ -2051,7 +2125,7 @@ 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); + (state.isStreaming !== false || state.activeInputSeq !== undefined); if ( options?.sendStopOnAbort !== false && outstanding && @@ -2145,7 +2219,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, @@ -2220,7 +2299,8 @@ export class TriggerChatTransport implements ChatTransport { // A passive abort closes this view, not the remote turn. Only a // settled subscription changes the turn's persisted streaming state. if ( - state.isStreaming && + state.isStreaming !== false && + currentSubscription?.sessionSettled && !combinedSignal.aborted && this.activeStreams.get(chatId) === internalAbort ) { @@ -2454,6 +2534,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); 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.test.ts b/packages/trigger-sdk/test/use-load-transcript.test.ts index 911f7c7ad10..c31b56f0f3c 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,117 @@ 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( + seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "9007199254740993" }, recover) + ).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?.("9007199254740994")).toBe(false); + expect(saved).toHaveLength(1); + }); + + it.each([undefined, "", "NaN", "-1", "1e20", "9007199254740993x", "9007199254740992", "42"])( + "keeps sends blocked for an invalid or stale checkpoint: %s", + (lastOutEventId) => { + const { transport, saved } = blockedTransport(); + const before = transport.getSession("chat-1"); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect(seedTranscriptCursor(transport, "chat-1", { lastOutEventId }, recover)).toBe(false); + expect(transport.getSession("chat-1")).toEqual(before); + expect(saved).toEqual([]); + } + ); + + it("rejects a malformed persisted cursor", () => { + const { transport } = blockedTransport({ lastEventId: "invalid" }); + expect(transport.prepareTranscriptRecovery("chat-1")?.("9007199254740993")).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")?.("42")).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(seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "50" }, stale)).toBe(false); + expect(transport.getSession("chat-1")?.lastEventId).toBeUndefined(); + expect(seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "42" }, current)).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?.("9007199254740993")).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?.("43")).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?.("9007199254740993")).toBe(false); + 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(); + }); +}); From ad4e2efc2e2103b111a9b09ec13011e9036bc430 Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Fri, 18 Sep 2026 19:48:12 -0700 Subject: [PATCH 3/5] fix(sdk): require input evidence for stopped transcript recovery --- .changeset/chat-stop-successor-boundary.md | 2 +- packages/trigger-sdk/src/v3/chat-react.ts | 10 +- packages/trigger-sdk/src/v3/chat-stop.test.ts | 110 +++++++++++++++++- packages/trigger-sdk/src/v3/chat.ts | 109 +++++++++++++---- .../test/use-load-transcript.test.ts | 66 +++++++++-- 5 files changed, 252 insertions(+), 45 deletions(-) diff --git a/.changeset/chat-stop-successor-boundary.md b/.changeset/chat-stop-successor-boundary.md index e29c95ce622..6fdb4bcd1b5 100644 --- a/.changeset/chat-stop-successor-boundary.md +++ b/.changeset/chat-stop-successor-boundary.md @@ -4,4 +4,4 @@ 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 without replaying old output. +Loading a fresh transcript through `useLoadTranscript` restores blocked sessions only after its saved input cursor covers the stopped turn. diff --git a/packages/trigger-sdk/src/v3/chat-react.ts b/packages/trigger-sdk/src/v3/chat-react.ts index 9de0eca1f97..acc844fa174 100644 --- a/packages/trigger-sdk/src/v3/chat-react.ts +++ b/packages/trigger-sdk/src/v3/chat-react.ts @@ -65,8 +65,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`). - * For a blocked session, a fresh load with a newer saved cursor also clears - * the transcript reload requirement. A stale result leaves sends blocked. + * 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. */ @@ -85,11 +85,11 @@ export type UseLoadTranscriptOptions = { export function seedTranscriptCursor( transport: Pick, chatId: string, - cursors: { lastOutEventId?: string } | undefined, - completeRecovery?: (lastEventId: string | undefined) => boolean + cursors: { lastOutEventId?: string; lastInEventId?: string } | undefined, + completeRecovery?: (lastEventId: string | undefined, lastInEventId?: string) => boolean ): boolean { const lastEventId = cursors?.lastOutEventId; - if (completeRecovery) return completeRecovery(lastEventId); + if (completeRecovery) return completeRecovery(lastEventId, cursors?.lastInEventId); if (!lastEventId) return false; transport.seedResumeCursor(chatId, lastEventId); return true; diff --git a/packages/trigger-sdk/src/v3/chat-stop.test.ts b/packages/trigger-sdk/src/v3/chat-stop.test.ts index 0c810d98412..4339524e0ea 100644 --- a/packages/trigger-sdk/src/v3/chat-stop.test.ts +++ b/packages/trigger-sdk/src/v3/chat-stop.test.ts @@ -501,6 +501,106 @@ describe("Stop with a successor 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("5", "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("5", "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("11", "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("17", "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("17", "11")).toBe(false); + await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); + }); + it("accepts a sequence-free response without a stopped boundary", async () => { includeSequence = false; const stream = await send(); @@ -532,7 +632,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession(hydrate); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("11")).toBe(true); + expect(recover("11", "11")).toBe(true); transport.seedResumeCursor("chat", "11"); expect(saved).toMatchObject({ lastEventId: "11", @@ -558,7 +658,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession("constructor"); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover(cursor)).toBe(false); + expect(recover(cursor, "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" @@ -585,7 +685,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession("constructor"); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("11")).toBe(true); + expect(recover("11", "11")).toBe(true); if (hydrate !== "none") { const session = transport.getSession("chat"); if (!session) throw new Error("Expected persisted state"); @@ -611,7 +711,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession("constructor"); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("11")).toBe(true); + expect(recover("11", "11")).toBe(true); resumeAfterStoppedCheckpoint = true; emptyRecoveredOutput = true; const resumed = await transport.reconnectToStream({ chatId: "chat" }); @@ -642,7 +742,7 @@ describe("Stop with a successor response", () => { holdStop = true; const stopped = transport.stopGeneration("chat"); await vi.waitFor(() => expect(pendingStop).toBeDefined()); - expect(recover("11")).toBe(false); + expect(recover("11", "11")).toBe(false); await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); const pending = pendingStop!; appendResponse(pending.response, pending.seq); diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index f052a60dd4a..76380c98e4f 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -502,6 +502,8 @@ export type ChatSessionPersistedState = { 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. */ @@ -724,6 +726,9 @@ type ChatSessionState = { /** `.in` seq of the turn the gate supersedes; only its boundary (or a later one) clears the gate. */ supersededInputSeq?: number; requiresTranscriptReload?: boolean; + transcriptRecoveryInputSeq?: number; + /** Identifies the stopped boundary for pending Stop acknowledgments. Never persisted. */ + stoppedBoundary?: symbol; /** 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. */ @@ -824,6 +829,7 @@ export class TriggerChatTransport implements ChatTransport { skipToTurnComplete: session.skipToTurnComplete, supersededInputSeq: session.supersededInputSeq, requiresTranscriptReload: session.requiresTranscriptReload, + transcriptRecoveryInputSeq: session.transcriptRecoveryInputSeq, outstandingTurnAbandoned: session.outstandingTurnAbandoned, skipSettledPeek: session.skipSettledPeek, closed: session.closed, @@ -1347,9 +1353,11 @@ export class TriggerChatTransport implements ChatTransport { !state.outstandingTurnAbandoned && (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) { @@ -1371,16 +1379,15 @@ export class TriggerChatTransport implements ChatTransport { const partId = crypto.randomUUID(); const serializedBody = this.serializeInputChunk({ kind: "stop" }); - const send = async (token: string) => { - await this.appendInputChunk(chatId, token, serializedBody, partId); - }; + const send = (token: string) => this.appendInputChunk(chatId, token, serializedBody, partId); try { - await this.sendWithEvents( + 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. @@ -1396,9 +1403,7 @@ export class TriggerChatTransport implements ChatTransport { clearSupersedeGate = (chatId: string): void => { const state = this.sessions.get(chatId); if (!state) return; - state.transcriptRecovery = undefined; - state.skipToTurnComplete = false; - state.supersededInputSeq = undefined; + this.clearStoppedBoundary(state); state.activeInputSeq = undefined; state.outstandingTurnAbandoned = true; state.skipSettledPeek = false; @@ -1497,6 +1502,7 @@ export class TriggerChatTransport implements ChatTransport { skipToTurnComplete: session.skipToTurnComplete, supersededInputSeq: session.supersededInputSeq, requiresTranscriptReload: session.requiresTranscriptReload, + transcriptRecoveryInputSeq: session.transcriptRecoveryInputSeq, outstandingTurnAbandoned: session.outstandingTurnAbandoned, skipSettledPeek: session.skipSettledPeek, }) @@ -1527,19 +1533,26 @@ export class TriggerChatTransport implements ChatTransport { /** * Capture recovery before a fresh transcript load starts. The returned callback - * accepts a newer saved cursor once, while the blocked session remains unchanged. + * accepts a newer saved cursor once, with input evidence for the stopped boundary. * Ordinary transcript loads do not reset session state. */ prepareTranscriptRecovery = ( chatId: string - ): ((lastEventId: string | undefined) => boolean) | undefined => { + ): ((lastEventId: string | undefined, lastInEventId?: string) => 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 } = state; - return (loadedEventId) => { + const { + lastEventId, + activeInputSeq, + skipToTurnComplete, + supersededInputSeq, + transcriptRecoveryInputSeq, + } = state; + const stoppedInputSeq = supersededInputSeq ?? transcriptRecoveryInputSeq; + return (loadedEventId, loadedInEventId) => { if (state.transcriptRecovery !== token) return false; state.transcriptRecovery = undefined; if ( @@ -1551,6 +1564,13 @@ export class TriggerChatTransport implements ChatTransport { state.activeInputSeq !== activeInputSeq || state.skipToTurnComplete !== skipToTurnComplete || state.supersededInputSeq !== supersededInputSeq || + state.transcriptRecoveryInputSeq !== transcriptRecoveryInputSeq || + stoppedInputSeq === undefined || + !Number.isSafeInteger(stoppedInputSeq) || + stoppedInputSeq < 0 || + loadedInEventId === undefined || + !/^\d+$/.test(loadedInEventId) || + BigInt(loadedInEventId) < BigInt(stoppedInputSeq) || loadedEventId === undefined || !/^\d+$/.test(loadedEventId) || (lastEventId !== undefined && @@ -1561,8 +1581,7 @@ export class TriggerChatTransport implements ChatTransport { state.lastEventId = loadedEventId; state.requiresTranscriptReload = false; - state.skipToTurnComplete = false; - state.supersededInputSeq = undefined; + this.clearStoppedBoundary(state); state.activeInputSeq = undefined; state.outstandingTurnAbandoned = false; state.skipSettledPeek = true; @@ -1721,6 +1740,7 @@ export class TriggerChatTransport implements ChatTransport { this.activeStreams.clear(); for (const state of this.sessions.values()) { state.transcriptRecovery = undefined; + state.stoppedBoundary = undefined; } this.coordinator?.dispose(); this.coordinator = null; @@ -1730,6 +1750,48 @@ 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); } @@ -1764,6 +1826,7 @@ export class TriggerChatTransport implements ChatTransport { skipToTurnComplete: state.skipToTurnComplete, supersededInputSeq: state.supersededInputSeq, requiresTranscriptReload: state.requiresTranscriptReload, + transcriptRecoveryInputSeq: state.transcriptRecoveryInputSeq, outstandingTurnAbandoned: state.outstandingTurnAbandoned, skipSettledPeek: state.skipSettledPeek, closed: state.closed, @@ -1787,8 +1850,7 @@ export class TriggerChatTransport implements ChatTransport { 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", @@ -2080,8 +2142,7 @@ export class TriggerChatTransport implements ChatTransport { state.lastEventId = undefined; state.activeInputSeq = undefined; state.isStreaming = false; - state.skipToTurnComplete = false; - state.supersededInputSeq = undefined; + this.clearStoppedBoundary(state); state.requiresTranscriptReload = false; state.outstandingTurnAbandoned = false; state.skipSettledPeek = false; @@ -2132,15 +2193,16 @@ export class TriggerChatTransport implements ChatTransport { !internalAbort.signal.aborted && this.activeStreams.get(chatId) === internalAbort ) { - state.skipToTurnComplete = true; - state.supersededInputSeq = state.activeInputSeq; + 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(); }, @@ -2436,8 +2498,7 @@ 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 diff --git a/packages/trigger-sdk/test/use-load-transcript.test.ts b/packages/trigger-sdk/test/use-load-transcript.test.ts index c31b56f0f3c..f00f0e672c9 100644 --- a/packages/trigger-sdk/test/use-load-transcript.test.ts +++ b/packages/trigger-sdk/test/use-load-transcript.test.ts @@ -86,7 +86,12 @@ describe("transcript recovery", () => { const recover = transport.prepareTranscriptRecovery("chat-1"); expect(recover).toBeDefined(); expect( - seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "9007199254740993" }, recover) + seedTranscriptCursor( + transport, + "chat-1", + { lastOutEventId: "9007199254740993", lastInEventId: "4" }, + recover + ) ).toBe(true); expect(transport.getSession("chat-1")).toMatchObject({ lastEventId: "9007199254740993", @@ -99,7 +104,7 @@ describe("transcript recovery", () => { isStreaming: undefined, }); expect(saved).toHaveLength(1); - expect(recover?.("9007199254740994")).toBe(false); + expect(recover?.("9007199254740994", "4")).toBe(false); expect(saved).toHaveLength(1); }); @@ -109,7 +114,9 @@ describe("transcript recovery", () => { const { transport, saved } = blockedTransport(); const before = transport.getSession("chat-1"); const recover = transport.prepareTranscriptRecovery("chat-1"); - expect(seedTranscriptCursor(transport, "chat-1", { lastOutEventId }, recover)).toBe(false); + expect( + seedTranscriptCursor(transport, "chat-1", { lastOutEventId, lastInEventId: "4" }, recover) + ).toBe(false); expect(transport.getSession("chat-1")).toEqual(before); expect(saved).toEqual([]); } @@ -117,13 +124,13 @@ describe("transcript recovery", () => { it("rejects a malformed persisted cursor", () => { const { transport } = blockedTransport({ lastEventId: "invalid" }); - expect(transport.prepareTranscriptRecovery("chat-1")?.("9007199254740993")).toBe(false); + expect(transport.prepareTranscriptRecovery("chat-1")?.("9007199254740993", "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")?.("42")).toBe(true); + expect(transport.prepareTranscriptRecovery("chat-1")?.("42", "4")).toBe(true); expect(transport.getSession("chat-1")?.lastEventId).toBe("42"); }); @@ -131,16 +138,25 @@ describe("transcript recovery", () => { const { transport } = blockedTransport({ lastEventId: undefined }); const stale = transport.prepareTranscriptRecovery("chat-1"); const current = transport.prepareTranscriptRecovery("chat-1"); - expect(seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "50" }, stale)).toBe(false); + expect( + seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "50", lastInEventId: "4" }, stale) + ).toBe(false); expect(transport.getSession("chat-1")?.lastEventId).toBeUndefined(); - expect(seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "42" }, current)).toBe(true); + expect( + seedTranscriptCursor( + transport, + "chat-1", + { lastOutEventId: "42", lastInEventId: "4" }, + current + ) + ).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?.("9007199254740993")).toBe(false); + expect(recover?.("9007199254740993", "4")).toBe(false); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); @@ -148,7 +164,7 @@ describe("transcript recovery", () => { const { transport } = blockedTransport({ lastEventId: undefined }); const recover = transport.prepareTranscriptRecovery("chat-1"); transport.seedResumeCursor("chat-1", "42"); - expect(recover?.("43")).toBe(false); + expect(recover?.("43", "4")).toBe(false); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); @@ -157,7 +173,37 @@ describe("transcript recovery", () => { const recover = transport.prepareTranscriptRecovery("chat-1"); if (operation === "abandon") transport.clearSupersedeGate("chat-1"); else transport.dispose(); - expect(recover?.("9007199254740993")).toBe(false); + expect(recover?.("9007199254740993", "4")).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + }); + + it.each([undefined, "", "NaN", "-1", "4x", "3"])( + "rejects a newer output checkpoint without stopped-input evidence: %s", + (lastInEventId) => { + const { transport, saved } = blockedTransport({ lastEventId: undefined }); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect( + seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "42", lastInEventId }, recover) + ).toBe(false); + expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); + expect(saved).toEqual([]); + } + ); + + it("keeps recovery blocked when the stopped input is unknown", () => { + const { transport } = blockedTransport({ + lastEventId: undefined, + supersededInputSeq: undefined, + }); + const recover = transport.prepareTranscriptRecovery("chat-1"); + expect( + seedTranscriptCursor( + transport, + "chat-1", + { lastOutEventId: "42", lastInEventId: "100" }, + recover + ) + ).toBe(false); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); From 75ec5792beb2d05eb0ee48605b94a700c67591e7 Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Fri, 18 Sep 2026 21:44:33 -0700 Subject: [PATCH 4/5] fix(sdk): simplify chat transcript recovery --- .changeset/chat-stop-successor-boundary.md | 1 + knip.json | 2 +- packages/trigger-sdk/package.json | 4 + packages/trigger-sdk/src/v3/chat-react.ts | 16 +- packages/trigger-sdk/src/v3/chat-stop.test.ts | 157 +++++++++++- packages/trigger-sdk/src/v3/chat.ts | 114 ++++----- .../test/use-load-transcript-react.test.ts | 230 ++++++++++++++++++ .../test/use-load-transcript.test.ts | 120 +++++---- pnpm-lock.yaml | 24 +- 9 files changed, 537 insertions(+), 131 deletions(-) create mode 100644 packages/trigger-sdk/test/use-load-transcript-react.test.ts diff --git a/.changeset/chat-stop-successor-boundary.md b/.changeset/chat-stop-successor-boundary.md index 6fdb4bcd1b5..f356ac582b0 100644 --- a/.changeset/chat-stop-successor-boundary.md +++ b/.changeset/chat-stop-successor-boundary.md @@ -5,3 +5,4 @@ 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 acc844fa174..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; }; @@ -79,17 +80,15 @@ 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. A captured recovery callback replaces ordinary cursor seeding. + * no cursor. * Returns whether the cursor was accepted. */ export function seedTranscriptCursor( transport: Pick, chatId: string, - cursors: { lastOutEventId?: string; lastInEventId?: string } | undefined, - completeRecovery?: (lastEventId: string | undefined, lastInEventId?: string) => boolean + cursors: TranscriptCursors | undefined ): boolean { const lastEventId = cursors?.lastOutEventId; - if (completeRecovery) return completeRecovery(lastEventId, cursors?.lastInEventId); if (!lastEventId) return false; transport.seedResumeCursor(chatId, lastEventId); return true; @@ -150,11 +149,12 @@ export function useLoadTranscript( .current({ chatId, ...(limit !== undefined ? { limit } : {}) }) .then((result) => { if (cancelled) return; - if (transport) { - const seeded = seedTranscriptCursor(transport, chatId, result.cursors, completeRecovery); - if (completeRecovery && !seeded) { + 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, diff --git a/packages/trigger-sdk/src/v3/chat-stop.test.ts b/packages/trigger-sdk/src/v3/chat-stop.test.ts index 4339524e0ea..53d6dcee8c8 100644 --- a/packages/trigger-sdk/src/v3/chat-stop.test.ts +++ b/packages/trigger-sdk/src/v3/chat-stop.test.ts @@ -4,6 +4,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { TriggerChatTransport, type ChatSessionPersistedState, + type ChatTransportEvent, type TriggerChatTransportOptions, } from "./chat.js"; @@ -505,7 +506,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession("constructor"); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("5", "9")).toBe(false); + expect(recover({ lastOutEventId: "5", lastInEventId: "9" })).toBe(false); await expect(send()).rejects.toThrow("Stopped chat response cannot be matched"); expect(inputSeq).toBe(13); }); @@ -542,11 +543,11 @@ describe("Stop with a successor response", () => { if (hydrate === "setSession") transport.setSession("chat", session); const stale = transport.prepareTranscriptRecovery("chat"); if (!stale) throw new Error("Expected transcript recovery"); - expect(stale("5", "9")).toBe(false); + 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("11", "10")).toBe(true); + 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)); @@ -575,7 +576,7 @@ describe("Stop with a successor response", () => { } const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("17", "11")).toBe(true); + 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"); @@ -597,8 +598,15 @@ describe("Stop with a successor response", () => { expect(await stopped).toBe(true); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("17", "11")).toBe(false); + 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 () => { @@ -632,7 +640,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession(hydrate); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("11", "11")).toBe(true); + expect(recover({ lastOutEventId: "11", lastInEventId: "11" })).toBe(true); transport.seedResumeCursor("chat", "11"); expect(saved).toMatchObject({ lastEventId: "11", @@ -658,7 +666,13 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession("constructor"); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover(cursor, "11")).toBe(false); + 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" @@ -685,7 +699,7 @@ describe("Stop with a successor response", () => { await hydrateBlockedSession("constructor"); const recover = transport.prepareTranscriptRecovery("chat"); if (!recover) throw new Error("Expected transcript recovery"); - expect(recover("11", "11")).toBe(true); + expect(recover({ lastOutEventId: "11", lastInEventId: "11" })).toBe(true); if (hydrate !== "none") { const session = transport.getSession("chat"); if (!session) throw new Error("Expected persisted state"); @@ -707,16 +721,21 @@ describe("Stop with a successor response", () => { } ); - it("retains unknown recovery state after a bounded empty response", async () => { + 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("11", "11")).toBe(true); + 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)).resolves.toBe(""); + 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); @@ -742,7 +761,7 @@ describe("Stop with a successor response", () => { holdStop = true; const stopped = transport.stopGeneration("chat"); await vi.waitFor(() => expect(pendingStop).toBeDefined()); - expect(recover("11", "11")).toBe(false); + 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); @@ -751,6 +770,120 @@ describe("Stop with a successor response", () => { 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) => { diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 76380c98e4f..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, @@ -710,36 +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; - requiresTranscriptReload?: boolean; - transcriptRecoveryInputSeq?: number; +type ChatSessionState = ChatSessionPersistedState & { /** Identifies the stopped boundary for pending Stop acknowledgments. Never persisted. */ stoppedBoundary?: symbol; - /** 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; /** Identifies the latest transcript load for this blocked session. Never persisted. */ transcriptRecovery?: symbol; - skipSettledPeek?: 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; }; /** @@ -821,20 +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, - skipToTurnComplete: session.skipToTurnComplete, - supersededInputSeq: session.supersededInputSeq, - requiresTranscriptReload: session.requiresTranscriptReload, - transcriptRecoveryInputSeq: session.transcriptRecoveryInputSeq, - outstandingTurnAbandoned: session.outstandingTurnAbandoned, - skipSettledPeek: session.skipSettledPeek, - closed: session.closed, - closedReason: session.closedReason, - }); + this.sessions.set(chatId, this.toPersisted(session)); } } } @@ -1492,22 +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, - skipToTurnComplete: session.skipToTurnComplete, - supersededInputSeq: session.supersededInputSeq, - requiresTranscriptReload: session.requiresTranscriptReload, - transcriptRecoveryInputSeq: session.transcriptRecoveryInputSeq, - outstandingTurnAbandoned: session.outstandingTurnAbandoned, - skipSettledPeek: session.skipSettledPeek, - }) - ); - 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); } /** @@ -1523,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; @@ -1535,10 +1488,11 @@ export class TriggerChatTransport implements ChatTransport { * 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 - ): ((lastEventId: string | undefined, lastInEventId?: string) => boolean) | undefined => { + ): ((cursors: TranscriptCursors | undefined) => boolean) | undefined => { const state = this.sessions.get(chatId); if (!state?.requiresTranscriptReload || state.closed) return undefined; @@ -1552,7 +1506,7 @@ export class TriggerChatTransport implements ChatTransport { transcriptRecoveryInputSeq, } = state; const stoppedInputSeq = supersededInputSeq ?? transcriptRecoveryInputSeq; - return (loadedEventId, loadedInEventId) => { + return (cursors) => { if (state.transcriptRecovery !== token) return false; state.transcriptRecovery = undefined; if ( @@ -1564,15 +1518,33 @@ export class TriggerChatTransport implements ChatTransport { state.activeInputSeq !== activeInputSeq || state.skipToTurnComplete !== skipToTurnComplete || state.supersededInputSeq !== supersededInputSeq || - state.transcriptRecoveryInputSeq !== transcriptRecoveryInputSeq || + state.transcriptRecoveryInputSeq !== transcriptRecoveryInputSeq + ) { + return false; + } + if ( stoppedInputSeq === undefined || !Number.isSafeInteger(stoppedInputSeq) || - stoppedInputSeq < 0 || + 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) || - BigInt(loadedInEventId) < BigInt(stoppedInputSeq) || loadedEventId === undefined || - !/^\d+$/.test(loadedEventId) || + !/^\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))) ) { @@ -2327,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 @@ -2358,6 +2333,16 @@ export class TriggerChatTransport implements ChatTransport { ); } + 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 ( @@ -2651,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/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 f00f0e672c9..a1c85a331b1 100644 --- a/packages/trigger-sdk/test/use-load-transcript.test.ts +++ b/packages/trigger-sdk/test/use-load-transcript.test.ts @@ -85,14 +85,7 @@ describe("transcript recovery", () => { const { transport, saved } = blockedTransport(); const recover = transport.prepareTranscriptRecovery("chat-1"); expect(recover).toBeDefined(); - expect( - seedTranscriptCursor( - transport, - "chat-1", - { lastOutEventId: "9007199254740993", lastInEventId: "4" }, - recover - ) - ).toBe(true); + expect(recover?.({ lastOutEventId: "9007199254740993", lastInEventId: "4" })).toBe(true); expect(transport.getSession("chat-1")).toMatchObject({ lastEventId: "9007199254740993", requiresTranscriptReload: false, @@ -104,33 +97,57 @@ describe("transcript recovery", () => { isStreaming: undefined, }); expect(saved).toHaveLength(1); - expect(recover?.("9007199254740994", "4")).toBe(false); + expect(recover?.({ lastOutEventId: "9007199254740994", lastInEventId: "4" })).toBe(false); expect(saved).toHaveLength(1); }); - it.each([undefined, "", "NaN", "-1", "1e20", "9007199254740993x", "9007199254740992", "42"])( - "keeps sends blocked for an invalid or stale checkpoint: %s", + 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( - seedTranscriptCursor(transport, "chat-1", { lastOutEventId, lastInEventId: "4" }, recover) - ).toBe(false); + 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")?.("9007199254740993", "4")).toBe(false); + 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")?.("42", "4")).toBe(true); + expect( + transport.prepareTranscriptRecovery("chat-1")?.({ lastOutEventId: "42", lastInEventId: "4" }) + ).toBe(true); expect(transport.getSession("chat-1")?.lastEventId).toBe("42"); }); @@ -138,25 +155,16 @@ describe("transcript recovery", () => { const { transport } = blockedTransport({ lastEventId: undefined }); const stale = transport.prepareTranscriptRecovery("chat-1"); const current = transport.prepareTranscriptRecovery("chat-1"); - expect( - seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "50", lastInEventId: "4" }, stale) - ).toBe(false); + expect(stale?.({ lastOutEventId: "50", lastInEventId: "4" })).toBe(false); expect(transport.getSession("chat-1")?.lastEventId).toBeUndefined(); - expect( - seedTranscriptCursor( - transport, - "chat-1", - { lastOutEventId: "42", lastInEventId: "4" }, - current - ) - ).toBe(true); + 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?.("9007199254740993", "4")).toBe(false); + expect(recover?.({ lastOutEventId: "9007199254740993", lastInEventId: "4" })).toBe(false); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); @@ -164,7 +172,7 @@ describe("transcript recovery", () => { const { transport } = blockedTransport({ lastEventId: undefined }); const recover = transport.prepareTranscriptRecovery("chat-1"); transport.seedResumeCursor("chat-1", "42"); - expect(recover?.("43", "4")).toBe(false); + expect(recover?.({ lastOutEventId: "43", lastInEventId: "4" })).toBe(false); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); @@ -173,37 +181,51 @@ describe("transcript recovery", () => { const recover = transport.prepareTranscriptRecovery("chat-1"); if (operation === "abandon") transport.clearSupersedeGate("chat-1"); else transport.dispose(); - expect(recover?.("9007199254740993", "4")).toBe(false); + expect(recover?.({ lastOutEventId: "9007199254740993", lastInEventId: "4" })).toBe(false); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); - it.each([undefined, "", "NaN", "-1", "4x", "3"])( - "rejects a newer output checkpoint without stopped-input evidence: %s", + 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( - seedTranscriptCursor(transport, "chat-1", { lastOutEventId: "42", lastInEventId }, recover) - ).toBe(false); + 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( - seedTranscriptCursor( - transport, - "chat-1", - { lastOutEventId: "42", lastInEventId: "100" }, - recover - ) - ).toBe(false); + expect(() => recover?.({ lastOutEventId: "42", lastInEventId: "100" })).toThrow( + "Transcript recovery requires a stopped input sequence" + ); expect(transport.getSession("chat-1")?.requiresTranscriptReload).toBe(true); }); @@ -215,4 +237,16 @@ describe("transcript recovery", () => { 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 From 4c100ce9108752878354db3f244b0db6b07334cc Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Mon, 21 Sep 2026 17:03:58 -0700 Subject: [PATCH 5/5] fix(sdk): keep idle Stop from discarding the next response --- .changeset/chat-stop-successor-boundary.md | 2 + packages/trigger-sdk/src/v3/chat-stop.test.ts | 196 ++++++++++++++++++ packages/trigger-sdk/src/v3/chat.test.ts | 23 +- packages/trigger-sdk/src/v3/chat.ts | 157 ++++++++------ 4 files changed, 307 insertions(+), 71 deletions(-) diff --git a/.changeset/chat-stop-successor-boundary.md b/.changeset/chat-stop-successor-boundary.md index f356ac582b0..6a2069ec1c6 100644 --- a/.changeset/chat-stop-successor-boundary.md +++ b/.changeset/chat-stop-successor-boundary.md @@ -3,6 +3,8 @@ --- Keep new chat responses intact after Stop, including slow Stop acknowledgments and page reloads. +Stop on an idle hydrated session no longer discards the next response. Early resumed Stop retains its protection. + 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/packages/trigger-sdk/src/v3/chat-stop.test.ts b/packages/trigger-sdk/src/v3/chat-stop.test.ts index 53d6dcee8c8..8a17804fa23 100644 --- a/packages/trigger-sdk/src/v3/chat-stop.test.ts +++ b/packages/trigger-sdk/src/v3/chat-stop.test.ts @@ -78,12 +78,14 @@ describe("Stop with a successor response", () => { let outputHeaders: { peek: boolean; timeout: number }[]; let inputSeq: number; let holdStop: boolean; + let holdMessages: boolean; let stopStatus: number; let settled: boolean; let resumeAfterStoppedCheckpoint: boolean; let emptyRecoveredOutput: boolean; let includeSequence: boolean; let pendingStop: { response: ServerResponse; seq: number } | undefined; + let pendingMessages: { response: ServerResponse; seq: number }[]; let saved: ChatSessionPersistedState | null; function createTransport( @@ -113,12 +115,14 @@ describe("Stop with a successor response", () => { outputHeaders = []; inputSeq = 10; holdStop = false; + holdMessages = false; stopStatus = 200; settled = false; resumeAfterStoppedCheckpoint = false; emptyRecoveredOutput = false; includeSequence = true; pendingStop = undefined; + pendingMessages = []; saved = null; server = createServer(async (request, response) => { if (request.method === "POST") { @@ -130,6 +134,8 @@ describe("Stop with a successor response", () => { const seq = inputSeq++; if (isStop && holdStop) { pendingStop = { response, seq }; + } else if (!isStop && holdMessages) { + pendingMessages.push({ response, seq }); } else { appendResponse(response, seq, isStop ? stopStatus : 200); } @@ -288,6 +294,162 @@ describe("Stop with a successor response", () => { } ); + it.each([ + ["constructor", undefined], + ["constructor", "1"], + ["setSession", undefined], + ["setSession", "1"], + ] as const)("does not gate idle %s hydration with cursor %s", async (hydrate, lastEventId) => { + const session = { publicAccessToken: "test-token", lastEventId }; + if (hydrate === "constructor") { + transport.dispose(); + transport = createTransport(session); + } else { + transport.setSession("chat", session); + } + expect(await transport.stopGeneration("chat")).toBe(true); + expect(transport.getSession("chat")?.skipToTurnComplete).not.toBe(true); + const next = await send(); + emit([...reply(2), complete(7, 11)]); + await expect(readText(next)).resolves.toBe("New response"); + expect(inputSeq).toBe(12); + }); + + it.each([ + ["message", "idle"], + ["action", "idle"], + ["message", "repeated Stop"], + ["action", "repeated Stop"], + ] as const)("retains Stop during a pending %s append (%s)", async (kind, state) => { + holdMessages = true; + const first = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" }); + await vi.waitFor(() => expect(pendingMessages).toHaveLength(1)); + expect(await transport.stopGeneration("chat")).toBe(true); + if (state === "repeated Stop") expect(await transport.stopGeneration("chat")).toBe(true); + expect(transport.getSession("chat")).toMatchObject({ + skipToTurnComplete: true, + transcriptRecoveryInputSeq: 11, + }); + expect(transport.getSession("chat")).not.toHaveProperty("pendingInputCount"); + holdMessages = false; + const pending = pendingMessages[0]!; + appendResponse(pending.response, pending.seq); + const firstStream = await first; + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + const firstResult = readText(firstStream); + const nextInput = inputSeq; + const next = await send(); + await expect(firstResult).resolves.toBe(""); + emit(oldTailAndReply(10, nextInput)); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it.each(["message", "action"] as const)( + "retains Stop from the first %s acknowledgment event", + async (kind) => { + let stopping: Promise | undefined; + let stopFirst = true; + transport.setOnEvent((event) => { + if (event.type === "message-sent" && event.source !== "stop" && stopFirst) { + stopFirst = false; + queueMicrotask(() => { + stopping = transport.stopGeneration("chat"); + }); + } + }); + const first = await (kind === "message" + ? send() + : transport.sendAction("chat", { type: "undo" })); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + await expect(stopping).resolves.toBe(true); + expect(transport.getSession("chat")?.skipToTurnComplete).toBe(true); + const stopped = readText(first); + const next = await send(); + await expect(stopped).resolves.toBe(""); + emit(oldTailAndReply(10, 12)); + await expect(readText(next)).resolves.toBe("New response"); + } + ); + + it.each(["message", "action"] as const)( + "clears pending activity after a failed %s append", + async (kind) => { + holdMessages = true; + const first = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" }); + const failed = expect(first).rejects.toThrow(); + await vi.waitFor(() => expect(pendingMessages).toHaveLength(1)); + const pending = pendingMessages[0]!; + appendResponse(pending.response, pending.seq, 400); + await failed; + holdMessages = false; + await transport.stopGeneration("chat"); + expect(transport.getSession("chat")?.skipToTurnComplete).not.toBe(true); + const next = await send(); + emit([...reply(2), complete(7, 12)]); + await expect(readText(next)).resolves.toBe("New response"); + } + ); + + it.each([ + ["message", "idle"], + ["action", "idle"], + ["message", "abandoned"], + ["action", "abandoned"], + ] as const)("keeps a rejected %s append idle after Stop (%s)", async (kind, state) => { + transport.setSession("chat", { publicAccessToken: "test-token", isStreaming: false }); + if (state === "abandoned") transport.clearSupersedeGate("chat"); + holdMessages = true; + const first = kind === "message" ? send() : transport.sendAction("chat", { type: "undo" }); + const failed = expect(first).rejects.toThrow(); + await vi.waitFor(() => expect(pendingMessages).toHaveLength(1)); + await transport.stopGeneration("chat"); + const pending = pendingMessages[0]!; + appendResponse(pending.response, pending.seq, 400); + await failed; + holdMessages = false; + expect(transport.getSession("chat")?.skipToTurnComplete).not.toBe(true); + const next = await send(); + emit([...reply(2), complete(7, 12)]); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("retains pending activity until every overlapping append finishes", async () => { + holdMessages = true; + const first = transport.sendAction("chat", { type: "first" }); + const firstFailure = expect(first).rejects.toThrow(); + const second = transport.sendAction("chat", { type: "second" }); + await vi.waitFor(() => expect(pendingMessages).toHaveLength(2)); + const rejected = pendingMessages[0]!; + appendResponse(rejected.response, rejected.seq, 400); + await firstFailure; + await transport.stopGeneration("chat"); + expect(transport.getSession("chat")?.skipToTurnComplete).toBe(true); + holdMessages = false; + const accepted = pendingMessages[1]!; + appendResponse(accepted.response, accepted.seq); + const stopped = readText(await second); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); + const next = await send(); + await expect(stopped).resolves.toBe(""); + emit(oldTailAndReply(11, 13)); + await expect(readText(next)).resolves.toBe("New response"); + }); + + it("does not mark an already-canceled reconnect as an outstanding turn", async () => { + const abort = new AbortController(); + abort.abort(); + const resumed = await transport.reconnectToStream({ + chatId: "chat", + abortSignal: abort.signal, + }); + if (!resumed) throw new Error("Expected a resumed stream"); + await expect(readText(resumed)).resolves.toBe(""); + expect(await transport.stopGeneration("chat")).toBe(true); + 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 an empty settled resume", async () => { settled = true; transport.setSession("chat", { publicAccessToken: "test-token", lastEventId: "1" }); @@ -303,6 +465,25 @@ describe("Stop with a successor response", () => { await expect(readText(next)).resolves.toBe("New response"); }); + it("does not transfer an unknown resumed turn to a replacement idle session", async () => { + 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)); + abort.abort(); + await expect(readText(resumed)).resolves.toBe(""); + const session = transport.getSession("chat")!; + expect(session).not.toHaveProperty("resumedUnknownTurn"); + transport.setSession("chat", session); + await transport.stopGeneration("chat"); + 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( @@ -560,6 +741,8 @@ describe("Stop with a successor response", () => { it.each([false, true])( "retains the first Stop input through repeated Stop (delayed first acknowledgment: %s)", async (delayed) => { + await transport.reconnectToStream({ chatId: "chat" }); + await vi.waitFor(() => expect(outputs).toHaveLength(1)); holdStop = delayed; const firstStop = transport.stopGeneration("chat"); if (delayed) await vi.waitFor(() => expect(pendingStop).toBeDefined()); @@ -770,6 +953,19 @@ describe("Stop with a successor response", () => { expect(transport.getSession("chat")).toMatchObject({ requiresTranscriptReload: true }); }); + it("retains the stopped boundary for recovered output before reconnect", async () => { + await hydrateBlockedSession("constructor"); + expect( + transport.prepareTranscriptRecovery("chat")?.({ lastOutEventId: "11", lastInEventId: "11" }) + ).toBe(true); + await transport.stopGeneration("chat"); + expect(transport.getSession("chat")?.skipToTurnComplete).toBe(true); + const next = await send(); + emit([complete(12, 12), ...reply(13), complete(18, 14)]); + await expect(readText(next)).resolves.toBe("New response"); + expect(inputSeq).toBe(15); + }); + it.each(["abort", "stop"] as const)("closes recovery quietly after %s", async (operation) => { await hydrateBlockedSession("constructor"); expect( diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index afe31e84d6d..351f58851a2 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1435,18 +1435,21 @@ describe("TriggerChatTransport", () => { expect(await drainChunks(next)).toEqual(sampleChunks); }); - it("does not gate a stop with no turn outstanding", async () => { - mockFetch([() => defaultSseResponse()]); - - const transport = await armedGate("chat-idle-stop", { - publicAccessToken: "p", - isStreaming: false, - }); + it.each([undefined, false])( + "does not gate an idle stop with isStreaming %s", + async (isStreaming) => { + mockFetch([() => defaultSseResponse()]); + + const transport = await armedGate("chat-idle-stop", { + publicAccessToken: "p", + isStreaming, + }); - const stream = await send(transport, "chat-idle-stop"); + const stream = await send(transport, "chat-idle-stop"); - expect(await drainChunks(stream)).toEqual(sampleChunks); - }); + expect(await drainChunks(stream)).toEqual(sampleChunks); + } + ); it("does not arm or write a stop when the abort lands after the boundary", async () => { const bodies: string[] = []; diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 4a277cc7645..e5bf2de4a61 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -712,6 +712,10 @@ export type TriggerChatTransportOptions = { * @internal */ type ChatSessionState = ChatSessionPersistedState & { + /** Counts unfinished turn-producing sends. Never persisted. */ + pendingInputCount?: number; + /** An opened resume still has unknown turn state. Passive reader cancellation retains this marker. */ + resumedUnknownTurn?: boolean; /** Identifies the stopped boundary for pending Stop acknowledgments. Never persisted. */ stoppedBoundary?: symbol; /** Identifies the latest transcript load for this blocked session. Never persisted. */ @@ -945,35 +949,37 @@ export class TriggerChatTransport implements ChatTransport { const sendChatMessage = (token: string) => this.appendInputChunk(chatId, token, serializedBody, partId); - const inSeq = await this.sendWithEvents( - chatId, - trigger, - { - messageId: messageId ?? messages.at(-1)?.id, - partId, - bodyBytes: byteLength(serializedBody), - }, - () => this.callWithAuthRetry(chatId, state, sendChatMessage) - ); + return this.withPendingInput(state, async () => { + const inSeq = await this.sendWithEvents( + chatId, + trigger, + { + messageId: messageId ?? messages.at(-1)?.id, + partId, + bodyBytes: byteLength(serializedBody), + }, + () => this.callWithAuthRetry(chatId, state, sendChatMessage) + ); - // Cancel any in-flight stream for this chat — the new turn supersedes it. - const activeStream = this.activeStreams.get(chatId); - if (activeStream) { - activeStream.abort(); - this.activeStreams.delete(chatId); - } + // Cancel any in-flight stream for this chat — the new turn supersedes it. + const activeStream = this.activeStreams.get(chatId); + if (activeStream) { + activeStream.abort(); + this.activeStreams.delete(chatId); + } - state.activeInputSeq = inSeq; - this.requireStoppedTurnCorrelation(chatId, state, inSeq); - state.isStreaming = true; - state.outstandingTurnAbandoned = false; - state.skipSettledPeek = false; - this.notifySessionChange(chatId, state); + 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. - return this.subscribeToSessionStream(state, abortSignal, chatId, { - sinceInSeq: inSeq, - sendStopOnAbort: true, + // Owning turn: aborting this live send stops the turn the user drives. + return this.subscribeToSessionStream(state, abortSignal, chatId, { + sinceInSeq: inSeq, + sendStopOnAbort: true, + }); }); }; @@ -1278,6 +1284,10 @@ export class TriggerChatTransport implements ChatTransport { if (state.isStreaming === false && !this.watchMode) return null; if (this.activeStreams.has(options.chatId)) return null; + if (state.isStreaming === undefined && !options.abortSignal?.aborted) { + state.resumedUnknownTurn = true; + } + const abortController = new AbortController(); this.activeStreams.set(options.chatId, abortController); @@ -1312,10 +1322,7 @@ export class TriggerChatTransport implements ChatTransport { // 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 !== false || state.activeInputSeq !== undefined) - ) { + if (this.hasOutstandingTurn(state)) { this.armStoppedBoundary(state); } const stoppedBoundary = state.skipToTurnComplete @@ -1338,6 +1345,7 @@ export class TriggerChatTransport implements ChatTransport { // otherwise a reload resumes mid-turn and replays the chunks the user // explicitly stopped. state.isStreaming = false; + state.resumedUnknownTurn = undefined; this.notifySessionChange(chatId, state); const partId = crypto.randomUUID(); @@ -1411,36 +1419,38 @@ export class TriggerChatTransport implements ChatTransport { const partId = crypto.randomUUID(); const send = (token: string) => this.appendInputChunk(chatId, token, body, partId); - const inSeq = await this.sendWithEvents( - chatId, - "action", - { partId, bodyBytes: byteLength(body) }, - () => this.callWithAuthRetry(chatId, state, send) - ); + return this.withPendingInput(state, async () => { + const inSeq = await this.sendWithEvents( + chatId, + "action", + { partId, bodyBytes: byteLength(body) }, + () => this.callWithAuthRetry(chatId, state, send) + ); - // Supersede any in-flight reader before subscribing — same as - // `sendMessages`. Two concurrent readers both write `state.lastEventId` - // and the slower one can regress the cursor, replaying records on the - // next reconnect. - const activeStream = this.activeStreams.get(chatId); - if (activeStream) { - activeStream.abort(); - this.activeStreams.delete(chatId); - } + // Supersede any in-flight reader before subscribing — same as + // `sendMessages`. Two concurrent readers both write `state.lastEventId` + // and the slower one can regress the cursor, replaying records on the + // next reconnect. + const activeStream = this.activeStreams.get(chatId); + if (activeStream) { + activeStream.abort(); + this.activeStreams.delete(chatId); + } - // 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); + // 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. - return this.subscribeToSessionStream(state, options?.abortSignal, chatId, { - sinceInSeq: inSeq, - sendStopOnAbort: true, + // Owning action: aborting this send stops the turn the user drives. + return this.subscribeToSessionStream(state, options?.abortSignal, chatId, { + sinceInSeq: inSeq, + sendStopOnAbort: true, + }); }); }; @@ -1713,6 +1723,7 @@ export class TriggerChatTransport implements ChatTransport { for (const state of this.sessions.values()) { state.transcriptRecovery = undefined; state.stoppedBoundary = undefined; + state.resumedUnknownTurn = undefined; } this.coordinator?.dispose(); this.coordinator = null; @@ -1722,7 +1733,29 @@ export class TriggerChatTransport implements ChatTransport { // Internal helpers // ------------------------------------------------------------------------- + private async withPendingInput(state: ChatSessionState, op: () => Promise): Promise { + state.pendingInputCount = (state.pendingInputCount ?? 0) + 1; + try { + return await op(); + } finally { + state.pendingInputCount--; + } + } + + private hasOutstandingTurn(state: ChatSessionState): boolean { + return ( + !state.outstandingTurnAbandoned && + (state.isStreaming === true || + state.activeInputSeq !== undefined || + (state.isStreaming === undefined && + ((state.pendingInputCount ?? 0) > 0 || + state.resumedUnknownTurn === true || + state.skipSettledPeek === true))) + ); + } + private armStoppedBoundary(state: ChatSessionState): symbol { + state.resumedUnknownTurn = undefined; state.skipToTurnComplete = true; state.supersededInputSeq = state.activeInputSeq; state.transcriptRecoveryInputSeq = undefined; @@ -1731,6 +1764,7 @@ export class TriggerChatTransport implements ChatTransport { } private clearStoppedBoundary(state: ChatSessionState): void { + state.resumedUnknownTurn = undefined; state.skipToTurnComplete = false; state.supersededInputSeq = undefined; state.transcriptRecoveryInputSeq = undefined; @@ -2156,12 +2190,9 @@ export class TriggerChatTransport implements ChatTransport { () => { // A late abort (unmount, or the consumer dropping a drained stream) // has no turn to stop: don't gate the next one, don't write a stop. - const outstanding = - !state.outstandingTurnAbandoned && - (state.isStreaming !== false || state.activeInputSeq !== undefined); if ( options?.sendStopOnAbort !== false && - outstanding && + this.hasOutstandingTurn(state) && !internalAbort.signal.aborted && this.activeStreams.get(chatId) === internalAbort ) { @@ -2352,6 +2383,7 @@ export class TriggerChatTransport implements ChatTransport { this.activeStreams.get(chatId) === internalAbort ) { state.isStreaming = false; + state.resumedUnknownTurn = undefined; this.notifySessionChange(chatId, state); } return null; @@ -2582,6 +2614,7 @@ export class TriggerChatTransport implements ChatTransport { sinceInSeq = undefined; state.skipSettledPeek = false; state.isStreaming = false; + state.resumedUnknownTurn = undefined; this.notifySessionChange(chatId, state); this.coordinator?.release(chatId); this.coordinator?.broadcastSession(chatId, { @@ -2609,6 +2642,7 @@ export class TriggerChatTransport implements ChatTransport { // Its first data record establishes an active turn for Stop. if (!state.outstandingTurnAbandoned && state.isStreaming !== true) { state.isStreaming = true; + state.resumedUnknownTurn = undefined; this.notifySessionChange(chatId, state); } if (!sawFirstChunk) { @@ -2639,6 +2673,7 @@ export class TriggerChatTransport implements ChatTransport { // 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; + state.resumedUnknownTurn = undefined; this.notifySessionChange(chatId, state); } this.emitEvent({