From 5a0c4e4a685036dcb34732461ffdc51e69c33054 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Tue, 21 Jul 2026 15:46:25 -0300 Subject: [PATCH 1/6] fix(server): stop reporting client-aborted renders through onRequestError --- .../vinext/src/server/app-page-response.ts | 2 + packages/vinext/src/server/app-rsc-errors.ts | 10 +++ packages/vinext/src/server/app-ssr-entry.ts | 12 +++- .../vinext/src/server/response-aborted.ts | 64 +++++++++++++++++++ tests/app-router-production-server.test.ts | 45 +++++++++++++ tests/app-rsc-errors.test.ts | 47 ++++++++++++++ .../app/slow-stream-abort-test/page.tsx | 19 ++++++ 7 files changed, 198 insertions(+), 1 deletion(-) create mode 100644 packages/vinext/src/server/response-aborted.ts create mode 100644 tests/fixtures/app-basic/app/slow-stream-abort-test/page.tsx diff --git a/packages/vinext/src/server/app-page-response.ts b/packages/vinext/src/server/app-page-response.ts index 6554884a95..857977598e 100644 --- a/packages/vinext/src/server/app-page-response.ts +++ b/packages/vinext/src/server/app-page-response.ts @@ -12,6 +12,7 @@ import { VINEXT_TIMING_HEADER, } from "./headers.js"; import { setCacheStateHeaders } from "./cache-headers.js"; +import { tagConsumerCancellation } from "./response-aborted.js"; import { mergeMiddlewareResponseHeaders } from "./middleware-response-headers.js"; import { VINEXT_RSC_CONTENT_TYPE, @@ -284,6 +285,7 @@ export function buildAppPageRscResponse( body: ReadableStream, options: BuildAppPageRscResponseOptions, ): Response { + body = tagConsumerCancellation(body); const headers = new Headers({ "Content-Type": VINEXT_RSC_CONTENT_TYPE, Vary: VINEXT_RSC_VARY_HEADER, diff --git a/packages/vinext/src/server/app-rsc-errors.ts b/packages/vinext/src/server/app-rsc-errors.ts index af6cb2da83..dbc2b6098b 100644 --- a/packages/vinext/src/server/app-rsc-errors.ts +++ b/packages/vinext/src/server/app-rsc-errors.ts @@ -1,5 +1,6 @@ import { resolveAppPageSpecialError } from "./app-page-execution.js"; import { isNavigationSignalError } from "../utils/navigation-signal.js"; +import { isResponseAbortedError } from "./response-aborted.js"; type DigestError = Error & { digest?: string }; const ORIGINAL_SERVER_ERROR = Symbol.for("vinext.originalServerError"); @@ -121,6 +122,15 @@ export function createRscOnErrorHandler( return wellKnownDigest; } + // Renders aborted because the response consumer went away (client + // disconnect / aborted navigation) are expected control flow, not request + // errors. The response boundary tags that cancellation (see + // response-aborted.ts); skip reporting, matching Next.js's isAbortError + // handling. The digest is moot since no client is listening anymore. + if (isResponseAbortedError(error)) { + return errorDigest(getThrownValueMessage(error)); + } + if ( nodeEnv !== "production" && error instanceof Error && diff --git a/packages/vinext/src/server/app-ssr-entry.ts b/packages/vinext/src/server/app-ssr-entry.ts index 591b8c0650..190709f6ab 100644 --- a/packages/vinext/src/server/app-ssr-entry.ts +++ b/packages/vinext/src/server/app-ssr-entry.ts @@ -53,6 +53,8 @@ import { AppRouterContext } from "vinext/shims/internal/app-router-context"; import { createClientReferencePreloader } from "./app-client-reference-preloader.js"; import { RSC_FORM_STATE_GLOBAL } from "./app-browser-hydration.js"; import { isPprFallbackShellAbortError } from "vinext/shims/ppr-fallback-shell"; +import { isResponseAbortedError } from "./response-aborted.js"; +import { pumpThrough } from "./stream-pump.js"; import DefaultGlobalError from "vinext/shims/default-global-error"; import { appendAssetDeploymentIdQuery } from "../utils/deployment-id.js"; import { ssrAppRouterInstance } from "./app-ssr-router-instance.js"; @@ -598,6 +600,13 @@ export async function handleSsr( return undefined; } + // Consumer cancellation tagged at the response boundary (client + // disconnect): not a real render failure, keep it out of the + // error meta stream. + if (isResponseAbortedError(error)) { + return undefined; + } + errorMetaRenderer.capture(error); if (error && typeof error === "object" && "digest" in error) { @@ -717,7 +726,8 @@ export async function handleSsr( } const finalStream = deferUntilStreamConsumed( - htmlStream.pipeThrough( + pumpThrough( + htmlStream, createTickBufferedTransform( rscEmbed, getInsertedHTML, diff --git a/packages/vinext/src/server/response-aborted.ts b/packages/vinext/src/server/response-aborted.ts new file mode 100644 index 0000000000..baec467b9f --- /dev/null +++ b/packages/vinext/src/server/response-aborted.ts @@ -0,0 +1,64 @@ +const RESPONSE_ABORTED_NAME = "ResponseAborted"; + +/** + * Cancellation reason used when the response consumer goes away (client + * disconnect, aborted navigation). Spec-compliant server runtimes cancel the + * response body stream in that situation, usually with no reason; React then + * aborts the in-flight render with "The render was aborted by the server + * without a reason." and every aborted task reaches the render `onError`, + * where it would be reported through `onRequestError` as if it were a real + * failure. + * + * Tagging the cancellation at the response boundary lets the error handlers + * classify these aborts as expected control flow instead. Mirrors Next.js's + * `ResponseAborted` (`server/web/spec-extension/adapters/next-request.ts`) and + * its `isAbortError` handling in `server/pipe-readable.ts`, which swallows + * exactly this class of error. + */ +export class ResponseAbortedError extends Error { + constructor() { + super("The client closed the connection before the render completed."); + this.name = RESPONSE_ABORTED_NAME; + } +} + +export function isResponseAbortedError(error: unknown): boolean { + return error instanceof Error && error.name === RESPONSE_ABORTED_NAME; +} + +/** + * Wrap a render stream that is about to become a Response body so that a + * consumer cancellation with no reason reaches the underlying render as a + * tagged {@link ResponseAbortedError} instead of an anonymous abort. + */ +export function tagConsumerCancellation( + stream: ReadableStream, +): ReadableStream { + const reader = stream.getReader(); + let cancelled = false; + return new ReadableStream( + { + pull(controller) { + return reader.read().then( + ({ done, value }) => { + if (cancelled) return; + if (done) { + controller.close(); + } else { + controller.enqueue(value); + } + }, + (error) => { + if (cancelled) return; + controller.error(error); + }, + ); + }, + cancel(reason) { + cancelled = true; + return reader.cancel(reason ?? new ResponseAbortedError()); + }, + }, + { highWaterMark: 0 }, + ); +} diff --git a/tests/app-router-production-server.test.ts b/tests/app-router-production-server.test.ts index d35b65a78a..63e3209648 100644 --- a/tests/app-router-production-server.test.ts +++ b/tests/app-router-production-server.test.ts @@ -1,6 +1,7 @@ import fs from "node:fs"; import http from "node:http"; import os from "node:os"; +import { pathToFileURL } from "node:url"; import path from "node:path"; import { createBuilder } from "vite"; import { afterAll, beforeAll, describe, expect, it } from "vite-plus/test"; @@ -1458,6 +1459,50 @@ describe("App Router Production server (startProdServer)", () => { expect(html).not.toContain("Primary not-found metadata"); }); + it("does not report client-aborted renders via instrumentation", async () => { + const resetRes = await fetch(`${baseUrl}/api/instrumentation-test`, { + method: "DELETE", + }); + expect(resetRes.status).toBe(200); + + // Spec-compliant runtimes (Nitro/h3, Workers, Deno) cancel the response + // body stream when the client disconnects mid-stream. Node's fromWeb + + // pipeline path in prod-server does not, so exercise the entry directly. + const entryModule = await import( + pathToFileURL(path.join(outDir, "server", "index.js")).href + ); + const entryHandler = + typeof entryModule.default === "function" + ? entryModule.default + : entryModule.default.fetch.bind(entryModule.default); + const rscHeaders = { Accept: "text/x-component", RSC: "1" }; + let response: Response = await entryHandler( + new Request(`${baseUrl}/slow-stream-abort-test`, { headers: rscHeaders }), + ); + if (response.status === 307) { + const location = response.headers.get("location")!; + response = await entryHandler( + new Request(new URL(location, baseUrl), { headers: rscHeaders }), + ); + } + expect(response.status).toBe(200); + const reader = response.body!.getReader(); + const shellChunk = await reader.read(); + expect(shellChunk.done).toBe(false); + await reader.cancel(); + + // Give the suspended section time to settle server-side after the abort. + await new Promise((resolve) => setTimeout(resolve, 2500)); + + const stateRes = await fetch(`${baseUrl}/api/instrumentation-test`); + expect(stateRes.status).toBe(200); + const state = await stateRes.json(); + + expect( + state.errors.map((error: { message: string }) => error.message), + ).toEqual([]); + }); + it("reports server component render errors via instrumentation in production", async () => { const resetRes = await fetch(`${baseUrl}/api/instrumentation-test`, { method: "DELETE", diff --git a/tests/app-rsc-errors.test.ts b/tests/app-rsc-errors.test.ts index c317f722c2..0294351787 100644 --- a/tests/app-rsc-errors.test.ts +++ b/tests/app-rsc-errors.test.ts @@ -5,6 +5,11 @@ import { getDigestForWellKnownError, sanitizeErrorForClient, } from "../packages/vinext/src/server/app-rsc-errors.js"; +import { + isResponseAbortedError, + ResponseAbortedError, + tagConsumerCancellation, +} from "../packages/vinext/src/server/response-aborted.js"; type DigestCarrier = Error & { digest: unknown }; @@ -55,6 +60,48 @@ describe("app RSC error primitives", () => { expect(expectDigestError(sanitized).digest).toBe("existing-digest"); }); + it("does not report renders aborted by response consumer cancellation", () => { + const reportRequestError = vi.fn(); + const onError = createRscOnErrorHandler({ + errorContext: { routerKind: "App Router", routePath: "/aborted", routeType: "render" }, + nodeEnv: "production", + reportRequestError, + requestInfo: { path: "/aborted", method: "GET", headers: {} }, + }); + + const digest = onError(new ResponseAbortedError()); + + expect(reportRequestError).not.toHaveBeenCalled(); + expect(typeof digest).toBe("string"); + }); + + it("tags reason-less consumer cancellation before it reaches the render stream", async () => { + let cancelReason: unknown = "unset"; + const inner = new ReadableStream({ + cancel(reason) { + cancelReason = reason; + }, + }); + + await tagConsumerCancellation(inner).cancel(); + + expect(isResponseAbortedError(cancelReason)).toBe(true); + }); + + it("forwards an explicit cancellation reason unchanged", async () => { + let cancelReason: unknown = "unset"; + const inner = new ReadableStream({ + cancel(reason) { + cancelReason = reason; + }, + }); + + const explicit = new Error("runtime supplied reason"); + await tagConsumerCancellation(inner).cancel(explicit); + + expect(cancelReason).toBe(explicit); + }); + it("reports the original server error when the client transport error is sanitized", () => { const original = new Error("metadata secret"); original.stack = "original stack"; diff --git a/tests/fixtures/app-basic/app/slow-stream-abort-test/page.tsx b/tests/fixtures/app-basic/app/slow-stream-abort-test/page.tsx new file mode 100644 index 0000000000..8662f03a5f --- /dev/null +++ b/tests/fixtures/app-basic/app/slow-stream-abort-test/page.tsx @@ -0,0 +1,19 @@ +import { Suspense } from "react"; + +export const dynamic = "force-dynamic"; + +async function SlowSection() { + await new Promise((resolve) => setTimeout(resolve, 1500)); + return

slow section done

; +} + +export default function SlowStreamAbortTestPage() { + return ( +
+

slow stream abort test

+ loading slow section...

}> + +
+
+ ); +} From 92de6f4a1e2f241d2fc356b02d2292f45807e51e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Tue, 21 Jul 2026 15:46:25 -0300 Subject: [PATCH 2/6] fix(server): own stream pump promises instead of pipeThrough --- packages/vinext/src/server/app-page-render.ts | 15 ++++---- .../src/server/defer-until-stream-consumed.ts | 14 ++++---- .../vinext/src/server/rsc-stream-hints.ts | 4 +-- packages/vinext/src/server/stream-pump.ts | 34 +++++++++++++++++++ 4 files changed, 50 insertions(+), 17 deletions(-) create mode 100644 packages/vinext/src/server/stream-pump.ts diff --git a/packages/vinext/src/server/app-page-render.ts b/packages/vinext/src/server/app-page-render.ts index c3cb463bc2..8ed27496d2 100644 --- a/packages/vinext/src/server/app-page-render.ts +++ b/packages/vinext/src/server/app-page-render.ts @@ -584,19 +584,18 @@ function wrapRscResponseForDevErrorReporting( } }; - const cleanup = new TransformStream({ - flush() { - onConsumed(); - }, - }); - - const piped = originalBody.pipeThrough(cleanup); - const reader = piped.getReader(); + // A manual passthrough instead of pipeThrough(new TransformStream(...)): + // the transform's only job was firing `onConsumed` on clean drain, which + // the `done` branch below covers. It also avoids the internal pipeTo + // promise, which rejects unhandled when the consumer cancels the stream + // with a reason. + const reader = originalBody.getReader(); const wrappedStream = new ReadableStream({ pull(controller) { return reader.read().then( ({ done, value }) => { if (done) { + onConsumed(); controller.close(); } else { controller.enqueue(value); diff --git a/packages/vinext/src/server/defer-until-stream-consumed.ts b/packages/vinext/src/server/defer-until-stream-consumed.ts index cf802b5bea..6414d5908b 100644 --- a/packages/vinext/src/server/defer-until-stream-consumed.ts +++ b/packages/vinext/src/server/defer-until-stream-consumed.ts @@ -13,18 +13,18 @@ export function deferUntilStreamConsumed( } }; - const cleanup = new TransformStream({ - flush() { - once(); - }, - }); - - const reader = stream.pipeThrough(cleanup).getReader(); + // A manual passthrough instead of pipeThrough(new TransformStream(...)): + // the transform's only job was firing `onFlush` on clean drain, which the + // `done` branch below covers. It also avoids the internal pipeTo promise, + // which rejects unhandled when the consumer cancels the stream with a + // reason. + const reader = stream.getReader(); return new ReadableStream({ pull(controller) { return reader.read().then( ({ done, value }) => { if (done) { + once(); controller.close(); } else { controller.enqueue(value); diff --git a/packages/vinext/src/server/rsc-stream-hints.ts b/packages/vinext/src/server/rsc-stream-hints.ts index 2c21e12132..7a3ba3accd 100644 --- a/packages/vinext/src/server/rsc-stream-hints.ts +++ b/packages/vinext/src/server/rsc-stream-hints.ts @@ -1,3 +1,4 @@ +import { pumpThrough } from "./stream-pump.js"; const REACT_FLIGHT_STYLESHEET_PRELOAD_HINT = /^([0-9a-f]*:HL\[.*?),"stylesheet"(\]|,)/; const STYLESHEET_TO_STYLE_JSON_PADDING = " ".repeat("stylesheet".length - "style".length); @@ -112,8 +113,7 @@ export function normalizeReactFlightPreloadHints( let rawBytesRemaining = 0; let passThrough = false; - return stream.pipeThrough( - new TransformStream({ + return pumpThrough(stream, new TransformStream({ transform(chunk, controller) { if (passThrough) { controller.enqueue(chunk); diff --git a/packages/vinext/src/server/stream-pump.ts b/packages/vinext/src/server/stream-pump.ts new file mode 100644 index 0000000000..4c43c7bd8d --- /dev/null +++ b/packages/vinext/src/server/stream-pump.ts @@ -0,0 +1,34 @@ +/** + * `stream.pipeThrough(transform)` equivalent whose pump owns every promise. + * + * The internal pipeTo loop of `pipeThrough` can leave an in-flight + * `writer.write` rejection unhandled when the consumer cancels the readable + * side between chunks (observable as an unhandled rejection carrying the + * cancel reason). This pump awaits or settles every read/write explicitly, and + * on failure propagates the reason both ways (cancel upstream, abort the + * transform), so cancellation reasons reach the source untouched. + */ +export function pumpThrough( + stream: ReadableStream, + transform: TransformStream, +): ReadableStream { + const writer = transform.writable.getWriter(); + const reader = stream.getReader(); + + void (async () => { + try { + for (;;) { + const { done, value } = await reader.read(); + if (done) { + await writer.close(); + return; + } + await writer.write(value); + } + } catch (error) { + await Promise.allSettled([reader.cancel(error), writer.abort(error)]); + } + })(); + + return transform.readable; +} From eded7a196da7b3fa97abb833687259fa4af1a738 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Tue, 21 Jul 2026 19:32:31 -0300 Subject: [PATCH 3/6] fix(server): address review feedback on abort classification and pump cancellation --- .../vinext/src/server/response-aborted.ts | 34 ++++++++++-- packages/vinext/src/server/stream-pump.ts | 7 +++ tests/app-rsc-errors.test.ts | 53 +++++++++++++++++++ 3 files changed, 90 insertions(+), 4 deletions(-) diff --git a/packages/vinext/src/server/response-aborted.ts b/packages/vinext/src/server/response-aborted.ts index baec467b9f..893b0d51c6 100644 --- a/packages/vinext/src/server/response-aborted.ts +++ b/packages/vinext/src/server/response-aborted.ts @@ -1,4 +1,5 @@ const RESPONSE_ABORTED_NAME = "ResponseAborted"; +const RESPONSE_ABORTED_BRAND = Symbol.for("vinext.responseAborted"); /** * Cancellation reason used when the response consumer goes away (client @@ -16,14 +17,31 @@ const RESPONSE_ABORTED_NAME = "ResponseAborted"; * exactly this class of error. */ export class ResponseAbortedError extends Error { - constructor() { - super("The client closed the connection before the render completed."); + readonly [RESPONSE_ABORTED_BRAND] = true; + + constructor(cause?: unknown) { + super( + "The client closed the connection before the render completed.", + ...(cause !== undefined ? [{ cause }] : []), + ); this.name = RESPONSE_ABORTED_NAME; } } export function isResponseAbortedError(error: unknown): boolean { - return error instanceof Error && error.name === RESPONSE_ABORTED_NAME; + return ( + error instanceof Error && + (error as Partial>)[RESPONSE_ABORTED_BRAND] === + true + ); +} + +function isDomAbortError(reason: unknown): boolean { + return ( + typeof DOMException !== "undefined" && + reason instanceof DOMException && + reason.name === "AbortError" + ); } /** @@ -56,7 +74,15 @@ export function tagConsumerCancellation( }, cancel(reason) { cancelled = true; - return reader.cancel(reason ?? new ResponseAbortedError()); + if (reason == null) { + return reader.cancel(new ResponseAbortedError()); + } + // Runtimes that propagate a standard aborted signal cancel with a + // DOMException named AbortError; that is still a consumer abort. + if (isDomAbortError(reason)) { + return reader.cancel(new ResponseAbortedError(reason)); + } + return reader.cancel(reason); }, }, { highWaterMark: 0 }, diff --git a/packages/vinext/src/server/stream-pump.ts b/packages/vinext/src/server/stream-pump.ts index 4c43c7bd8d..9a17185a5c 100644 --- a/packages/vinext/src/server/stream-pump.ts +++ b/packages/vinext/src/server/stream-pump.ts @@ -15,6 +15,13 @@ export function pumpThrough( const writer = transform.writable.getWriter(); const reader = stream.getReader(); + // Cancelling the transform's readable errors its writable. Without this + // watcher the pump would only notice on the next chunk, keeping a stalled + // render's resources alive after the consumer went away. + void writer.closed.catch((reason: unknown) => { + void reader.cancel(reason).catch(() => {}); + }); + void (async () => { try { for (;;) { diff --git a/tests/app-rsc-errors.test.ts b/tests/app-rsc-errors.test.ts index 0294351787..2d61657413 100644 --- a/tests/app-rsc-errors.test.ts +++ b/tests/app-rsc-errors.test.ts @@ -10,6 +10,7 @@ import { ResponseAbortedError, tagConsumerCancellation, } from "../packages/vinext/src/server/response-aborted.js"; +import { pumpThrough } from "../packages/vinext/src/server/stream-pump.js"; type DigestCarrier = Error & { digest: unknown }; @@ -88,6 +89,58 @@ describe("app RSC error primitives", () => { expect(isResponseAbortedError(cancelReason)).toBe(true); }); + it("still reports errors that merely share the ResponseAborted name", () => { + const reportRequestError = vi.fn(); + const onError = createRscOnErrorHandler({ + errorContext: { routerKind: "App Router", routePath: "/fake", routeType: "render" }, + nodeEnv: "production", + reportRequestError, + requestInfo: { path: "/fake", method: "GET", headers: {} }, + }); + + const impostor = new Error("looks aborted but is a real failure"); + impostor.name = "ResponseAborted"; + onError(impostor); + + expect(reportRequestError).toHaveBeenCalledOnce(); + }); + + it("tags consumer cancellation carrying a standard AbortError reason", async () => { + let cancelReason: unknown = "unset"; + const inner = new ReadableStream({ + cancel(reason) { + cancelReason = reason; + }, + }); + + const abortReason = new DOMException("The operation was aborted.", "AbortError"); + await tagConsumerCancellation(inner).cancel(abortReason); + + expect(isResponseAbortedError(cancelReason)).toBe(true); + expect((cancelReason as Error).cause).toBe(abortReason); + }); + + it("cancels the pump source promptly when the readable side is cancelled mid-wait", async () => { + let sourceCancelled: unknown = null; + const neverEndingSource = new ReadableStream({ + pull() { + return new Promise(() => {}); + }, + cancel(reason) { + sourceCancelled = reason; + }, + }); + + const piped = pumpThrough(neverEndingSource, new TransformStream()); + const reader = piped.getReader(); + const pending = reader.read(); + await reader.cancel(new Error("consumer went away")); + await pending.catch(() => {}); + await new Promise((resolve) => setTimeout(resolve, 20)); + + expect(sourceCancelled).toBeInstanceOf(Error); + }); + it("forwards an explicit cancellation reason unchanged", async () => { let cancelReason: unknown = "unset"; const inner = new ReadableStream({ From 4bcb9b92b8b874bae0449bf7d13ed23bc7a5a9d4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jo=C3=A3o=20Lucas?= Date: Tue, 21 Jul 2026 19:59:12 -0300 Subject: [PATCH 4/6] fix(server): tag server action aborts and settle cancellation during pending embed flush --- .../src/server/app-server-action-execution.ts | 6 +++- packages/vinext/src/server/app-ssr-entry.ts | 5 ++- packages/vinext/src/server/app-ssr-stream.ts | 30 +++++++++++++---- tests/app-ssr-stream.test.ts | 33 +++++++++++++++++++ 4 files changed, 65 insertions(+), 9 deletions(-) diff --git a/packages/vinext/src/server/app-server-action-execution.ts b/packages/vinext/src/server/app-server-action-execution.ts index a5ce3794fd..ba20ef1e86 100644 --- a/packages/vinext/src/server/app-server-action-execution.ts +++ b/packages/vinext/src/server/app-server-action-execution.ts @@ -41,6 +41,7 @@ import { applyEdgeRuntimeHeader } from "./app-page-response.js"; import { resolveAppPageActionRerenderTarget } from "./app-page-request.js"; import { resolveAppPageNavigationParams } from "./app-page-element-builder.js"; import { deferUntilStreamConsumed } from "./app-page-stream.js"; +import { tagConsumerCancellation } from "./response-aborted.js"; import { buildAppPageTags } from "./implicit-tags.js"; import { mergeMiddlewareResponseHeaders } from "./middleware-response-headers.js"; import { getSetCookieName } from "./cookie-utils.js"; @@ -489,7 +490,10 @@ function createServerActionRscResponse( return new Response(body, init); } - return new Response(deferUntilStreamConsumed(body, clearRequestContext), init); + return new Response( + tagConsumerCancellation(deferUntilStreamConsumed(body, clearRequestContext)), + init, + ); } function isRequestBodyTooLarge(error: unknown): boolean { diff --git a/packages/vinext/src/server/app-ssr-entry.ts b/packages/vinext/src/server/app-ssr-entry.ts index 190709f6ab..6dc54fd742 100644 --- a/packages/vinext/src/server/app-ssr-entry.ts +++ b/packages/vinext/src/server/app-ssr-entry.ts @@ -738,7 +738,10 @@ export async function handleSsr( options?.scriptNonce, ), ), - cleanup, + () => { + rscEmbed.abort?.(); + cleanup(); + }, ); return { diff --git a/packages/vinext/src/server/app-ssr-stream.ts b/packages/vinext/src/server/app-ssr-stream.ts index bc712fe072..216de95c82 100644 --- a/packages/vinext/src/server/app-ssr-stream.ts +++ b/packages/vinext/src/server/app-ssr-stream.ts @@ -15,6 +15,12 @@ import { NAVIGATION_RUNTIME_SYMBOL_DESCRIPTION } from "../client/navigation-runt type RscEmbedTransform = { flush(): string; finalize(): Promise; + /** + * Stop reading the embed stream so pending finalize()/getRawBuffer() + * callers settle. Used when the response consumer goes away while the + * auxiliary Flight stream is still suspended; a no-op after natural end. + */ + abort?(reason?: unknown): void; /** Resolves when all raw bytes from the embed stream have been read. */ getRawBuffer(): Promise; }; @@ -131,6 +137,10 @@ export function createRscEmbedTransform( const pumpPromise = pumpReader(); return { + abort(reason?: unknown): void { + void reader.cancel(reason).catch(() => {}); + }, + flush(): string { if (pendingChunks.length === 0) return ""; @@ -621,14 +631,20 @@ export function createTickBufferedTransform( } const finalScripts = await rscEmbed.finalize(); - if (finalScripts) { - controller.enqueue(encoder.encode(finalScripts)); - } + try { + if (finalScripts) { + controller.enqueue(encoder.encode(finalScripts)); + } - // Emit `` last so the document always terminates with a - // well-formed close, after any trailing flight chunks / preinit scripts. - // Mirrors Next.js's `createMoveSuffixStream` behaviour (#1532). - controller.enqueue(encoder.encode(DOCUMENT_CLOSE_SUFFIX)); + // Emit `` last so the document always terminates with + // a well-formed close, after any trailing flight chunks / preinit + // scripts. Mirrors Next.js's `createMoveSuffixStream` behaviour + // (#1532). + controller.enqueue(encoder.encode(DOCUMENT_CLOSE_SUFFIX)); + } catch { + // The readable side was cancelled while finalize() was pending + // (consumer went away mid-flush); nobody is reading anymore. + } }, }); } diff --git a/tests/app-ssr-stream.test.ts b/tests/app-ssr-stream.test.ts index 6cda68f033..f881957eb9 100644 --- a/tests/app-ssr-stream.test.ts +++ b/tests/app-ssr-stream.test.ts @@ -8,6 +8,8 @@ import { fixPreloadAs, waitAtLeastOneReactRenderTask, } from "../packages/vinext/src/server/app-ssr-stream.js"; +import { deferUntilStreamConsumed } from "../packages/vinext/src/server/defer-until-stream-consumed.js"; +import { pumpThrough } from "../packages/vinext/src/server/stream-pump.js"; it("serializes dynamic stale time into the hydration bootstrap", () => { expect( @@ -800,4 +802,35 @@ describe("createTickBufferedTransform inline CSS", () => { expect(out).toContain("