Skip to content

Commit fb25c01

Browse files
kathiekiwiTrigger.dev RepoOps
authored andcommitted
fix(core,webapp): stop a failed session.out write from killing the process
fix(core,webapp): stop a failed session.out write from killing the process `StreamsWriterV2` and `SessionStreamInstance` stored the promise from `initializeServerStream()` without a rejection handler. A writer created detached and awaited only later turned an early append failure into a floating rejection; without an `unhandledRejection` listener Node escalated it and the webapp exited, and the same path could fail a run whose append had exhausted its retries. The stored promise now carries a handler while `wait()` still rejects, the webapp logs unhandled rejections next to its `uncaughtException` handler, the chat server logs a dropped session output write instead of swallowing it, and the v1 writer gets the same guard. Regression tests run against a real HTTP/2 fake stream server. Mono-RevId: 4a4e2d30dfec23de72c7e2525dcd0589c6e5f29d
1 parent 35e57e7 commit fb25c01

7 files changed

Lines changed: 152 additions & 5 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/sdk": patch
3+
"@trigger.dev/core": patch
4+
---
5+
6+
A failed write to a realtime or chat session stream no longer crashes the process running it, and a dropped chat session output write is now logged instead of swallowed.

‎apps/webapp/app/entry.server.tsx‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -290,6 +290,17 @@ function logError(error: unknown, request?: Request) {
290290
console.error(error);
291291
}
292292

