diff --git a/packages/react-router-node/.changes/patch.streaming-client-disconnect.md b/packages/react-router-node/.changes/patch.streaming-client-disconnect.md new file mode 100644 index 0000000000..388e1cfe57 --- /dev/null +++ b/packages/react-router-node/.changes/patch.streaming-client-disconnect.md @@ -0,0 +1 @@ +fix: Prevent client disconnects during streaming from crashing the Node process diff --git a/packages/react-router-node/__tests__/fixtures/stream-cancelled-node-readable.ts b/packages/react-router-node/__tests__/fixtures/stream-cancelled-node-readable.ts new file mode 100644 index 0000000000..0f405e3e23 --- /dev/null +++ b/packages/react-router-node/__tests__/fixtures/stream-cancelled-node-readable.ts @@ -0,0 +1,34 @@ +import { PassThrough, Writable } from "node:stream"; + +import { + createReadableStreamFromReadable, + writeReadableStreamToWritable, +} from "../../stream.ts"; + +let source = new PassThrough(); +let readable = createReadableStreamFromReadable(source); +let writable = new Writable({ + write(_chunk, _encoding, callback) { + callback(); + }, +}); + +source.write(Buffer.from("first chunk")); +let writePromise = writeReadableStreamToWritable(readable, writable); + +await new Promise((resolve) => setImmediate(resolve)); +writable.emit("close"); + +try { + await writePromise; +} catch (error) { + if ( + !(error instanceof Error) || + error.message !== "Writable closed before stream finished" + ) { + throw error; + } +} + +await new Promise((resolve) => setImmediate(resolve)); +console.log("process survived"); diff --git a/packages/react-router-node/__tests__/fixtures/stream-closed-writable.ts b/packages/react-router-node/__tests__/fixtures/stream-closed-writable.ts new file mode 100644 index 0000000000..b58bfcef55 --- /dev/null +++ b/packages/react-router-node/__tests__/fixtures/stream-closed-writable.ts @@ -0,0 +1,35 @@ +import { Writable } from "node:stream"; + +import { writeReadableStreamToWritable } from "../../stream.ts"; + +let controller!: ReadableStreamDefaultController; +let readable = new ReadableStream({ + start(readableController) { + controller = readableController; + }, +}); +let writable = new Writable({ + write(_chunk, _encoding, callback) { + callback(); + }, +}); + +controller.enqueue(new Uint8Array(1)); +let writePromise = writeReadableStreamToWritable(readable, writable); + +await new Promise((resolve) => setImmediate(resolve)); +writable.emit("close"); + +try { + await writePromise; +} catch (error) { + if ( + !(error instanceof Error) || + error.message !== "Writable closed before stream finished" + ) { + throw error; + } +} + +await new Promise((resolve) => setImmediate(resolve)); +console.log("process survived"); diff --git a/packages/react-router-node/__tests__/fixtures/stream-pending-writable-error.ts b/packages/react-router-node/__tests__/fixtures/stream-pending-writable-error.ts new file mode 100644 index 0000000000..a44cf1d7ab --- /dev/null +++ b/packages/react-router-node/__tests__/fixtures/stream-pending-writable-error.ts @@ -0,0 +1,38 @@ +import { Writable } from "node:stream"; + +import { writeReadableStreamToWritable } from "../../stream.ts"; + +let controller!: ReadableStreamDefaultController; +let readable = new ReadableStream({ + start(readableController) { + controller = readableController; + }, +}); + +let finishDestroy!: () => void; +let writable = new Writable({ + write(_chunk, _encoding, callback) { + callback(); + }, + destroy(error, callback) { + finishDestroy = () => callback(error); + }, +}); + +let writePromise = writeReadableStreamToWritable(readable, writable); +let writableError = new Error("Writable failed"); + +// Node marks the writable as destroyed before its destroy callback completes. +writable.destroy(writableError); + +// Let the pending read complete while the writable's error is not yet emitted. +controller.enqueue(new Uint8Array(1)); + +// Wait for the stream pump to observe the destroyed state and clean up. +await writePromise.catch(() => {}); + +// Completing destruction schedules the writable's error after the stream +// monitor's cleanup path has run. +finishDestroy(); +await new Promise((resolve) => setImmediate(resolve)); +console.log("process survived"); diff --git a/packages/react-router-node/__tests__/stream-test.ts b/packages/react-router-node/__tests__/stream-test.ts index 459e80ad88..45e6139788 100644 --- a/packages/react-router-node/__tests__/stream-test.ts +++ b/packages/react-router-node/__tests__/stream-test.ts @@ -2,7 +2,10 @@ * @jest-environment node */ +import { spawnSync } from "node:child_process"; +import path from "node:path"; import { Writable } from "node:stream"; +import { fileURLToPath } from "node:url"; import { writeAsyncIterableToWritable, @@ -45,6 +48,33 @@ function withTimeout(promise: Promise, ms: number): Promise { ); } +function runFixtureProcess(fixtureName: string) { + let fixture = path.join( + path.dirname(fileURLToPath(import.meta.url)), + "fixtures", + fixtureName, + ); + let result = spawnSync( + process.execPath, + ["--experimental-strip-types", "--no-warnings", fixture], + { encoding: "utf8", timeout: 5_000 }, + ); + + return { + status: result.status, + signal: result.signal, + stdout: result.stdout, + stderr: result.stderr, + }; +} + +let survivedProcess = { + status: 0, + signal: null, + stdout: "process survived\n", + stderr: "", +}; + describe("writeReadableStreamToWritable", () => { it("respects writable backpressure", async () => { let highWaterMark = 16; @@ -92,6 +122,24 @@ describe("writeReadableStreamToWritable", () => { "Writable failed", ); }); + + it("does not crash when a destination writable closes mid-stream", () => { + expect(runFixtureProcess("stream-closed-writable.ts")).toEqual( + survivedProcess, + ); + }); + + it("does not crash when a destination close cancels a Node readable", () => { + expect(runFixtureProcess("stream-cancelled-node-readable.ts")).toEqual( + survivedProcess, + ); + }); + + it("does not crash while a destroyed writable has an error pending", () => { + expect(runFixtureProcess("stream-pending-writable-error.ts")).toEqual( + survivedProcess, + ); + }); }); describe("writeAsyncIterableToWritable", () => { diff --git a/packages/react-router-node/stream.ts b/packages/react-router-node/stream.ts index 40086489cf..46f3c977db 100644 --- a/packages/react-router-node/stream.ts +++ b/packages/react-router-node/stream.ts @@ -33,11 +33,13 @@ export async function writeReadableStreamToWritable( } } catch (error: unknown) { try { - reader.cancel(error).catch(() => {}); + reader + .cancel(error instanceof WritableClosedError ? undefined : error) + .catch(() => {}); } catch { // Ignore cancellation errors so we preserve the original write failure. } - writable.destroy(error as Error); + destroyWritable(writable, error as Error); throw error; } finally { writableError.cleanup(); @@ -55,8 +57,24 @@ interface WritableErrorMonitor { throwIfClosed(): void; } +class WritableClosedError extends Error {} + +function destroyWritable(writable: Writable, error: Error) { + if (error instanceof WritableClosedError || writable.destroyed) { + return; + } + + // The write promise carries this error to the caller. Also consume the + // asynchronous error event from destroy so it cannot crash the process + // after the writable error monitor has been cleaned up. + writable.once("error", () => {}); + writable.destroy(error); +} + function monitorWritableError(writable: Writable): WritableErrorMonitor { + let writableStartedDestroyed = writable.destroyed; let settled = false; + let writableErrorEmitted = false; let writableError: Error | undefined; let rejectWritableError!: (error: Error) => void; let writableErrorPromise = new Promise((_, reject) => { @@ -65,6 +83,17 @@ function monitorWritableError(writable: Writable): WritableErrorMonitor { writableErrorPromise.catch(() => {}); function cleanup() { + // `destroy(error)` sets these properties before a potentially async + // destroy callback emits the error, so keep listening during that gap. + if ( + !writableStartedDestroyed && + writable.destroyed && + writable.errored && + !writableErrorEmitted + ) { + return; + } + writable.off("error", onError); writable.off("close", onClose); } @@ -81,11 +110,12 @@ function monitorWritableError(writable: Writable): WritableErrorMonitor { } function onError(error: Error) { + writableErrorEmitted = true; reject(error); } function onClose() { - reject(new Error("Writable closed before stream finished")); + reject(new WritableClosedError("Writable closed before stream finished")); } writable.once("error", onError); @@ -102,7 +132,9 @@ function monitorWritableError(writable: Writable): WritableErrorMonitor { } if (writable.destroyed || writable.writableEnded) { - throw new Error("Cannot write to a destroyed or ended writable stream"); + throw new WritableClosedError( + "Cannot write to a destroyed or ended writable stream", + ); } }, }; @@ -166,7 +198,7 @@ export async function writeAsyncIterableToWritable( // Ignore return errors so we preserve the original write failure. } } - writable.destroy(error); + destroyWritable(writable, error); throw error; } finally { writableError.cleanup();