Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
fix: Prevent client disconnects during streaming from crashing the Node process
Original file line number Diff line number Diff line change
@@ -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");
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
import { Writable } from "node:stream";

import { writeReadableStreamToWritable } from "../../stream.ts";

let controller!: ReadableStreamDefaultController<Uint8Array>;
let readable = new ReadableStream<Uint8Array>({
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");
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
import { Writable } from "node:stream";

import { writeReadableStreamToWritable } from "../../stream.ts";

let controller!: ReadableStreamDefaultController<Uint8Array>;
let readable = new ReadableStream<Uint8Array>({
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");
48 changes: 48 additions & 0 deletions packages/react-router-node/__tests__/stream-test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -45,6 +48,33 @@ function withTimeout<T>(promise: Promise<T>, ms: number): Promise<T> {
);
}

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;
Expand Down Expand Up @@ -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", () => {
Expand Down
42 changes: 37 additions & 5 deletions packages/react-router-node/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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<never>((_, reject) => {
Expand All @@ -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);
}
Expand All @@ -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);
Expand All @@ -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",
);
}
},
};
Expand Down Expand Up @@ -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();
Expand Down