diff --git a/CHANGELOG.md b/CHANGELOG.md index de2fa1fe14..b935013783 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -161,7 +161,7 @@ Node apps opting-into the Web Streams API _might_ even see a small performance b ### Unstable Changes -⚠️ _[Unstable features](https://reactrouter.com/community/api-development-strategy#unstable-flags) are not recommended for production use_ +⚠️ _[Unstable features](https://reactrouter.com/community/api-development-strategy#unstable-flags) are not recommended for production use_ - `@react-router/dev` - Add the [`future.unstable_enableNodeReadableStream`](https://reactrouter.com/upgrading/future#futureunstable_enablenodereadablestream) flag to opt Node Framework mode apps into using `renderToReadableStream` instead of `renderToPipeableStream` ([#15290](https://github.com/remix-run/react-router/pull/15290)) - This flag has no effect if you have your own `entry.server.tsx` diff --git a/contributors.yml b/contributors.yml index f6f3163a20..b0d86fab06 100644 --- a/contributors.yml +++ b/contributors.yml @@ -281,6 +281,7 @@ - m-kawafuji - m-shojaei - machour +- MahinAnowar - majamarijan - Malien - Manc diff --git a/packages/react-router-dev/CHANGELOG.md b/packages/react-router-dev/CHANGELOG.md index f2032f1de7..886f5ded0c 100644 --- a/packages/react-router-dev/CHANGELOG.md +++ b/packages/react-router-dev/CHANGELOG.md @@ -16,7 +16,7 @@ ### Unstable Changes -⚠️ _[Unstable features](https://reactrouter.com/community/api-development-strategy#unstable-flags) are not recommended for production use_ +⚠️ _[Unstable features](https://reactrouter.com/community/api-development-strategy#unstable-flags) are not recommended for production use_ - Add the [`future.unstable_enableNodeReadableStream`](https://reactrouter.com/upgrading/future#futureunstable_enablenodereadablestream) flag to opt Node Framework mode apps into using `renderToReadableStream` instead of `renderToPipeableStream` ([#15290](https://github.com/remix-run/react-router/pull/15290)) - This flag has no effect if you have your own `entry.server.tsx` diff --git a/packages/react-router/.changes/patch.rsc-html-stream-cancel.md b/packages/react-router/.changes/patch.rsc-html-stream-cancel.md new file mode 100644 index 0000000000..322d3e5488 --- /dev/null +++ b/packages/react-router/.changes/patch.rsc-html-stream-cancel.md @@ -0,0 +1 @@ +Fix server crash (`TypeError: Invalid state: Unable to enqueue`) when a request is aborted while the RSC HTML stream has a pending flush — `injectRSCPayload` now handles cancellation of its readable side, clears the pending flush, and cancels the underlying RSC payload stream diff --git a/packages/react-router/__tests__/rsc/html-stream-test.ts b/packages/react-router/__tests__/rsc/html-stream-test.ts new file mode 100644 index 0000000000..4999fa0a4a --- /dev/null +++ b/packages/react-router/__tests__/rsc/html-stream-test.ts @@ -0,0 +1,222 @@ +import { injectRSCPayload } from "../../lib/rsc/html-stream/server"; +import { routeRSCServerRequest } from "../../lib/rsc/server.ssr"; + +const encoder = new TextEncoder(); +const decoder = new TextDecoder(); + +function createDeferred() { + let resolve!: (value: T | PromiseLike) => void; + let reject!: (reason?: unknown) => void; + let promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} + +function createRSCStream({ + keepOpen = false, + chunks = ['S1:"hello"'], +}: { + keepOpen?: boolean; + chunks?: string[]; +} = {}) { + let cancelled = false; + let cancelReason: unknown; + let stream = new ReadableStream({ + start(controller) { + chunks.forEach((chunk) => controller.enqueue(encoder.encode(chunk))); + if (!keepOpen) { + controller.close(); + } + }, + cancel(reason) { + cancelled = true; + cancelReason = reason; + }, + }); + return { + stream, + isCancelled: () => cancelled, + cancelReason: () => cancelReason, + }; +} + +async function withUnhandledRejections(run: () => Promise) { + let unhandledRejections: unknown[] = []; + let onUnhandledRejection = (reason: unknown) => + unhandledRejections.push(reason); + process.on("unhandledRejection", onUnhandledRejection); + try { + await run(); + } finally { + process.off("unhandledRejection", onUnhandledRejection); + } + return unhandledRejections; +} + +function tick() { + return new Promise((resolve) => setTimeout(resolve, 20)); +} + +async function withTimeout(promise: Promise, message: string) { + let timeout: ReturnType; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timeout = setTimeout(() => reject(new Error(message)), 1000); + }), + ]); + } finally { + clearTimeout(timeout!); + } +} + +async function readStream(stream: ReadableStream) { + let reader = stream.getReader(); + let chunks: Uint8Array[] = []; + + while (true) { + let { done, value } = await reader.read(); + if (done) { + break; + } + chunks.push(value); + } + + let length = chunks.reduce((sum, chunk) => sum + chunk.length, 0); + let merged = new Uint8Array(length); + let offset = 0; + for (let chunk of chunks) { + merged.set(chunk, offset); + offset += chunk.length; + } + + return decoder.decode(merged); +} + +describe("injectRSCPayload", () => { + it("streams buffered HTML, RSC payload chunks, and the HTML trailer", async () => { + let html = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("")); + controller.enqueue(encoder.encode("hi")); + controller.close(); + }, + }); + + let result = await readStream( + html.pipeThrough(injectRSCPayload(createRSCStream().stream)), + ); + + expect(result).toBe( + 'hi', + ); + }); + + it("does not crash when the readable side is cancelled while a flush is pending", async () => { + let rsc = createRSCStream({ keepOpen: true }); + let transform = injectRSCPayload(rsc.stream); + let writer = transform.writable.getWriter(); + let reader = transform.readable.getReader(); + let reason = new Error("client aborted"); + + let unhandledRejections = await withUnhandledRejections(async () => { + // Schedule the buffered flush (`setTimeout(..., 0)`) by writing a chunk, + // then cancel the readable side (client aborted the request) before the + // timer fires. + writer + .write(encoder.encode("hi")) + .catch(() => {}); + await reader.cancel(reason); + + // Let the pending flush timer fire. + await tick(); + }); + + // Without a cancel handler the pending flush enqueues into a cancelled + // stream, and the rejection from the timer callback kills the process. + expect(unhandledRejections).toEqual([]); + expect(rsc.isCancelled()).toBe(true); + expect(rsc.cancelReason()).toBe(reason); + }); + + it("does not crash when the readable side is cancelled while the RSC payload is streaming", async () => { + let rsc = createRSCStream({ keepOpen: true }); + let transform = injectRSCPayload(rsc.stream); + let writer = transform.writable.getWriter(); + let reader = transform.readable.getReader(); + let reason = new Error("client aborted"); + + let unhandledRejections = await withUnhandledRejections(async () => { + writer + .write(encoder.encode("hi")) + .catch(() => {}); + // Read the flushed HTML and the first RSC script chunk so the RSC + // payload stream is being consumed, then abort. + await reader.read(); + await reader.read(); + await reader.cancel(reason); + + await tick(); + }); + + expect(unhandledRejections).toEqual([]); + expect(rsc.isCancelled()).toBe(true); + expect(rsc.cancelReason()).toBe(reason); + }); +}); + +describe("routeRSCServerRequest", () => { + it("does not crash when an RSC Framework document response is cancelled while payload injection has a pending flush", async () => { + let htmlPulled = createDeferred(); + let htmlCancelled = createDeferred(); + let response = await routeRSCServerRequest({ + request: new Request("https://remix.run/"), + serverResponse: new Response(createRSCStream().stream), + createFromReadableStream: async (body) => { + await readStream(body); + return { type: "render" } as never; + }, + async renderHTML(getPayload) { + await getPayload(); + let sent = false; + return new ReadableStream({ + pull(controller) { + if (sent) { + return; + } + sent = true; + controller.enqueue(encoder.encode("hi")); + htmlPulled.resolve(); + }, + cancel(reason) { + htmlCancelled.resolve(reason); + }, + }); + }, + }); + let reader = response.body!.getReader(); + let reason = new Error("client aborted"); + + let unhandledRejections = await withUnhandledRejections(async () => { + let read = reader.read().catch(() => {}); + + await htmlPulled.promise; + await Promise.resolve(); + await withTimeout( + reader.cancel(reason), + "Timed out cancelling document response body", + ); + await withTimeout(read, "Timed out settling pending document body read"); + + await tick(); + }); + + expect(unhandledRejections).toEqual([]); + await expect( + withTimeout(htmlCancelled.promise, "Timed out cancelling HTML stream"), + ).resolves.toBe(reason); + }); +}); diff --git a/packages/react-router/lib/rsc/html-stream/server.ts b/packages/react-router/lib/rsc/html-stream/server.ts index 38cd5bdb3a..e418aa5c00 100644 --- a/packages/react-router/lib/rsc/html-stream/server.ts +++ b/packages/react-router/lib/rsc/html-stream/server.ts @@ -10,6 +10,8 @@ export function injectRSCPayload(rscStream: ReadableStream) { (resolve) => (resolveFlightDataPromise = resolve), ); let startedRSC = false; + let cancelled = false; + let rscReader: ReadableStreamDefaultReader | null = null; // Buffer all HTML chunks enqueued during the current tick of the event loop (roughly) // and write them to the output stream all at once. This ensures that we don't generate @@ -31,7 +33,9 @@ export function injectRSCPayload(rscStream: ReadableStream) { timeout = null; } - return new TransformStream({ + let transformer: Transformer & { + cancel?: (reason: unknown) => void | Promise; + } = { transform(chunk, controller) { buffered.push(chunk); if (timeout) { @@ -39,10 +43,16 @@ export function injectRSCPayload(rscStream: ReadableStream) { } timeout = setTimeout(async () => { + // The readable side may have been cancelled (e.g., the client aborted + // the request) while this flush was pending — enqueueing would throw. + if (cancelled) { + return; + } flushBufferedChunks(controller); if (!startedRSC) { startedRSC = true; - writeRSCStream(rscStream, controller) + rscReader = rscStream.getReader(); + writeRSCStream(rscReader, controller, () => cancelled) .catch((err) => controller.error(err)) .then(resolveFlightDataPromise); } @@ -56,18 +66,36 @@ export function injectRSCPayload(rscStream: ReadableStream) { } controller.enqueue(encoder.encode("")); }, - }); + async cancel(reason) { + cancelled = true; + if (timeout) { + clearTimeout(timeout); + timeout = null; + } + buffered.length = 0; + if (rscReader) { + await rscReader.cancel(reason).catch(() => {}); + } else { + await rscStream.cancel(reason).catch(() => {}); + } + resolveFlightDataPromise(); + }, + }; + return new TransformStream(transformer); } async function writeRSCStream( - rscStream: ReadableStream, + reader: ReadableStreamDefaultReader, controller: TransformStreamDefaultController, + isCancelled: () => boolean, ) { let decoder = new TextDecoder("utf-8", { fatal: true }); - const reader = rscStream.getReader(); try { let read: ReadableStreamReadResult; while ((read = await reader.read()) && !read.done) { + if (isCancelled()) { + return; + } const chunk = read.value; // Try decoding the chunk to send as a string. // If that fails (e.g. binary data that is invalid unicode), write as base64. @@ -92,7 +120,7 @@ async function writeRSCStream( } let remaining = decoder.decode(); - if (remaining.length) { + if (remaining.length && !isCancelled()) { writeChunk(JSON.stringify(remaining), controller); } }