293+
// Without a listener Node escalates a rejection into the uncaughtException
294+
// handler below, which exits the process. Log it (Sentry via Logger.onError)
295+
// and keep serving. Wrapped in singleton() so Remix's dev-mode CJS reloads
296+
// don't stack duplicate listeners.
297+
singleton("UnhandledRejectionHandler", () => {
298+
process.on("unhandledRejection", (reason) => {
299+
logger.error("unhandledRejection", { error: reason });
300+
});
301+
return true;
302+
});
303+
293304
process.on("uncaughtException", (error, origin) => {
294305
if (
295306
error instanceof Prisma.PrismaClientKnownRequestError ||

‎packages/core/src/v3/realtimeStreams/sessionStreamInstance.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,9 @@ export class SessionStreamInstance<T> implements StreamsWriter {
4040

4141
constructor(private options: SessionStreamInstanceOptions<T>) {
4242
this.streamPromise = this.initializeWriter();
43+
// Same detached-writer guard as `StreamsWriterV2`: the error is still
44+
// surfaced by `wait()` / `stream`.
45+
this.streamPromise.catch(() => {});
4346
}
4447

4548
private async initializeWriter(): Promise<StreamsWriterV2<T>> {

‎packages/core/src/v3/realtimeStreams/streamsWriterV1.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,9 @@ export class StreamsWriterV1<T> implements StreamsWriter {
5454
this.startBuffering();
5555

5656
this.streamPromise = this.initializeServerStream();
57+
// Same detached-writer guard as `StreamsWriterV2`: `wait()` still surfaces
58+
// the error, but an unawaited failure never reaches Node as unhandled.
59+
this.streamPromise.catch(() => {});
5760
}
5861

5962
private generateClientId(): string {

‎packages/core/src/v3/realtimeStreams/streamsWriterV2.test.ts‎

Lines changed: 114 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,10 @@
1+
import http2 from "node:http2";
2+
import net from "node:net";
3+
import type { AddressInfo } from "node:net";
14
import { describe, expect, it } from "vitest";
25

36
import { ChatChunkTooLargeError, isChatChunkTooLargeError } from "../errors.js";
4-
import { encodeChunkOrError } from "./streamsWriterV2.js";
7+
import { StreamsWriterV2, encodeChunkOrError } from "./streamsWriterV2.js";
58

69
// The size cap and discriminant extraction are the only S2-independent bits
710
// of `StreamsWriterV2` that benefit from unit coverage. Both live in the
@@ -75,3 +78,113 @@ describe("isChatChunkTooLargeError", () => {
7578
expect(isChatChunkTooLargeError(undefined)).toBe(false);
7679
});
7780
});
81+
82+
// A session `.out` stream that S2 reports as not found (the head-start drain
83+
// can append before the stream is visible). The writer is detached — nothing
84+
// awaits it until the handover flush — so a rejecting append used to escape
85+
// as an unhandled rejection and take the whole process down.
86+
describe("StreamsWriterV2 against a stream S2 does not have", () => {
87+
it("keeps the append failure on `wait()` instead of emitting an unhandled rejection", async () => {
88+
const server = http2.createServer((_req, res) => {
89+
res.writeHead(404, { "content-type": "application/json" });
90+
res.end(
91+
JSON.stringify({ message: "stream sessions/chat_x/out not found", code: "not_found" })
92+
);
93+
});
94+
const sessions = new Set<http2.ServerHttp2Session>();
95+
server.on("session", (session) => {
96+
sessions.add(session);
97+
session.on("close", () => sessions.delete(session));
98+
});
99+
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", () => resolve()));
100+
const { port } = server.address() as AddressInfo;
101+
102+
const rejections: unknown[] = [];
103+
const onUnhandledRejection = (reason: unknown) => rejections.push(reason);
104+
process.on("unhandledRejection", onUnhandledRejection);
105+
106+
try {
107+
const writer = new StreamsWriterV2<string>({
108+
basin: "trigger-local",
109+
stream: "sessions/chat_x/out",
110+
accessToken: "token",
111+
endpoint: `http://127.0.0.1:${port}/v1`,
112+
source: new ReadableStream<string>({
113+
start(controller) {
114+
controller.enqueue("hello");
115+
controller.close();
116+
},
117+
}),
118+
flushIntervalMs: 5,
119+
});
120+
121+
// Long enough for the append to fail and for Node to flag an unhandled
122+
// rejection, without anything having awaited the writer yet.
123+
await new Promise((resolve) => setTimeout(resolve, 1_000));
124+
125+
expect(rejections).toEqual([]);
126+
await expect(writer.wait()).rejects.toThrow(/not found/);
127+
} finally {
128+
process.off("unhandledRejection", onUnhandledRejection);
129+
for (const session of sessions) session.destroy();
130+
await new Promise<void>((resolve) => server.close(() => resolve()));
131+
}
132+
}, 20_000);
133+
134+
// Same escape route when the SDK's own append retries run out. Reached via
135+
// the retry loop's session-creation branch, not the ack-timeout branch.
136+
it("keeps a retry-exhausted append on `wait()`", async () => {
137+
// Holds the port and refuses every connection, so the retry budget is
138+
// spent on connection failures with no window for another listener.
139+
let connectionAttempts = 0;
140+
const server = net.createServer((socket) => {
141+
connectionAttempts++;
142+
socket.destroy();
143+
});
144+
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", () => resolve()));
145+
const { port } = server.address() as AddressInfo;
146+
147+
const rejections: unknown[] = [];
148+
const onUnhandledRejection = (reason: unknown) => rejections.push(reason);
149+
process.on("unhandledRejection", onUnhandledRejection);
150+
151+
try {
152+
const writer = new StreamsWriterV2<string>({
153+
basin: "trigger-local",
154+
stream: "sessions/chat_x/out",
155+
accessToken: "token",
156+
endpoint: `http://127.0.0.1:${port}/v1`,
157+
source: new ReadableStream<string>({
158+
start(controller) {
159+
controller.enqueue("hello");
160+
controller.close();
161+
},
162+
}),
163+
flushIntervalMs: 5,
164+
});
165+
166+
// Side channel: wait for the budget (3 attempts) to be spent without
167+
// touching the writer, then let the abort propagate.
168+
const deadline = Date.now() + 5_000;
169+
while (connectionAttempts < 3 && Date.now() < deadline) {
170+
await new Promise((resolve) => setTimeout(resolve, 10));
171+
}
172+
expect(connectionAttempts).toBeGreaterThanOrEqual(3);
173+
await new Promise((resolve) => setTimeout(resolve, 100));
174+
175+
expect(rejections).toEqual([]);
176+
// Rejects in a microtask if it already failed; the 0ms timer wins if the
177+
// writer is somehow still pending, so a slow runner fails instead of
178+
// passing vacuously.
179+
await expect(
180+
Promise.race([
181+
writer.wait(),
182+
new Promise((_, reject) => setTimeout(() => reject(new Error("still pending")), 0)),
183+
])
184+
).rejects.toThrow(/Max attempts \(3\) exhausted/);
185+
} finally {
186+
process.off("unhandledRejection", onUnhandledRejection);
187+
await new Promise<void>((resolve) => server.close(() => resolve()));
188+
}
189+
}, 20_000);
190+
});

‎packages/core/src/v3/realtimeStreams/streamsWriterV2.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,10 @@ export class StreamsWriterV2<T = any> implements StreamsWriter {
108108
this.consumerStream = consumerStream;
109109

110110
this.streamPromise = this.initializeServerStream();
111+
// Detached writers (e.g. the head-start drain) call `wait()` long after a
112+
// failed append rejects; without this the rejection is unhandled and Node
113+
// takes the process down. `wait()` still surfaces the error.
114+
this.streamPromise.catch(() => {});
111115
}
112116

113117
private handleAbort(): void {

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

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -749,10 +749,17 @@ async function openHandoverSession(opts: {
749749
if (!sessionWriter) return;
750750
try {
751751
await sessionWriter.wait();
752-
} catch {
753-
// Drop write errors — the customer's response stream is the
754-
// source of truth for what the user sees. Durability/resume
755-
// best-effort.
752+
} catch (error) {
753+
// Dropped, not thrown: the customer's response stream is the source of
754+
// truth for what the user sees. Logged so the loss is diagnosable.
755+
const { message, code, status, origin } = (error ?? {}) as Record<string, unknown>;
756+
console.warn("[chat.handover] session.out write failed", {
757+
chatId,
758+
message,
759+
code,
760+
status,
761+
origin,
762+
});
756763
}
757764
};
758765

0 commit comments

Comments
 (0)