Skip to content

Commit 10f1d02

Browse files
committed
fix: stop retrying failed chat streams indefinitely
1 parent 8067b1f commit 10f1d02

5 files changed

Lines changed: 285 additions & 9 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
"@trigger.dev/core": patch
3+
"@trigger.dev/sdk": patch
4+
---
5+
6+
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.
Lines changed: 106 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,106 @@
1+
import { createServer, type Server, type ServerResponse } from "node:http";
2+
import { afterEach, beforeEach, describe, expect, it } from "vitest";
3+
import { SSEStreamSubscription } from "./runStream.js";
4+
5+
describe("SSE retry exhaustion", () => {
6+
let server: Server;
7+
let url: string;
8+
let abort: AbortController;
9+
let attempts: number;
10+
let respond: (response: ServerResponse) => void;
11+
12+
beforeEach(async () => {
13+
attempts = 0;
14+
abort = new AbortController();
15+
server = createServer((_request, response) => {
16+
attempts++;
17+
respond(response);
18+
});
19+
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
20+
const address = server.address();
21+
if (!address || typeof address === "string") throw new Error("Expected a TCP address");
22+
url = `http://127.0.0.1:${address.port}`;
23+
});
24+
25+
afterEach(async () => {
26+
abort.abort();
27+
server.closeAllConnections();
28+
await new Promise<void>((resolve) => server.close(() => resolve()));
29+
});
30+
31+
async function open(options: { fetchTimeoutMs?: number; stallTimeoutMs?: number } = {}) {
32+
return (
33+
await new SSEStreamSubscription(url, {
34+
signal: abort.signal,
35+
maxRetries: 2,
36+
retryDelayMs: 1,
37+
retryJitter: 0,
38+
...options,
39+
}).subscribe()
40+
).getReader();
41+
}
42+
43+
it.each(["fetch", "stall"] as const)(
44+
"reports exhausted %s timeouts as failures",
45+
async (failure) => {
46+
respond = (response) => {
47+
if (failure === "stall") {
48+
response.writeHead(200, { "Content-Type": "text/event-stream" });
49+
response.flushHeaders();
50+
}
51+
};
52+
const reader = await open({
53+
fetchTimeoutMs: failure === "fetch" ? 100 : 1_000,
54+
stallTimeoutMs: 100,
55+
});
56+
57+
await expect(reader.read()).rejects.toMatchObject({
58+
name: "Error",
59+
message: "Stream connection retries exhausted",
60+
});
61+
expect(attempts).toBe(3);
62+
}
63+
);
64+
65+
it.each([
66+
["comment", ": keepalive\n\n"],
67+
["keepalive event", "event: keepalive\ndata: {}\n\n"],
68+
["empty batch", 'event: batch\ndata: {"records":[]}\n\n'],
69+
])("does not reset the retry budget after a %s", async (_name, payload) => {
70+
respond = (response) => {
71+
response.writeHead(200, {
72+
"Content-Type": "text/event-stream",
73+
"X-Stream-Version": "v2",
74+
});
75+
response.write(payload);
76+
};
77+
const reader = await open({ stallTimeoutMs: 100 });
78+
79+
await expect(reader.read()).rejects.toThrow("Stream connection retries exhausted");
80+
expect(attempts).toBe(3);
81+
});
82+
83+
it("restores the retry budget after a decoded record", async () => {
84+
respond = (response) => {
85+
if (attempts !== 3) {
86+
response.writeHead(503).end();
87+
return;
88+
}
89+
response.writeHead(200, { "Content-Type": "text/event-stream" });
90+
response.write('id: 1\ndata: {"hello":1}\n\n');
91+
};
92+
const reader = await open({ stallTimeoutMs: 100 });
93+
94+
expect(await reader.read()).toMatchObject({ done: false, value: { chunk: { hello: 1 } } });
95+
await expect(reader.read()).rejects.toMatchObject({ status: 503 });
96+
expect(attempts).toBe(5);
97+
});
98+
99+
it("closes without retries when the caller cancels", async () => {
100+
respond = () => abort.abort();
101+
const reader = await open();
102+
103+
expect(await reader.read()).toEqual({ done: true, value: undefined });
104+
expect(attempts).toBe(1);
105+
});
106+
});

