Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Next Next commit
fix: stop retrying failed chat streams indefinitely
  • Loading branch information
gtremper committed Sep 17, 2026
commit 10f1d02f534ff353e421ad3dd94d2ee5f9f61d9a
6 changes: 6 additions & 0 deletions .changeset/quiet-chat-stream-retries.md
Original file line number Diff line number Diff line change
@@ -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.
106 changes: 106 additions & 0 deletions packages/core/src/v3/apiClient/runStream-retries.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>((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<void>((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);
});
});
14 changes: 9 additions & 5 deletions packages/core/src/v3/apiClient/runStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
Expand Down
154 changes: 154 additions & 0 deletions packages/trigger-sdk/src/v3/chat-retries.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>((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<void>((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<string>((resolve, reject) => {
releaseToken = resolve;
rejectToken = reject;
});
let notifyRefresh!: () => void;
const refreshing = new Promise<void>((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();
}
);
});
14 changes: 10 additions & 4 deletions packages/trigger-sdk/src/v3/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2090,10 +2090,11 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
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;
Expand Down Expand Up @@ -2163,7 +2164,7 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {

// 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);
}
Expand Down Expand Up @@ -2440,6 +2441,11 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
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,
Expand Down
Loading