From 10f1d02f534ff353e421ad3dd94d2ee5f9f61d9a Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Thu, 17 Sep 2026 11:23:41 -0700 Subject: [PATCH 1/4] fix: stop retrying failed chat streams indefinitely --- .changeset/quiet-chat-stream-retries.md | 6 + .../v3/apiClient/runStream-retries.test.ts | 106 ++++++++++++ packages/core/src/v3/apiClient/runStream.ts | 14 +- .../trigger-sdk/src/v3/chat-retries.test.ts | 154 ++++++++++++++++++ packages/trigger-sdk/src/v3/chat.ts | 14 +- 5 files changed, 285 insertions(+), 9 deletions(-) create mode 100644 .changeset/quiet-chat-stream-retries.md create mode 100644 packages/core/src/v3/apiClient/runStream-retries.test.ts create mode 100644 packages/trigger-sdk/src/v3/chat-retries.test.ts diff --git a/.changeset/quiet-chat-stream-retries.md b/.changeset/quiet-chat-stream-retries.md new file mode 100644 index 00000000000..98fbc1a4bf9 --- /dev/null +++ b/.changeset/quiet-chat-stream-retries.md @@ -0,0 +1,6 @@ +--- +"@trigger.dev/core": patch +"@trigger.dev/sdk": patch +--- + +Chat streams now stop after five failed connection retries and report a terminal error instead of remaining active indefinitely. Internal timeout exhaustion reports an error, while caller cancellation still closes cleanly. Watch subscriptions continue to retry without a fixed limit. diff --git a/packages/core/src/v3/apiClient/runStream-retries.test.ts b/packages/core/src/v3/apiClient/runStream-retries.test.ts new file mode 100644 index 00000000000..70edf14a49b --- /dev/null +++ b/packages/core/src/v3/apiClient/runStream-retries.test.ts @@ -0,0 +1,106 @@ +import { createServer, type Server, type ServerResponse } from "node:http"; +import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { SSEStreamSubscription } from "./runStream.js"; + +describe("SSE retry exhaustion", () => { + let server: Server; + let url: string; + let abort: AbortController; + let attempts: number; + let respond: (response: ServerResponse) => void; + + beforeEach(async () => { + attempts = 0; + abort = new AbortController(); + server = createServer((_request, response) => { + attempts++; + respond(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"); + url = `http://127.0.0.1:${address.port}`; + }); + + afterEach(async () => { + abort.abort(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + }); + + async function open(options: { fetchTimeoutMs?: number; stallTimeoutMs?: number } = {}) { + return ( + await new SSEStreamSubscription(url, { + signal: abort.signal, + maxRetries: 2, + retryDelayMs: 1, + retryJitter: 0, + ...options, + }).subscribe() + ).getReader(); + } + + it.each(["fetch", "stall"] as const)( + "reports exhausted %s timeouts as failures", + async (failure) => { + respond = (response) => { + if (failure === "stall") { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.flushHeaders(); + } + }; + const reader = await open({ + fetchTimeoutMs: failure === "fetch" ? 100 : 1_000, + stallTimeoutMs: 100, + }); + + await expect(reader.read()).rejects.toMatchObject({ + name: "Error", + message: "Stream connection retries exhausted", + }); + expect(attempts).toBe(3); + } + ); + + it.each([ + ["comment", ": keepalive\n\n"], + ["keepalive event", "event: keepalive\ndata: {}\n\n"], + ["empty batch", 'event: batch\ndata: {"records":[]}\n\n'], + ])("does not reset the retry budget after a %s", async (_name, payload) => { + respond = (response) => { + response.writeHead(200, { + "Content-Type": "text/event-stream", + "X-Stream-Version": "v2", + }); + response.write(payload); + }; + const reader = await open({ stallTimeoutMs: 100 }); + + await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + expect(attempts).toBe(3); + }); + + it("restores the retry budget after a decoded record", async () => { + respond = (response) => { + if (attempts !== 3) { + response.writeHead(503).end(); + return; + } + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.write('id: 1\ndata: {"hello":1}\n\n'); + }; + const reader = await open({ stallTimeoutMs: 100 }); + + expect(await reader.read()).toMatchObject({ done: false, value: { chunk: { hello: 1 } } }); + await expect(reader.read()).rejects.toMatchObject({ status: 503 }); + expect(attempts).toBe(5); + }); + + it("closes without retries when the caller cancels", async () => { + respond = () => abort.abort(); + const reader = await open(); + + expect(await reader.read()).toEqual({ done: true, value: undefined }); + expect(attempts).toBe(1); + }); +}); diff --git a/packages/core/src/v3/apiClient/runStream.ts b/packages/core/src/v3/apiClient/runStream.ts index 5ae5b961fb0..8d3723ca25c 100644 --- a/packages/core/src/v3/apiClient/runStream.ts +++ b/packages/core/src/v3/apiClient/runStream.ts @@ -273,8 +273,7 @@ export class SSEStreamSubscription implements StreamSubscription { // the connection is established, force a reconnect. Catches // silent-dead-socket cases (mobile OS killed the TCP socket but // the read just blocks). Disabled (`0`) by default; opt in - // explicitly. Servers that emit periodic keepalive comments - // reset the timer naturally. + // explicitly. Only decoded records reset the timer. stallTimeoutMs?: number; // HTTP statuses that should NOT be retried — fail the stream // permanently. Defaults cover the permanent client-error set: @@ -461,7 +460,6 @@ export class SSEStreamSubscription implements StreamSubscription { const streamVersion = response.headers.get("X-Stream-Version") ?? "v1"; this.sessionSettled = response.headers.get("X-Session-Settled") === "true"; - this.retryCount = 0; // reset on success armStall(); // Dedup window for record ids. Bounded with FIFO eviction so a @@ -576,8 +574,10 @@ export class SSEStreamSubscription implements StreamSubscription { return; } - armStall(); // any chunk (including server keepalives) resets the silence timer + armStall(); // Each decoded record resets the silence timer. this.authRefreshed = false; + // Headers alone do not establish stream recovery. + this.retryCount = 0; controller.enqueue(value); } } catch (error) { @@ -645,7 +645,11 @@ export class SSEStreamSubscription implements StreamSubscription { } if (this.retryCount >= this.maxRetries) { - const finalError = error || new Error("Max retries reached"); + // Internal timeouts are failures, not caller cancellation. + const finalError = + error?.name === "AbortError" + ? new Error("Stream connection retries exhausted") + : error || new Error("Max retries reached"); controller.error(finalError); this.options.onError?.(finalError); return; diff --git a/packages/trigger-sdk/src/v3/chat-retries.test.ts b/packages/trigger-sdk/src/v3/chat-retries.test.ts new file mode 100644 index 00000000000..8192bd89add --- /dev/null +++ b/packages/trigger-sdk/src/v3/chat-retries.test.ts @@ -0,0 +1,154 @@ +import { createServer, type IncomingMessage, type Server, type ServerResponse } from "node:http"; +import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { createChatTransport, type ChatTransportEvent, type TriggerChatTransport } from "./chat.js"; + +describe("Chat subscription retry exhaustion", () => { + let server: Server; + let baseURL: string; + let transport: TriggerChatTransport; + let attempts: number; + let respond: (response: ServerResponse, request: IncomingMessage) => void; + let events: ChatTransportEvent[]; + + beforeEach(async () => { + attempts = 0; + events = []; + server = createServer((request, response) => { + attempts++; + respond(response, request); + }); + 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 = createChatTransport({ + task: "chat-task", + baseURL, + sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } }, + accessToken: () => "test-token", + onEvent: (event) => events.push(event), + }); + }); + + afterEach(async () => { + transport.dispose(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + }); + + it("clears persisted streaming state after a terminal authorization failure", async () => { + respond = (response) => response.writeHead(401).end(); + const stream = await transport.reconnectToStream({ chatId: "chat" }); + if (!stream) throw new Error("Expected a resumed stream"); + + await expect(stream.getReader().read()).rejects.toMatchObject({ status: 401 }); + expect(attempts).toBe(2); + expect(transport.getSession("chat")?.isStreaming).toBe(false); + expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull(); + expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1); + }); + + it("limits failed connections and reports a terminal stream error", async () => { + respond = (response) => response.writeHead(503).end(); + const stream = await transport.reconnectToStream({ chatId: "chat" }); + if (!stream) throw new Error("Expected a resumed stream"); + + await expect(stream.getReader().read()).rejects.toMatchObject({ status: 503 }); + expect(attempts).toBe(6); + expect(transport.getSession("chat")?.isStreaming).toBe(false); + expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull(); + expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1); + }, 25_000); + + it("preserves unlimited retries for watch subscriptions", async () => { + transport.dispose(); + transport = createChatTransport({ + task: "chat-task", + baseURL, + watch: true, + sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } }, + accessToken: () => "test-token", + }); + respond = (response) => { + if (attempts <= 6) { + response.writeHead(503).end(); + return; + } + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.write('id: 1\ndata: {"type":"start","messageId":"assistant"}\n\n'); + }; + const stream = await transport.reconnectToStream({ chatId: "chat" }); + if (!stream) throw new Error("Expected a resumed stream"); + const reader = stream.getReader(); + + expect(await reader.read()).toMatchObject({ done: false, value: { type: "start" } }); + expect(attempts).toBe(7); + await reader.cancel(); + }, 12_000); + + it.each(["resolve", "reject"] as const)( + "keeps the new stream after a late token refresh: %s", + async (outcome) => { + let releaseToken!: (token: string) => void; + let rejectToken!: (error: Error) => void; + const token = new Promise((resolve, reject) => { + releaseToken = resolve; + rejectToken = reject; + }); + let notifyRefresh!: () => void; + const refreshing = new Promise((resolve) => { + notifyRefresh = resolve; + }); + transport.dispose(); + transport = createChatTransport({ + task: "chat-task", + baseURL, + sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } }, + accessToken: () => { + notifyRefresh(); + return token; + }, + }); + respond = (response, request) => { + if (request.method === "POST") { + response.writeHead(200, { "Content-Type": "application/json" }).end('{"seq_num":50}'); + } else if (attempts === 1 || request.headers.authorization === "Bearer refreshed-token") { + response.writeHead(401).end(); + } else { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.write('id: 51\ndata: {"type":"start","messageId":"replacement"}\n\n'); + } + }; + const oldStream = await transport.reconnectToStream({ chatId: "chat" }); + if (!oldStream) throw new Error("Expected a resumed stream"); + const oldRead = oldStream + .getReader() + .read() + .catch((error: unknown) => error); + await refreshing; + const replacement = await transport.sendMessages({ + chatId: "chat", + trigger: "submit-message", + messageId: "user", + messages: [{ id: "user", role: "user", parts: [{ type: "text", text: "Continue" }] }], + abortSignal: undefined, + }); + const reader = replacement.getReader(); + expect(await reader.read()).toMatchObject({ + done: false, + value: { messageId: "replacement" }, + }); + + if (outcome === "resolve") { + releaseToken("refreshed-token"); + expect(await oldRead).toEqual({ done: true, value: undefined }); + } else { + const error = new Error("Token refresh failed"); + rejectToken(error); + expect(await oldRead).toBe(error); + } + expect(transport.getSession("chat")?.isStreaming).toBe(true); + await reader.cancel(); + } + ); +}); diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index bc352f9e198..b52dca62951 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -2090,10 +2090,11 @@ export class TriggerChatTransport implements ChatTransport { signal: combinedSignal, timeoutInSeconds: this.streamTimeoutSeconds, lastEventId: state.lastEventId, - // Catch silent-dead-socket: if no chunk (or server - // keepalive) arrives in 60s, force reconnect. Sized - // generously over typical agent thinking pauses. + // Reconnect if no decoded record arrives for 60 seconds. stallTimeoutMs: 60_000, + // Normal chat streams must reach a terminal error. Watch subscriptions stay open. + maxRetries: this.watchMode ? Infinity : 5, + retryDelayMs: this.watchMode ? undefined : 1_000, fetchClient: sseFetchClient, }); currentSubscription = subscription; @@ -2163,7 +2164,7 @@ 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) { + if (state.isStreaming && this.activeStreams.get(chatId) === internalAbort) { state.isStreaming = false; this.notifySessionChange(chatId, state); } @@ -2440,6 +2441,11 @@ export class TriggerChatTransport implements ChatTransport { return; } const errorStatus = (error as { status?: unknown }).status; + // A superseded stream cannot settle the replacement stream. + if (this.activeStreams.get(chatId) === internalAbort) { + state.isStreaming = false; + this.notifySessionChange(chatId, state); + } this.emitEvent({ type: "stream-error", chatId, From 8149b9649ea81b954442f69b9fa220b74f227995 Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Thu, 17 Sep 2026 13:24:36 -0700 Subject: [PATCH 2/4] fix: preserve network recovery with a separate stall limit --- .changeset/quiet-chat-stream-retries.md | 2 +- .../v3/apiClient/runStream-retries.test.ts | 104 ++++++++++++++++-- packages/core/src/v3/apiClient/runStream.ts | 16 ++- .../trigger-sdk/src/v3/chat-retries.test.ts | 67 ++++++----- packages/trigger-sdk/src/v3/chat.ts | 14 +-- 5 files changed, 145 insertions(+), 58 deletions(-) diff --git a/.changeset/quiet-chat-stream-retries.md b/.changeset/quiet-chat-stream-retries.md index 98fbc1a4bf9..2f3f247b67b 100644 --- a/.changeset/quiet-chat-stream-retries.md +++ b/.changeset/quiet-chat-stream-retries.md @@ -3,4 +3,4 @@ "@trigger.dev/sdk": patch --- -Chat streams now stop after five failed connection retries and report a terminal error instead of remaining active indefinitely. Internal timeout exhaustion reports an error, while caller cancellation still closes cleanly. Watch subscriptions continue to retry without a fixed limit. +Chat streams now report an error after five retries of a connected stream that sends no records. Network failures and browser wakeups retain automatic recovery. Healthy tool calls with no records for about six minutes also reach this silence limit. Watch subscriptions remain unlimited, and caller cancellation still closes cleanly. diff --git a/packages/core/src/v3/apiClient/runStream-retries.test.ts b/packages/core/src/v3/apiClient/runStream-retries.test.ts index 70edf14a49b..9ec5fd26372 100644 --- a/packages/core/src/v3/apiClient/runStream-retries.test.ts +++ b/packages/core/src/v3/apiClient/runStream-retries.test.ts @@ -8,6 +8,7 @@ describe("SSE retry exhaustion", () => { let abort: AbortController; let attempts: number; let respond: (response: ServerResponse) => void; + let subscription: SSEStreamSubscription; beforeEach(async () => { attempts = 0; @@ -28,16 +29,22 @@ describe("SSE retry exhaustion", () => { await new Promise((resolve) => server.close(() => resolve())); }); - async function open(options: { fetchTimeoutMs?: number; stallTimeoutMs?: number } = {}) { - return ( - await new SSEStreamSubscription(url, { - signal: abort.signal, - maxRetries: 2, - retryDelayMs: 1, - retryJitter: 0, - ...options, - }).subscribe() - ).getReader(); + async function open( + options: { + fetchTimeoutMs?: number; + stallTimeoutMs?: number; + maxRetries?: number; + maxStallRetries?: number; + } = {} + ) { + subscription = new SSEStreamSubscription(url, { + signal: abort.signal, + maxRetries: 2, + retryDelayMs: 1, + retryJitter: 0, + ...options, + }); + return (await subscription.subscribe()).getReader(); } it.each(["fetch", "stall"] as const)( @@ -74,7 +81,7 @@ describe("SSE retry exhaustion", () => { }); response.write(payload); }; - const reader = await open({ stallTimeoutMs: 100 }); + const reader = await open({ stallTimeoutMs: 100, maxRetries: Infinity, maxStallRetries: 2 }); await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); expect(attempts).toBe(3); @@ -103,4 +110,79 @@ describe("SSE retry exhaustion", () => { expect(await reader.read()).toEqual({ done: true, value: undefined }); expect(attempts).toBe(1); }); + + it("limits silent stalls without a general retry limit", async () => { + respond = (response) => { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.flushHeaders(); + }; + const reader = await open({ stallTimeoutMs: 100, maxRetries: Infinity, maxStallRetries: 2 }); + + await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + expect(attempts).toBe(3); + }); + + it("restores the stall budget only after a decoded record", async () => { + respond = (response) => { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.flushHeaders(); + if (attempts === 3) response.write('id: 1\ndata: {"hello":1}\n\n'); + }; + const reader = await open({ stallTimeoutMs: 100, maxRetries: Infinity, maxStallRetries: 2 }); + + expect(await reader.read()).toMatchObject({ done: false, value: { chunk: { hello: 1 } } }); + await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + expect(attempts).toBe(5); + }); + + it.each(["http", "fetch", "body", "wake"] as const)( + "does not charge %s failures to the stall budget", + async (failure) => { + respond = (response) => { + if (attempts === 5) { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.end('id: 1\ndata: {"hello":1}\n\n'); + } else if (failure === "http") { + response.writeHead(503).end(); + } else if (failure !== "fetch") { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.write(": keepalive\n\n"); + setTimeout(() => { + if (failure === "wake") subscription.forceReconnect(); + else response.destroy(); + }, 10); + } + }; + const reader = await open({ + maxRetries: Infinity, + maxStallRetries: 0, + fetchTimeoutMs: 100, + stallTimeoutMs: 1_000, + }); + + expect(await reader.read()).toMatchObject({ done: false, value: { chunk: { hello: 1 } } }); + expect(attempts).toBe(5); + } + ); + + it("retains the stall budget across connection failures and wakeups", async () => { + respond = (response) => { + if (attempts === 2) { + response.writeHead(503).end(); + } else if (attempts !== 4) { + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.flushHeaders(); + if (attempts === 3) setTimeout(() => subscription.forceReconnect(), 10); + } + }; + const reader = await open({ + maxRetries: Infinity, + maxStallRetries: 1, + fetchTimeoutMs: 100, + stallTimeoutMs: 100, + }); + + await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + expect(attempts).toBe(5); + }); }); diff --git a/packages/core/src/v3/apiClient/runStream.ts b/packages/core/src/v3/apiClient/runStream.ts index 8d3723ca25c..ad4721d1638 100644 --- a/packages/core/src/v3/apiClient/runStream.ts +++ b/packages/core/src/v3/apiClient/runStream.ts @@ -219,6 +219,7 @@ export class SSEStreamSubscription implements StreamSubscription { private lastEventId: string | undefined; private from: "beginning" | "latest"; private retryCount = 0; + private stallCount = 0; private maxRetries: number; private retryDelayMs: number; private maxRetryDelayMs: number; @@ -275,6 +276,9 @@ export class SSEStreamSubscription implements StreamSubscription { // the read just blocks). Disabled (`0`) by default; opt in // explicitly. Only decoded records reset the timer. stallTimeoutMs?: number; + // Reconnects after stall timeouts before the stream errors. + // Only decoded records restore this budget. Defaults to Infinity. + maxStallRetries?: number; // HTTP statuses that should NOT be retried — fail the stream // permanently. Defaults cover the permanent client-error set: // `400` (bad request), `404` (stream gone), `409` (conflict), @@ -402,7 +406,11 @@ export class SSEStreamSubscription implements StreamSubscription { const armStall = () => { if (this.stallTimeoutMs <= 0) return; clearTimeout(stallTimer); - stallTimer = setTimeout(() => this.internalAbort?.abort(), this.stallTimeoutMs); + stallTimer = setTimeout(() => { + if (!this.internalAbort || this.internalAbort.signal.aborted) return; + this.stallCount++; + this.internalAbort.abort(); + }, this.stallTimeoutMs); }; // Idempotent — both the catch (before recursion) and the finally @@ -578,6 +586,7 @@ export class SSEStreamSubscription implements StreamSubscription { this.authRefreshed = false; // Headers alone do not establish stream recovery. this.retryCount = 0; + this.stallCount = 0; controller.enqueue(value); } } catch (error) { @@ -644,7 +653,10 @@ export class SSEStreamSubscription implements StreamSubscription { return; } - if (this.retryCount >= this.maxRetries) { + if ( + this.retryCount >= this.maxRetries || + this.stallCount > (this.options.maxStallRetries ?? Infinity) + ) { // Internal timeouts are failures, not caller cancellation. const finalError = error?.name === "AbortError" diff --git a/packages/trigger-sdk/src/v3/chat-retries.test.ts b/packages/trigger-sdk/src/v3/chat-retries.test.ts index 8192bd89add..961befdc4f2 100644 --- a/packages/trigger-sdk/src/v3/chat-retries.test.ts +++ b/packages/trigger-sdk/src/v3/chat-retries.test.ts @@ -48,43 +48,38 @@ describe("Chat subscription retry exhaustion", () => { expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1); }); - it("limits failed connections and reports a terminal stream error", async () => { - respond = (response) => response.writeHead(503).end(); - const stream = await transport.reconnectToStream({ chatId: "chat" }); - if (!stream) throw new Error("Expected a resumed stream"); - - await expect(stream.getReader().read()).rejects.toMatchObject({ status: 503 }); - expect(attempts).toBe(6); - expect(transport.getSession("chat")?.isStreaming).toBe(false); - expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull(); - expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1); - }, 25_000); - - it("preserves unlimited retries for watch subscriptions", async () => { - transport.dispose(); - transport = createChatTransport({ - task: "chat-task", - baseURL, - watch: true, - sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } }, - accessToken: () => "test-token", - }); - respond = (response) => { - if (attempts <= 6) { - response.writeHead(503).end(); - return; - } - response.writeHead(200, { "Content-Type": "text/event-stream" }); - response.write('id: 1\ndata: {"type":"start","messageId":"assistant"}\n\n'); - }; - const stream = await transport.reconnectToStream({ chatId: "chat" }); - if (!stream) throw new Error("Expected a resumed stream"); - const reader = stream.getReader(); + it.each([false, true])( + "recovers after six connection failures (watch: %s)", + async (watch) => { + transport.dispose(); + transport = createChatTransport({ + task: "chat-task", + baseURL, + watch, + sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } }, + accessToken: () => "test-token", + onEvent: (event) => events.push(event), + }); + respond = (response) => { + if (attempts <= 6) { + response.writeHead(503).end(); + return; + } + response.writeHead(200, { "Content-Type": "text/event-stream" }); + response.write('id: 1\ndata: {"type":"start","messageId":"assistant"}\n\n'); + }; + const stream = await transport.reconnectToStream({ chatId: "chat" }); + if (!stream) throw new Error("Expected a resumed stream"); + const reader = stream.getReader(); - expect(await reader.read()).toMatchObject({ done: false, value: { type: "start" } }); - expect(attempts).toBe(7); - await reader.cancel(); - }, 12_000); + expect(await reader.read()).toMatchObject({ done: false, value: { type: "start" } }); + expect(attempts).toBe(7); + expect(transport.getSession("chat")?.isStreaming).toBe(true); + expect(events.filter((event) => event.type === "stream-error")).toHaveLength(0); + await reader.cancel(); + }, + 30_000 + ); it.each(["resolve", "reject"] as const)( "keeps the new stream after a late token refresh: %s", diff --git a/packages/trigger-sdk/src/v3/chat.ts b/packages/trigger-sdk/src/v3/chat.ts index b52dca62951..b3c10316df9 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -2092,9 +2092,12 @@ export class TriggerChatTransport implements ChatTransport { lastEventId: state.lastEventId, // Reconnect if no decoded record arrives for 60 seconds. stallTimeoutMs: 60_000, - // Normal chat streams must reach a terminal error. Watch subscriptions stay open. - maxRetries: this.watchMode ? Infinity : 5, - retryDelayMs: this.watchMode ? undefined : 1_000, + // Bound connected silence while preserving recovery from network failures. + ...(!this.watchMode && { + maxStallRetries: 5, + retryDelayMs: 1_000, + maxRetryDelayMs: 5_000, + }), fetchClient: sseFetchClient, }); currentSubscription = subscription; @@ -2152,11 +2155,6 @@ 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)." ); From 68398fd72e7ac9c3fd0e27b5f797cddf5b75dd5a Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Thu, 17 Sep 2026 14:44:48 -0700 Subject: [PATCH 3/4] fix: distinguish stall exhaustion errors --- .../core/src/v3/apiClient/runStream-retries.test.ts | 8 ++++---- packages/core/src/v3/apiClient/runStream.ts | 13 +++++++------ 2 files changed, 11 insertions(+), 10 deletions(-) diff --git a/packages/core/src/v3/apiClient/runStream-retries.test.ts b/packages/core/src/v3/apiClient/runStream-retries.test.ts index 9ec5fd26372..fea0559187d 100644 --- a/packages/core/src/v3/apiClient/runStream-retries.test.ts +++ b/packages/core/src/v3/apiClient/runStream-retries.test.ts @@ -83,7 +83,7 @@ describe("SSE retry exhaustion", () => { }; const reader = await open({ stallTimeoutMs: 100, maxRetries: Infinity, maxStallRetries: 2 }); - await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + await expect(reader.read()).rejects.toThrow("Stream stalled: no records received"); expect(attempts).toBe(3); }); @@ -118,7 +118,7 @@ describe("SSE retry exhaustion", () => { }; const reader = await open({ stallTimeoutMs: 100, maxRetries: Infinity, maxStallRetries: 2 }); - await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + await expect(reader.read()).rejects.toThrow("Stream stalled: no records received"); expect(attempts).toBe(3); }); @@ -131,7 +131,7 @@ describe("SSE retry exhaustion", () => { const reader = await open({ stallTimeoutMs: 100, maxRetries: Infinity, maxStallRetries: 2 }); expect(await reader.read()).toMatchObject({ done: false, value: { chunk: { hello: 1 } } }); - await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + await expect(reader.read()).rejects.toThrow("Stream stalled: no records received"); expect(attempts).toBe(5); }); @@ -182,7 +182,7 @@ describe("SSE retry exhaustion", () => { stallTimeoutMs: 100, }); - await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted"); + await expect(reader.read()).rejects.toThrow("Stream stalled: no records received"); expect(attempts).toBe(5); }); }); diff --git a/packages/core/src/v3/apiClient/runStream.ts b/packages/core/src/v3/apiClient/runStream.ts index ad4721d1638..4b600fe7f48 100644 --- a/packages/core/src/v3/apiClient/runStream.ts +++ b/packages/core/src/v3/apiClient/runStream.ts @@ -221,6 +221,7 @@ export class SSEStreamSubscription implements StreamSubscription { private retryCount = 0; private stallCount = 0; private maxRetries: number; + private maxStallRetries: number; private retryDelayMs: number; private maxRetryDelayMs: number; private retryJitter: number; @@ -296,6 +297,7 @@ export class SSEStreamSubscription implements StreamSubscription { this.lastEventId = options.lastEventId; this.from = options.from ?? "beginning"; this.maxRetries = options.maxRetries ?? Infinity; + this.maxStallRetries = options.maxStallRetries ?? Infinity; this.retryDelayMs = options.retryDelayMs ?? 100; this.maxRetryDelayMs = options.maxRetryDelayMs ?? 5000; this.retryJitter = options.retryJitter ?? 0.5; @@ -653,13 +655,12 @@ export class SSEStreamSubscription implements StreamSubscription { return; } - if ( - this.retryCount >= this.maxRetries || - this.stallCount > (this.options.maxStallRetries ?? Infinity) - ) { + const stallsExhausted = this.stallCount > this.maxStallRetries; + if (this.retryCount >= this.maxRetries || stallsExhausted) { // Internal timeouts are failures, not caller cancellation. - const finalError = - error?.name === "AbortError" + const finalError = stallsExhausted + ? new Error("Stream stalled: no records received") + : error?.name === "AbortError" ? new Error("Stream connection retries exhausted") : error || new Error("Max retries reached"); controller.error(finalError); From 7825d1b84eb481a298132a0f4c3d05b6ac91f5fb Mon Sep 17 00:00:00 2001 From: Graham Tremper Date: Thu, 17 Sep 2026 14:51:25 -0700 Subject: [PATCH 4/4] docs: name the stall exhaustion error in release notes --- .changeset/quiet-chat-stream-retries.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changeset/quiet-chat-stream-retries.md b/.changeset/quiet-chat-stream-retries.md index 2f3f247b67b..01776a1defa 100644 --- a/.changeset/quiet-chat-stream-retries.md +++ b/.changeset/quiet-chat-stream-retries.md @@ -3,4 +3,4 @@ "@trigger.dev/sdk": patch --- -Chat streams now report an error after five retries of a connected stream that sends no records. Network failures and browser wakeups retain automatic recovery. Healthy tool calls with no records for about six minutes also reach this silence limit. Watch subscriptions remain unlimited, and caller cancellation still closes cleanly. +Chat streams now report `Stream stalled: no records received` after five retries of a connected stream that sends no records. Network failures and browser wakeups retain automatic recovery. Healthy tool calls with no records for about six minutes also reach this silence limit. Watch subscriptions remain unlimited, and caller cancellation still closes cleanly.