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
Prev Previous commit
Next Next commit
fix: preserve network recovery with a separate stall limit
  • Loading branch information
gtremper committed Sep 17, 2026
commit 8149b9649ea81b954442f69b9fa220b74f227995
2 changes: 1 addition & 1 deletion .changeset/quiet-chat-stream-retries.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
104 changes: 93 additions & 11 deletions packages/core/src/v3/apiClient/runStream-retries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -28,16 +29,22 @@ describe("SSE retry exhaustion", () => {
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();
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)(
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
});
});
16 changes: 14 additions & 2 deletions packages/core/src/v3/apiClient/runStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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"
Expand Down
67 changes: 31 additions & 36 deletions packages/trigger-sdk/src/v3/chat-retries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
14 changes: 6 additions & 8 deletions packages/trigger-sdk/src/v3/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2092,9 +2092,12 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
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;
Expand Down Expand Up @@ -2152,11 +2155,6 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
!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)."
);
Expand Down
Loading