diff --git a/.changeset/quiet-chat-stream-retries.md b/.changeset/quiet-chat-stream-retries.md new file mode 100644 index 00000000000..01776a1defa --- /dev/null +++ b/.changeset/quiet-chat-stream-retries.md @@ -0,0 +1,6 @@ +--- +"@trigger.dev/core": patch +"@trigger.dev/sdk": patch +--- + +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. 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..fea0559187d --- /dev/null +++ b/packages/core/src/v3/apiClient/runStream-retries.test.ts @@ -0,0 +1,188 @@ +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; + let subscription: SSEStreamSubscription; + + 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; + 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)( + "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, maxRetries: Infinity, maxStallRetries: 2 }); + + await expect(reader.read()).rejects.toThrow("Stream stalled: no records received"); + 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); + }); + + 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 stalled: no records received"); + 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 stalled: no records received"); + 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 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 5ae5b961fb0..4b600fe7f48 100644 --- a/packages/core/src/v3/apiClient/runStream.ts +++ b/packages/core/src/v3/apiClient/runStream.ts @@ -219,7 +219,9 @@ export class SSEStreamSubscription implements StreamSubscription { private lastEventId: string | undefined; private from: "beginning" | "latest"; private retryCount = 0; + private stallCount = 0; private maxRetries: number; + private maxStallRetries: number; private retryDelayMs: number; private maxRetryDelayMs: number; private retryJitter: number; @@ -273,9 +275,11 @@ 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; + // 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), @@ -293,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; @@ -403,7 +408,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 @@ -461,7 +470,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 +584,11 @@ 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; + this.stallCount = 0; controller.enqueue(value); } } catch (error) { @@ -644,8 +655,14 @@ export class SSEStreamSubscription implements StreamSubscription { return; } - if (this.retryCount >= this.maxRetries) { - const finalError = error || new Error("Max retries reached"); + const stallsExhausted = this.stallCount > this.maxStallRetries; + if (this.retryCount >= this.maxRetries || stallsExhausted) { + // Internal timeouts are failures, not caller cancellation. + 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); 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..961befdc4f2 --- /dev/null +++ b/packages/trigger-sdk/src/v3/chat-retries.test.ts @@ -0,0 +1,149 @@ +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.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); + 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", + 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..b3c10316df9 100644 --- a/packages/trigger-sdk/src/v3/chat.ts +++ b/packages/trigger-sdk/src/v3/chat.ts @@ -2090,10 +2090,14 @@ 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, + // Bound connected silence while preserving recovery from network failures. + ...(!this.watchMode && { + maxStallRetries: 5, + retryDelayMs: 1_000, + maxRetryDelayMs: 5_000, + }), fetchClient: sseFetchClient, }); currentSubscription = subscription; @@ -2151,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)." ); @@ -2163,7 +2162,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 +2439,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,