packages/core/src/v3/apiClient/runStream.ts

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -273,8 +273,7 @@ export class SSEStreamSubscription implements StreamSubscription {
273273
// the connection is established, force a reconnect. Catches
274274
// silent-dead-socket cases (mobile OS killed the TCP socket but
275275
// the read just blocks). Disabled (`0`) by default; opt in
276-
// explicitly. Servers that emit periodic keepalive comments
277-
// reset the timer naturally.
276+
// explicitly. Only decoded records reset the timer.
278277
stallTimeoutMs?: number;
279278
// HTTP statuses that should NOT be retried — fail the stream
280279
// permanently. Defaults cover the permanent client-error set:
@@ -461,7 +460,6 @@ export class SSEStreamSubscription implements StreamSubscription {
461460

462461
const streamVersion = response.headers.get("X-Stream-Version") ?? "v1";
463462
this.sessionSettled = response.headers.get("X-Session-Settled") === "true";
464-
this.retryCount = 0; // reset on success
465463
armStall();
466464

467465
// Dedup window for record ids. Bounded with FIFO eviction so a
@@ -576,8 +574,10 @@ export class SSEStreamSubscription implements StreamSubscription {
576574
return;
577575
}
578576

579-
armStall(); // any chunk (including server keepalives) resets the silence timer
577+
armStall(); // Each decoded record resets the silence timer.
580578
this.authRefreshed = false;
579+
// Headers alone do not establish stream recovery.
580+
this.retryCount = 0;
581581
controller.enqueue(value);
582582
}
583583
} catch (error) {
@@ -645,7 +645,11 @@ export class SSEStreamSubscription implements StreamSubscription {
645645
}
646646

647647
if (this.retryCount >= this.maxRetries) {
648-
const finalError = error || new Error("Max retries reached");
648+
// Internal timeouts are failures, not caller cancellation.
649+
const finalError =
650+
error?.name === "AbortError"
651+
? new Error("Stream connection retries exhausted")
652+
: error || new Error("Max retries reached");
649653
controller.error(finalError);
650654
this.options.onError?.(finalError);
651655
return;
Lines changed: 154 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,154 @@
1+
import { createServer, type IncomingMessage, type Server, type ServerResponse } from "node:http";
2+
import { afterEach, beforeEach, describe, expect, it } from "vitest";
3+
import { createChatTransport, type ChatTransportEvent, type TriggerChatTransport } from "./chat.js";
4+
5+
describe("Chat subscription retry exhaustion", () => {
6+
let server: Server;
7+
let baseURL: string;
8+
let transport: TriggerChatTransport;
9+
let attempts: number;
10+
let respond: (response: ServerResponse, request: IncomingMessage) => void;
11+
let events: ChatTransportEvent[];
12+
13+
beforeEach(async () => {
14+
attempts = 0;
15+
events = [];
16+
server = createServer((request, response) => {
17+
attempts++;
18+
respond(response, request);
19+
});
20+
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
21+
const address = server.address();
22+
if (!address || typeof address === "string") throw new Error("Expected a TCP address");
23+
baseURL = `http://127.0.0.1:${address.port}`;
24+
transport = createChatTransport({
25+
task: "chat-task",
26+
baseURL,
27+
sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } },
28+
accessToken: () => "test-token",
29+
onEvent: (event) => events.push(event),
30+
});
31+
});
32+
33+
afterEach(async () => {
34+
transport.dispose();
35+
server.closeAllConnections();
36+
await new Promise<void>((resolve) => server.close(() => resolve()));
37+
});
38+
39+
it("clears persisted streaming state after a terminal authorization failure", async () => {
40+
respond = (response) => response.writeHead(401).end();
41+
const stream = await transport.reconnectToStream({ chatId: "chat" });
42+
if (!stream) throw new Error("Expected a resumed stream");
43+
44+
await expect(stream.getReader().read()).rejects.toMatchObject({ status: 401 });
45+
expect(attempts).toBe(2);
46+
expect(transport.getSession("chat")?.isStreaming).toBe(false);
47+
expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull();
48+
expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1);
49+
});
50+
51+
it("limits failed connections and reports a terminal stream error", async () => {
52+
respond = (response) => response.writeHead(503).end();
53+
const stream = await transport.reconnectToStream({ chatId: "chat" });
54+
if (!stream) throw new Error("Expected a resumed stream");
55+
56+
await expect(stream.getReader().read()).rejects.toMatchObject({ status: 503 });
57+
expect(attempts).toBe(6);
58+
expect(transport.getSession("chat")?.isStreaming).toBe(false);
59+
expect(await transport.reconnectToStream({ chatId: "chat" })).toBeNull();
60+
expect(events.filter((event) => event.type === "stream-error")).toHaveLength(1);
61+
}, 25_000);
62+
63+
it("preserves unlimited retries for watch subscriptions", async () => {
64+
transport.dispose();
65+
transport = createChatTransport({
66+
task: "chat-task",
67+
baseURL,
68+
watch: true,
69+
sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } },
70+
accessToken: () => "test-token",
71+
});
72+
respond = (response) => {
73+
if (attempts <= 6) {
74+
response.writeHead(503).end();
75+
return;
76+
}
77+
response.writeHead(200, { "Content-Type": "text/event-stream" });
78+
response.write('id: 1\ndata: {"type":"start","messageId":"assistant"}\n\n');
79+
};
80+
const stream = await transport.reconnectToStream({ chatId: "chat" });
81+
if (!stream) throw new Error("Expected a resumed stream");
82+
const reader = stream.getReader();
83+
84+
expect(await reader.read()).toMatchObject({ done: false, value: { type: "start" } });
85+
expect(attempts).toBe(7);
86+
await reader.cancel();
87+
}, 12_000);
88+
89+
it.each(["resolve", "reject"] as const)(
90+
"keeps the new stream after a late token refresh: %s",
91+
async (outcome) => {
92+
let releaseToken!: (token: string) => void;
93+
let rejectToken!: (error: Error) => void;
94+
const token = new Promise<string>((resolve, reject) => {
95+
releaseToken = resolve;
96+
rejectToken = reject;
97+
});
98+
let notifyRefresh!: () => void;
99+
const refreshing = new Promise<void>((resolve) => {
100+
notifyRefresh = resolve;
101+
});
102+
transport.dispose();
103+
transport = createChatTransport({
104+
task: "chat-task",
105+
baseURL,
106+
sessions: { chat: { publicAccessToken: "test-token", isStreaming: true } },
107+
accessToken: () => {
108+
notifyRefresh();
109+
return token;
110+
},
111+
});
112+
respond = (response, request) => {
113+
if (request.method === "POST") {
114+
response.writeHead(200, { "Content-Type": "application/json" }).end('{"seq_num":50}');
115+
} else if (attempts === 1 || request.headers.authorization === "Bearer refreshed-token") {
116+
response.writeHead(401).end();
117+
} else {
118+
response.writeHead(200, { "Content-Type": "text/event-stream" });
119+
response.write('id: 51\ndata: {"type":"start","messageId":"replacement"}\n\n');
120+
}
121+
};
122+
const oldStream = await transport.reconnectToStream({ chatId: "chat" });
123+
if (!oldStream) throw new Error("Expected a resumed stream");
124+
const oldRead = oldStream
125+
.getReader()
126+
.read()
127+
.catch((error: unknown) => error);
128+
await refreshing;
129+
const replacement = await transport.sendMessages({
130+
chatId: "chat",
131+
trigger: "submit-message",
132+
messageId: "user",
133+
messages: [{ id: "user", role: "user", parts: [{ type: "text", text: "Continue" }] }],
134+
abortSignal: undefined,
135+
});
136+
const reader = replacement.getReader();
137+
expect(await reader.read()).toMatchObject({
138+
done: false,
139+
value: { messageId: "replacement" },
140+
});
141+
142+
if (outcome === "resolve") {
143+
releaseToken("refreshed-token");
144+
expect(await oldRead).toEqual({ done: true, value: undefined });
145+
} else {
146+
const error = new Error("Token refresh failed");
147+
rejectToken(error);
148+
expect(await oldRead).toBe(error);
149+
}
150+
expect(transport.getSession("chat")?.isStreaming).toBe(true);
151+
await reader.cancel();
152+
}
153+
);
154+
});

