From 27b04fb9108bce90ec10389dd84dc52aa9ced189 Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 09:56:28 +0000 Subject: [PATCH 01/10] fix(sdk): watch-mode chat subscriptions survive quiet windows (TRI-13065) --- .changeset/watch-mode-keepalive.md | 5 + packages/trigger-sdk/src/v3/chat.test.ts | 127 +++++++++++++++++- packages/trigger-sdk/src/v3/chat.ts | 8 +- .../test/chat-transport-events.test.ts | 13 +- 4 files changed, 147 insertions(+), 6 deletions(-) create mode 100644 .changeset/watch-mode-keepalive.md diff --git a/.changeset/watch-mode-keepalive.md b/.changeset/watch-mode-keepalive.md new file mode 100644 index 00000000000..49e8fafb6cf --- /dev/null +++ b/.changeset/watch-mode-keepalive.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Watch-mode chat subscriptions now stay connected across quiet periods. diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index 21aa004c9a9..42da9c83951 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1218,6 +1218,122 @@ describe("TriggerChatTransport", () => { }); }); + describe("watch mode across long-poll window boundaries", () => { + function settled(response: Response): Response { + const headers = new Headers(response.headers); + headers.set("X-Session-Settled", "true"); + return new Response(response.body, { status: 200, headers }); + } + + it("resubscribes after a completed turn and receives a later wake", async () => { + let subscribeCount = 0; + global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionOutSubscribeUrl(urlStr)) { + subscribeCount++; + // Window 1: a turn completes, then the body EOFs with no + // settled header — the quiet long-poll boundary. + return subscribeCount === 1 + ? defaultSseResponse([ + { type: "text-delta", id: "p1", delta: "turn1" }, + { type: "trigger:turn-complete" }, + ]) + : settled(defaultSseResponse([{ type: "text-delta", id: "p2", delta: "wake" }])); + } + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-watch-eof": { publicAccessToken: "p", isStreaming: true } }, + }); + + const stream = await transport.reconnectToStream({ chatId: "chat-watch-eof" }); + const chunks = await drainChunks(stream!); + + expect(subscribeCount).toBe(2); + expect(chunks).toEqual([ + { type: "text-delta", id: "p1", delta: "turn1" }, + { type: "text-delta", id: "p2", delta: "wake" }, + ]); + }); + + it("stops when the server says the session settled", async () => { + let subscribeCount = 0; + global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionOutSubscribeUrl(urlStr)) { + subscribeCount++; + return settled( + defaultSseResponse([ + { type: "text-delta", id: "p1", delta: "last" }, + { type: "trigger:turn-complete" }, + ]) + ); + } + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-watch-settled": { publicAccessToken: "p", isStreaming: true } }, + }); + + const stream = await transport.reconnectToStream({ chatId: "chat-watch-settled" }); + const chunks = await drainChunks(stream!); + + expect(subscribeCount).toBe(1); + expect(chunks).toHaveLength(1); + expect(transport.getSession("chat-watch-settled")?.isStreaming).toBe(false); + }); + + it("stops promptly when aborted during backoff", async () => { + vi.useFakeTimers(); + try { + let subscribeCount = 0; + global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionStreamAppendUrl(urlStr)) return defaultAppendResponse(); + if (isSessionOutSubscribeUrl(urlStr)) { + subscribeCount++; + // Every window is quiet: EOF with no records, never settled. + return defaultSseResponse([]); + } + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const abortController = new AbortController(); + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-watch-abort": { publicAccessToken: "p", isStreaming: true } }, + }); + + const stream = await transport.reconnectToStream({ + chatId: "chat-watch-abort", + abortSignal: abortController.signal, + }); + const drained = drainChunks(stream!); + await vi.advanceTimersByTimeAsync(10_000); + // The budget doesn't apply in watch mode, so it is still reconnecting. + expect(subscribeCount).toBeGreaterThan(6); + + const countAtAbort = subscribeCount; + abortController.abort(); + await drained; + await vi.advanceTimersByTimeAsync(10_000); + expect(subscribeCount).toBe(countAtAbort); + } finally { + vi.useRealTimers(); + } + }); + }); + describe("multi-tab coordination", () => { it("isReadOnly defaults to false when multiTab is disabled", () => { const transport = new TriggerChatTransport({ @@ -1488,9 +1604,18 @@ describe("TriggerChatTransport", () => { { type: "text-delta", id: "p2", delta: "Again" }, { type: "trigger:turn-complete" }, ]; + let subscribeCount = 0; global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { const urlStr = typeof url === "string" ? url : url.toString(); - if (isSessionOutSubscribeUrl(urlStr)) return defaultSseResponse(turn1); + if (isSessionOutSubscribeUrl(urlStr)) { + subscribeCount++; + if (subscribeCount === 1) return defaultSseResponse(turn1); + // Watch mode reconnects past the body EOF; settle so the drain ends. + const response = defaultSseResponse([]); + const headers = new Headers(response.headers); + headers.set("X-Session-Settled", "true"); + return new Response(response.body, { status: 200, headers }); + } throw new Error(`Unexpected URL: ${urlStr}`); }); diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 137222b1bae..07e89fc6d9a 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -1780,11 +1780,13 @@ export class TriggerChatTransport implements ChatTransport { let eofResubscribes = 0; const resumeAfterEof = async () => { + // Watch mode is a standing subscription: it outlives turn-complete + // (which clears `isStreaming`) and idle windows EOF by design, so the + // give-up budget doesn't apply. Only abort or a settled session ends it. while ( - state.isStreaming && + (this.watchMode || (state.isStreaming && eofResubscribes < MAX_EOF_RESUBSCRIBES)) && !currentSubscription?.sessionSettled && - !combinedSignal.aborted && - eofResubscribes < MAX_EOF_RESUBSCRIBES + !combinedSignal.aborted ) { eofResubscribes++; // Sleep, but wake immediately on abort — otherwise a stop lands diff --git a/packages/trigger-sdk/test/chat-transport-events.test.ts b/packages/trigger-sdk/test/chat-transport-events.test.ts index ded1895dfd8..39f4e53d722 100644 --- a/packages/trigger-sdk/test/chat-transport-events.test.ts +++ b/packages/trigger-sdk/test/chat-transport-events.test.ts @@ -211,11 +211,20 @@ describe("transport stream events", () => { ``, ].join("\n"); + let subscribes = 0; const { transport, events } = makeTransport({ watch: true, sessions: { c1: { publicAccessToken: "tok_test", isStreaming: true } }, - fetch: async (_url, _init, ctx) => - ctx.endpoint === "in" ? jsonOk() : sseResponse(TWO_TURNS), + fetch: async (_url, _init, ctx) => { + if (ctx.endpoint === "in") return jsonOk(); + if (subscribes++ > 0) { + // Watch mode reconnects past the body EOF; settle so the read ends. + const settled = sseResponse(""); + settled.headers.set("X-Session-Settled", "true"); + return settled; + } + return sseResponse(TWO_TURNS); + }, }); const stream = await transport.reconnectToStream({ chatId: "c1" }); From 574a84d659a03d3be00048d925fb292689b7898e Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 10:43:20 +0000 Subject: [PATCH 02/10] fix(sdk): stopping a turn is an owning-transport action, not an abort side effect (TRI-13070) --- .changeset/watch-mode-keepalive.md | 2 + packages/trigger-sdk/src/v3/chat.test.ts | 88 ++++++++++++++++++++++++ packages/trigger-sdk/src/v3/chat.ts | 21 +++++- 3 files changed, 108 insertions(+), 3 deletions(-) diff --git a/.changeset/watch-mode-keepalive.md b/.changeset/watch-mode-keepalive.md index 49e8fafb6cf..e50a04f2a8c 100644 --- a/.changeset/watch-mode-keepalive.md +++ b/.changeset/watch-mode-keepalive.md @@ -3,3 +3,5 @@ --- Watch-mode chat subscriptions now stay connected across quiet periods. + +Read-only chat subscriptions no longer stop a turn when they disconnect. diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index 42da9c83951..f9851f52ade 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1334,6 +1334,94 @@ describe("TriggerChatTransport", () => { }); }); + describe("reconnectToStream stop-on-abort ownership (TRI-13070)", () => { + // A quiet stream: EOF, no records, never settled — the subscription + // stays alive (watch mode) so an abort mid-flight exercises the stop path. + function quietWatchTransport(): { + transport: TriggerChatTransport; + appends: () => number; + } { + let appendCount = 0; + global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionStreamAppendUrl(urlStr)) { + appendCount++; + return defaultAppendResponse(); + } + if (isSessionOutSubscribeUrl(urlStr)) return defaultSseResponse([]); + throw new Error(`Unexpected URL: ${urlStr}`); + }); + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-own": { publicAccessToken: "p", isStreaming: true } }, + }); + return { transport, appends: () => appendCount }; + } + + it("passive subscriber aborting writes no stop chunk to .in", async () => { + vi.useFakeTimers(); + try { + const { transport, appends } = quietWatchTransport(); + const abort = new AbortController(); + const stream = await transport.reconnectToStream({ + chatId: "chat-own", + abortSignal: abort.signal, + }); + const drained = drainChunks(stream!); + await vi.advanceTimersByTimeAsync(1_000); + abort.abort(); + await drained; + await vi.advanceTimersByTimeAsync(1_000); + expect(appends()).toBe(0); + } finally { + vi.useRealTimers(); + } + }); + + it("owning subscriber with stopOnAbort:true sends a stop chunk on abort", async () => { + vi.useFakeTimers(); + try { + const { transport, appends } = quietWatchTransport(); + const abort = new AbortController(); + const stream = await transport.reconnectToStream({ + chatId: "chat-own", + abortSignal: abort.signal, + stopOnAbort: true, + }); + const drained = drainChunks(stream!); + await vi.advanceTimersByTimeAsync(1_000); + abort.abort(); + await drained; + await vi.advanceTimersByTimeAsync(1_000); + expect(appends()).toBe(1); + } finally { + vi.useRealTimers(); + } + }); + + it("abortSignal presence alone (stopOnAbort unset) sends no stop", async () => { + vi.useFakeTimers(); + try { + const { transport, appends } = quietWatchTransport(); + const abort = new AbortController(); + const stream = await transport.reconnectToStream({ + chatId: "chat-own", + abortSignal: abort.signal, + }); + const drained = drainChunks(stream!); + await vi.advanceTimersByTimeAsync(1_000); + abort.abort(); + await drained; + await vi.advanceTimersByTimeAsync(1_000); + expect(appends()).toBe(0); + } finally { + vi.useRealTimers(); + } + }); + }); + describe("multi-tab coordination", () => { it("isReadOnly defaults to false when multiTab is disabled", () => { const transport = new TriggerChatTransport({ diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 07e89fc6d9a..34cc2ad2320 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -873,7 +873,11 @@ export class TriggerChatTransport implements ChatTransport { state.isStreaming = true; this.notifySessionChange(chatId, state); - return this.subscribeToSessionStream(state, abortSignal, chatId, { sinceInSeq: inSeq }); + // Owning turn: aborting this live send stops the turn the user drives. + return this.subscribeToSessionStream(state, abortSignal, chatId, { + sinceInSeq: inSeq, + sendStopOnAbort: true, + }); }; /** @@ -1146,6 +1150,13 @@ export class TriggerChatTransport implements ChatTransport { options: { chatId: string; abortSignal?: AbortSignal | undefined; + /** + * Whether aborting this subscription sends `{kind:"stop"}` on `.in`. + * A subscription ending is not session ownership — a passive/watch + * reader unmounting must never stop a turn it doesn't drive. Only + * pass `true` from a caller that owns the live turn. @default false + */ + stopOnAbort?: boolean; } & ChatRequestOptions ): Promise | null> => { const state = this.sessions.get(options.chatId); @@ -1163,7 +1174,7 @@ export class TriggerChatTransport implements ChatTransport { return this.subscribeToSessionStream(state, abortSignal, options.chatId, { resumed: true, - sendStopOnAbort: !!options.abortSignal, + sendStopOnAbort: options.stopOnAbort ?? false, // Reconnect-on-reload opts into the server's settled-peek shortcut // so the SSE doesn't hang for 60s when no turn is in flight. Active // send-a-message paths must keep wait=60 to avoid racing the @@ -1266,7 +1277,11 @@ export class TriggerChatTransport implements ChatTransport { state.isStreaming = true; this.notifySessionChange(chatId, state); - return this.subscribeToSessionStream(state, undefined, chatId, { sinceInSeq: inSeq }); + // Owning action: aborting this send stops the turn the user drives. + return this.subscribeToSessionStream(state, undefined, chatId, { + sinceInSeq: inSeq, + sendStopOnAbort: true, + }); }; // ------------------------------------------------------------------------- From fef7b438f78830c19972dab99988110458a231e2 Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 14:46:29 +0000 Subject: [PATCH 03/10] fix(sdk): surface an error when a watch-mode turn is truncated by reconnect exhaustion --- .changeset/chat-truncated-turn-error.md | 5 +++++ packages/trigger-sdk/src/v3/chat.ts | 14 ++++++++++++++ 2 files changed, 19 insertions(+) create mode 100644 .changeset/chat-truncated-turn-error.md diff --git a/.changeset/chat-truncated-turn-error.md b/.changeset/chat-truncated-turn-error.md new file mode 100644 index 00000000000..f0805d00b8b --- /dev/null +++ b/.changeset/chat-truncated-turn-error.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Chat in the browser now shows an error when a reply is cut off by a lost connection that can't be re-established, instead of presenting the partial reply as if it had finished. diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 34cc2ad2320..fb4b2df3f28 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -1821,6 +1821,20 @@ export class TriggerChatTransport implements ChatTransport { if (opened) return opened; } + // A settled session or an abort ends the turn 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. + if ( + state.isStreaming && + !currentSubscription?.sessionSettled && + !combinedSignal.aborted + ) { + throw new Error( + "Chat stream ended before the turn completed (reconnect budget exhausted)." + ); + } + // Settled close, or the turn is gone — tell the UI instead of // leaving it spinning on a stream nobody will finish. if (state.isStreaming) { From a95973c88d3189a45a01da9bba0f9858978c49ec Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 14:47:03 +0000 Subject: [PATCH 04/10] fix(sdk): jitter the chat reconnect backoff --- packages/trigger-sdk/src/v3/chat.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index fb4b2df3f28..2ab7ede39df 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -1813,7 +1813,10 @@ export class TriggerChatTransport implements ChatTransport { combinedSignal.removeEventListener("abort", done); resolve(); }; - timer = setTimeout(done, Math.min(100 * 2 ** (eofResubscribes - 1), 5_000)); + // Jitter the backoff so many clients reconnecting after the same + // dropped window don't resubscribe in lockstep. + const backoff = Math.min(100 * 2 ** (eofResubscribes - 1), 5_000); + timer = setTimeout(done, backoff * (0.5 + Math.random() * 0.5)); combinedSignal.addEventListener("abort", done); }); if (combinedSignal.aborted) break; From bdf01b005667ba5f7760c4705ca313c98c02a2a2 Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 18:56:35 +0000 Subject: [PATCH 05/10] fix(sdk): keep watch-mode chat streams open across turns and clear state on give-up - reconnect no longer peek-settles in watch mode, so a settled peek between turns can't close the standing subscription before the next turn. - the returned stream now aborts its resubscribe loop when the reader is cancelled, instead of leaking it. - clear and persist isStreaming before the budget-exhaustion throw so a reload doesn't reopen a doomed subscription. --- .changeset/chat-truncated-turn-error.md | 2 ++ packages/trigger-sdk/src/v3/chat.ts | 16 ++++++++++++++-- 2 files changed, 16 insertions(+), 2 deletions(-) diff --git a/.changeset/chat-truncated-turn-error.md b/.changeset/chat-truncated-turn-error.md index f0805d00b8b..959231a8fc8 100644 --- a/.changeset/chat-truncated-turn-error.md +++ b/.changeset/chat-truncated-turn-error.md @@ -3,3 +3,5 @@ --- Chat in the browser now shows an error when a reply is cut off by a lost connection that can't be re-established, instead of presenting the partial reply as if it had finished. + +Watch-mode viewers now keep receiving later turns instead of the stream closing after the first turn completes. diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 2ab7ede39df..9dbeab2aeb3 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -1178,8 +1178,10 @@ export class TriggerChatTransport implements ChatTransport { // Reconnect-on-reload opts into the server's settled-peek shortcut // so the SSE doesn't hang for 60s when no turn is in flight. Active // send-a-message paths must keep wait=60 to avoid racing the - // freshly-triggered turn's first chunk. - peekSettled: true, + // freshly-triggered turn's first chunk. Watch mode must NOT peek: a + // settled peek between turns sets sessionSettled and closes the + // standing subscription, so the viewer never sees the next turn. + peekSettled: !this.watchMode, }); }; @@ -1833,6 +1835,11 @@ export class TriggerChatTransport implements ChatTransport { !currentSubscription?.sessionSettled && !combinedSignal.aborted ) { + // Clear + persist before throwing so the surfaced error leaves + // consistent state — otherwise a reload sees isStreaming: true + // and reopens a doomed subscription. + state.isStreaming = false; + this.notifySessionChange(chatId, state); throw new Error( "Chat stream ended before the turn completed (reconnect budget exhausted)." ); @@ -2039,6 +2046,11 @@ export class TriggerChatTransport implements ChatTransport { this.coordinator?.release(chatId); } }, + // A consumer that stops reading without aborting (drops the reader) + // would otherwise leave the resubscribe loop running forever. + cancel() { + internalAbort.abort(); + }, }); } } From e5ede6f62cd0a8fa0e63776ba3361108429f454e Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 18:56:44 +0000 Subject: [PATCH 06/10] test(sdk): cover watch-mode keepalive, reader-cancel, and budget-exhaustion error --- packages/trigger-sdk/src/v3/chat.test.ts | 100 ++++++++++++++++++++++- 1 file changed, 98 insertions(+), 2 deletions(-) diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index f9851f52ade..09296fc7f22 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1176,7 +1176,7 @@ describe("TriggerChatTransport", () => { expect(transport.getSession("chat-slow")?.isStreaming).toBe(false); }); - it("gives up after a bounded number of resubscribes", async () => { + it("surfaces an error after the resubscribe budget is exhausted", async () => { // Fake timers so the 100ms..1.6s backoffs don't cost real seconds. vi.useFakeTimers(); try { @@ -1205,12 +1205,18 @@ describe("TriggerChatTransport", () => { messages: [createUserMessage("hi")], abortSignal: undefined, }); + // A cut-off turn surfaces an error rather than reading as complete. + // Attach the rejection assertion before advancing timers so the + // rejection is never unhandled. const drained = drainChunks(stream); + const rejects = expect(drained).rejects.toThrow(/reconnect budget exhausted/i); await vi.advanceTimersByTimeAsync(10_000); - await drained; + await rejects; // One initial connect plus the five-attempt resubscribe budget. expect(subscribeCount).toBe(6); + // State is cleared before the throw, so a reload won't reopen a + // doomed subscription. expect(transport.getSession("chat-empty")?.isStreaming).toBe(false); } finally { vi.useRealTimers(); @@ -1260,6 +1266,96 @@ describe("TriggerChatTransport", () => { ]); }); + it("does not peek-settle an idle resubscribe, so the next turn is delivered", async () => { + // Watch mode must NOT send X-Peek-Settled between turns: a settled peek + // while no turn is in flight closes the standing subscription and the + // viewer never sees turn 2. This mock plays the server's peek shortcut — + // a peek request with nothing in flight settles — to prove the transport + // long-polls instead. + const subscribeHeaders: Headers[] = []; + global.fetch = vi.fn().mockImplementation(async (url: string | URL, init?: RequestInit) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionOutSubscribeUrl(urlStr)) { + subscribeHeaders.push(new Headers(init?.headers)); + const n = subscribeHeaders.length; + if (n === 1) { + // Turn 1 completes, then the body EOFs (no settled header). + return defaultSseResponse([ + { type: "text-delta", id: "p1", delta: "turn1" }, + { type: "trigger:turn-complete" }, + ]); + } + if (n === 2) { + // Idle resubscribe. If it peeked, the server settles and the + // subscription would close before turn 2; a long-poll delivers it. + if (init && new Headers(init.headers).get("X-Peek-Settled")) { + return settled(defaultSseResponse([])); + } + return defaultSseResponse([ + { type: "text-delta", id: "p2", delta: "turn2" }, + { type: "trigger:turn-complete" }, + ]); + } + // Turn 2 done — end the watch cleanly. + return settled(defaultSseResponse([])); + } + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-watch-turn2": { publicAccessToken: "p", isStreaming: true } }, + }); + + const stream = await transport.reconnectToStream({ chatId: "chat-watch-turn2" }); + const chunks = await drainChunks(stream!); + + expect(subscribeHeaders[1]?.get("X-Peek-Settled")).toBeNull(); + expect(chunks).toEqual([ + { type: "text-delta", id: "p1", delta: "turn1" }, + { type: "text-delta", id: "p2", delta: "turn2" }, + ]); + }); + + it("cancelling the reader stops the resubscribe loop", async () => { + // A consumer that stops reading without aborting must not leak the + // resubscribe loop — the stream's cancel() aborts it. + vi.useFakeTimers(); + try { + let subscribeCount = 0; + global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionOutSubscribeUrl(urlStr)) { + subscribeCount++; + // Quiet: EOF, no records, never settled — watch keeps resubscribing. + return defaultSseResponse([]); + } + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-watch-cancel": { publicAccessToken: "p", isStreaming: true } }, + }); + + const stream = await transport.reconnectToStream({ chatId: "chat-watch-cancel" }); + const reader = stream!.getReader(); + await vi.advanceTimersByTimeAsync(10_000); + expect(subscribeCount).toBeGreaterThan(1); + + const countAtCancel = subscribeCount; + await reader.cancel(); + await vi.advanceTimersByTimeAsync(10_000); + expect(subscribeCount).toBe(countAtCancel); + } finally { + vi.useRealTimers(); + } + }); + it("stops when the server says the session settled", async () => { let subscribeCount = 0; global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { From c40e5f08c21692a0181ebf09a507bee2a094833b Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Mon, 10 Aug 2026 19:06:28 +0000 Subject: [PATCH 07/10] fix(sdk): guard controller.close() so a clean reader cancel emits no stream-error A consumer cancelling the watch stream aborts the resubscribe loop, which reaches controller.close() on an already-closed controller. The resulting 'Invalid state' throw was surfaced as a bogus stream-error on every clean watch-viewer unmount. Wrap the remaining bare close sites to match the existing pattern. --- packages/trigger-sdk/src/v3/chat.test.ts | 8 +++++++- packages/trigger-sdk/src/v3/chat.ts | 18 +++++++++++++++--- 2 files changed, 22 insertions(+), 4 deletions(-) diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index 09296fc7f22..9b0bfd598a8 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; import type { UIMessage, UIMessageChunk } from "ai"; -import { TriggerChatTransport, createChatTransport } from "./chat.js"; +import { TriggerChatTransport, createChatTransport, type ChatTransportEvent } from "./chat.js"; // ─────────────────────────────────────────────────────────────────────────── // Test helpers @@ -1335,10 +1335,12 @@ describe("TriggerChatTransport", () => { throw new Error(`Unexpected URL: ${urlStr}`); }); + const events: ChatTransportEvent[] = []; const transport = new TriggerChatTransport({ task: "my-chat-task", accessToken: () => "pat", watch: true, + onEvent: (e) => events.push(e), sessions: { "chat-watch-cancel": { publicAccessToken: "p", isStreaming: true } }, }); @@ -1351,6 +1353,10 @@ describe("TriggerChatTransport", () => { await reader.cancel(); await vi.advanceTimersByTimeAsync(10_000); expect(subscribeCount).toBe(countAtCancel); + // A clean cancel must not surface a spurious stream-error (an + // unguarded controller.close() after cancel would throw "Invalid + // state" and leak it onto the telemetry channel). + expect(events.some((e) => e.type === "stream-error")).toBe(false); } finally { vi.useRealTimers(); } diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 9dbeab2aeb3..05242b96faa 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -1864,7 +1864,11 @@ export class TriggerChatTransport implements ChatTransport { const opened = (await openWithAuthRetry()) ?? (await resumeAfterEof()); if (opened === null) { - controller.close(); + try { + controller.close(); + } catch { + /* already closed by a consumer cancel */ + } return; } reader = opened.reader; @@ -1894,7 +1898,11 @@ export class TriggerChatTransport implements ChatTransport { if (next.done) { const resumed = await resumeAfterEof(); if (resumed === null) { - controller.close(); + try { + controller.close(); + } catch { + /* already closed by a consumer cancel */ + } return; } reader = resumed.reader; @@ -1907,7 +1915,11 @@ export class TriggerChatTransport implements ChatTransport { if (combinedSignal.aborted) { internalAbort.abort(); await reader.cancel(); - controller.close(); + try { + controller.close(); + } catch { + /* already closed by a consumer cancel */ + } return; } From d9dc61eddd88242f44bba680034a7eb59b0fda1b Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Tue, 11 Aug 2026 13:14:08 +0000 Subject: [PATCH 08/10] chore(changeset): consolidate the chat keepalive changesets into one --- .changeset/chat-truncated-turn-error.md | 7 ------- .changeset/watch-mode-keepalive.md | 4 +--- 2 files changed, 1 insertion(+), 10 deletions(-) delete mode 100644 .changeset/chat-truncated-turn-error.md diff --git a/.changeset/chat-truncated-turn-error.md b/.changeset/chat-truncated-turn-error.md deleted file mode 100644 index 959231a8fc8..00000000000 --- a/.changeset/chat-truncated-turn-error.md +++ /dev/null @@ -1,7 +0,0 @@ ---- -"@trigger.dev/sdk": patch ---- - -Chat in the browser now shows an error when a reply is cut off by a lost connection that can't be re-established, instead of presenting the partial reply as if it had finished. - -Watch-mode viewers now keep receiving later turns instead of the stream closing after the first turn completes. diff --git a/.changeset/watch-mode-keepalive.md b/.changeset/watch-mode-keepalive.md index e50a04f2a8c..ba6c0441724 100644 --- a/.changeset/watch-mode-keepalive.md +++ b/.changeset/watch-mode-keepalive.md @@ -2,6 +2,4 @@ "@trigger.dev/sdk": patch --- -Watch-mode chat subscriptions now stay connected across quiet periods. - -Read-only chat subscriptions no longer stop a turn when they disconnect. +Watch-mode chat streams now survive quiet windows and keep delivering later turns, and a reply cut off by a lost connection now shows an error instead of appearing finished. From b1814dc70f48e64a72b056838ed807f9c0c74595 Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Tue, 11 Aug 2026 22:04:18 +0000 Subject: [PATCH 09/10] fix(sdk): keep the successor stream's abort controller when an aborted stream tears down --- .changeset/chat-stream-supersede-race.md | 5 + packages/trigger-sdk/src/v3/chat.test.ts | 135 +++++++++++++++++++++++ packages/trigger-sdk/src/v3/chat.ts | 9 +- 3 files changed, 147 insertions(+), 2 deletions(-) create mode 100644 .changeset/chat-stream-supersede-race.md diff --git a/.changeset/chat-stream-supersede-race.md b/.changeset/chat-stream-supersede-race.md new file mode 100644 index 00000000000..bc36a92d4de --- /dev/null +++ b/.changeset/chat-stream-supersede-race.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/sdk": patch +--- + +Fixed a race where quickly restarting a chat stream could break stop and reconnect for the new stream. diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index 9b0bfd598a8..a0c3be9e346 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -132,6 +132,35 @@ function defaultSseResponse( }); } +/** + * An SSE response whose body stays open until the request signal aborts. + * Models a live subscription sitting on a quiet server. + */ +function openSseResponse(signal?: AbortSignal | null): Response { + const body = new ReadableStream({ + start(controller) { + const onAbort = () => { + const err = new Error("aborted"); + err.name = "AbortError"; + try { + controller.error(err); + } catch { + /* already errored */ + } + }; + if (signal?.aborted) onAbort(); + else signal?.addEventListener("abort", onAbort, { once: true }); + }, + }); + return new Response(body, { + status: 200, + headers: { + "content-type": "text/event-stream", + "X-Stream-Version": "v2", + }, + }); +} + function authError(status = 401): Response { return new Response(JSON.stringify({ error: "Unauthorized", name: "TriggerApiError", status }), { status, @@ -1524,6 +1553,112 @@ describe("TriggerChatTransport", () => { }); }); + describe("superseded stream teardown", () => { + it("keeps the successor's controller registered when the aborted stream tears down", async () => { + vi.useFakeTimers(); + try { + let appendCount = 0; + global.fetch = vi.fn().mockImplementation(async (url: string | URL) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionStreamAppendUrl(urlStr)) { + appendCount++; + return defaultAppendResponse(); + } + // Quiet stream: EOF, no records, never settled — watch keeps it open. + if (isSessionOutSubscribeUrl(urlStr)) return defaultSseResponse([]); + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + watch: true, + sessions: { "chat-race": { publicAccessToken: "p", isStreaming: true } }, + }); + + const send = () => + transport.sendMessages({ + trigger: "submit-message" as const, + chatId: "chat-race", + messageId: undefined, + messages: [createUserMessage("hi")], + abortSignal: undefined, + }); + + const first = drainChunks(await send()); + await vi.advanceTimersByTimeAsync(1_000); + + // Supersede: the new stream registers its controller synchronously, + // the aborted one tears down a microtask later. + const second = await send(); + let secondClosed = false; + const secondDrain = drainChunks(second).then(() => { + secondClosed = true; + }); + await first; + await vi.advanceTimersByTimeAsync(1_000); + + // stopGeneration posts the stop chunk either way — only the + // closing assertion proves it found the successor to abort. + appendCount = 0; + expect(await transport.stopGeneration("chat-race")).toBe(true); + await vi.advanceTimersByTimeAsync(1_000); + expect(appendCount).toBe(1); + expect(secondClosed).toBe(true); + + transport.dispose(); + await secondDrain; + } finally { + vi.useRealTimers(); + } + }); + + it("keeps the tab claim the successor took (multi-tab)", async () => { + vi.useFakeTimers(); + try { + global.fetch = vi.fn().mockImplementation(async (url: string | URL, init?: RequestInit) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionStreamAppendUrl(urlStr)) return defaultAppendResponse(); + // Open SSE that only ends when the subscription is aborted, so + // the superseded stream tears down while the successor is live. + if (isSessionOutSubscribeUrl(urlStr)) return openSseResponse(init?.signal); + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + multiTab: true, + sessions: { "chat-race-tab": { publicAccessToken: "p", isStreaming: true } }, + }); + + const send = () => + transport.sendMessages({ + trigger: "submit-message" as const, + chatId: "chat-race-tab", + messageId: undefined, + messages: [createUserMessage("hi")], + abortSignal: undefined, + }); + + const first = drainChunks(await send()); + await vi.advanceTimersByTimeAsync(1_000); + const secondDrain = drainChunks(await send()); + await first; + await vi.advanceTimersByTimeAsync(1_000); + + // The superseded stream must not release the claim its successor + // holds — otherwise this tab flips to read-only mid-turn. + expect(transport.hasClaim("chat-race-tab")).toBe(true); + + transport.dispose(); + await secondDrain; + } finally { + vi.useRealTimers(); + } + }); + }); + describe("multi-tab coordination", () => { it("isReadOnly defaults to false when multiTab is disabled", () => { const transport = new TriggerChatTransport({ diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 05242b96faa..4b0a9ef2268 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -2054,8 +2054,13 @@ export class TriggerChatTransport implements ChatTransport { controller.error(error); } finally { teardownWakeListeners(); - this.activeStreams.delete(chatId); - this.coordinator?.release(chatId); + // Only clear the entry (and drop the tab claim) if it is still + // ours — a superseding send registers its controller before this + // teardown runs, and owns the claim from then on. + if (this.activeStreams.get(chatId) === internalAbort) { + this.activeStreams.delete(chatId); + this.coordinator?.release(chatId); + } } }, // A consumer that stops reading without aborting (drops the reader) From a026776f3fbda88d69e3bed74e5d837c20689d35 Mon Sep 17 00:00:00 2001 From: Katia Bulatova Date: Tue, 11 Aug 2026 22:36:47 +0000 Subject: [PATCH 10/10] fix(sdk): release the tab claim when the user stops generation --- .changeset/chat-stream-supersede-race.md | 2 +- packages/trigger-sdk/src/v3/chat.test.ts | 43 ++++++++++++++++++++++++ packages/trigger-sdk/src/v3/chat.ts | 5 +++ 3 files changed, 49 insertions(+), 1 deletion(-) diff --git a/.changeset/chat-stream-supersede-race.md b/.changeset/chat-stream-supersede-race.md index bc36a92d4de..7e2b807d8c3 100644 --- a/.changeset/chat-stream-supersede-race.md +++ b/.changeset/chat-stream-supersede-race.md @@ -2,4 +2,4 @@ "@trigger.dev/sdk": patch --- -Fixed a race where quickly restarting a chat stream could break stop and reconnect for the new stream. +Fixed a race where quickly restarting a chat stream could break stop and reconnect for the new stream. Stopping a chat now also hands it back to your other tabs instead of leaving them read-only. diff --git a/packages/trigger-sdk/src/v3/chat.test.ts b/packages/trigger-sdk/src/v3/chat.test.ts index a0c3be9e346..71566457dd8 100644 --- a/packages/trigger-sdk/src/v3/chat.test.ts +++ b/packages/trigger-sdk/src/v3/chat.test.ts @@ -1657,6 +1657,49 @@ describe("TriggerChatTransport", () => { vi.useRealTimers(); } }); + + it("releases the tab claim when the user stops generation (multi-tab)", async () => { + vi.useFakeTimers(); + try { + global.fetch = vi.fn().mockImplementation(async (url: string | URL, init?: RequestInit) => { + const urlStr = typeof url === "string" ? url : url.toString(); + if (isSessionStreamAppendUrl(urlStr)) return defaultAppendResponse(); + if (isSessionOutSubscribeUrl(urlStr)) return openSseResponse(init?.signal); + throw new Error(`Unexpected URL: ${urlStr}`); + }); + + const transport = new TriggerChatTransport({ + task: "my-chat-task", + accessToken: () => "pat", + multiTab: true, + sessions: { "chat-stop-tab": { publicAccessToken: "p", isStreaming: true } }, + }); + + const drain = drainChunks( + await transport.sendMessages({ + trigger: "submit-message" as const, + chatId: "chat-stop-tab", + messageId: undefined, + messages: [createUserMessage("hi")], + abortSignal: undefined, + }) + ); + await vi.advanceTimersByTimeAsync(1_000); + expect(transport.hasClaim("chat-stop-tab")).toBe(true); + + expect(await transport.stopGeneration("chat-stop-tab")).toBe(true); + await vi.advanceTimersByTimeAsync(1_000); + + // The turn ends here with no successor stream, so the claim must be + // freed or other tabs stay read-only until this one closes. + expect(transport.hasClaim("chat-stop-tab")).toBe(false); + + transport.dispose(); + await drain; + } finally { + vi.useRealTimers(); + } + }); }); describe("multi-tab coordination", () => { diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index 4b0a9ef2268..b7075e6f71e 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -1218,6 +1218,11 @@ export class TriggerChatTransport implements ChatTransport { activeStream.abort(); this.activeStreams.delete(chatId); } + // Release here, not in the stream teardown: that only releases while it + // still owns the map entry, and we just deleted it. Unlike a supersede, + // no successor stream follows a stop, so the claim would never be freed + // and other tabs would stay read-only until this one closes. + this.coordinator?.release(chatId); // The turn won't reach its turn-complete on this client (we just // aborted the reader), so clear the streaming flag here and persist —