packages/trigger-sdk/src/v3/chat.ts

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2090,10 +2090,11 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
20902090
signal: combinedSignal,
20912091
timeoutInSeconds: this.streamTimeoutSeconds,
20922092
lastEventId: state.lastEventId,
2093-
// Catch silent-dead-socket: if no chunk (or server
2094-
// keepalive) arrives in 60s, force reconnect. Sized
2095-
// generously over typical agent thinking pauses.
2093+
// Reconnect if no decoded record arrives for 60 seconds.
20962094
stallTimeoutMs: 60_000,
2095+
// Normal chat streams must reach a terminal error. Watch subscriptions stay open.
2096+
maxRetries: this.watchMode ? Infinity : 5,
2097+
retryDelayMs: this.watchMode ? undefined : 1_000,
20972098
fetchClient: sseFetchClient,
20982099
});
20992100
currentSubscription = subscription;
@@ -2163,7 +2164,7 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
21632164

21642165
// Settled close, or the turn is gone — tell the UI instead of
21652166
// leaving it spinning on a stream nobody will finish.
2166-
if (state.isStreaming) {
2167+
if (state.isStreaming && this.activeStreams.get(chatId) === internalAbort) {
21672168
state.isStreaming = false;
21682169
this.notifySessionChange(chatId, state);
21692170
}
@@ -2440,6 +2441,11 @@ export class TriggerChatTransport implements ChatTransport<UIMessage> {
24402441
return;
24412442
}
24422443
const errorStatus = (error as { status?: unknown }).status;
2444+
// A superseded stream cannot settle the replacement stream.
2445+
if (this.activeStreams.get(chatId) === internalAbort) {
2446+
state.isStreaming = false;
2447+
this.notifySessionChange(chatId, state);
2448+
}
24432449
this.emitEvent({
24442450
type: "stream-error",
24452451
chatId,

0 commit comments

Comments
 (0)