diff --git a/packages/vinext/src/server/app-rsc-handler.ts b/packages/vinext/src/server/app-rsc-handler.ts index de89ff34e5..c5ee2740a1 100644 --- a/packages/vinext/src/server/app-rsc-handler.ts +++ b/packages/vinext/src/server/app-rsc-handler.ts @@ -32,6 +32,7 @@ import { closeAfterResponse, closeAfterResponseWithBody, createRequestContext, + preserveFullyBufferedBodyMetadata, runWithRequestContext, } from "vinext/shims/unified-request-context"; import { flattenErrorCauses } from "../utils/error-cause.js"; @@ -146,11 +147,14 @@ function applyMiddlewareContextToResponse( const headers = new Headers(response.headers); mergeMiddlewareResponseHeaders(headers, middlewareContext.headers); - return new Response(response.body, { - status: middlewareContext.status ?? response.status, - statusText: response.statusText, - headers, - }); + return preserveFullyBufferedBodyMetadata( + response, + new Response(response.body, { + status: middlewareContext.status ?? response.status, + statusText: response.statusText, + headers, + }), + ); } type DispatchMatchedPageOptions = { diff --git a/packages/vinext/src/server/metadata-route-response.ts b/packages/vinext/src/server/metadata-route-response.ts index afe20f18a9..8a369e3f04 100644 --- a/packages/vinext/src/server/metadata-route-response.ts +++ b/packages/vinext/src/server/metadata-route-response.ts @@ -10,6 +10,7 @@ import { type SitemapEntry, } from "./metadata-routes.js"; import { notFoundResponse } from "./http-error-responses.js"; +import { markFullyBufferedBody } from "vinext/shims/unified-request-context"; type AppPageParams = Record; type MetadataRouteFunction = (props: Record) => unknown; @@ -223,12 +224,15 @@ async function handleGeneratedSitemap( if (!isSitemapEntries(result)) { throw new TypeError("Metadata sitemap routes must return an array."); } - return new Response(sitemapToXml(result), { - headers: { - "Content-Type": route.contentType, - "Cache-Control": metadataRouteCacheHeader(route), - }, - }); + // Body serialized to a string here — fully materialized, no producer left. + return markFullyBufferedBody( + new Response(sitemapToXml(result), { + headers: { + "Content-Type": route.contentType, + "Cache-Control": metadataRouteCacheHeader(route), + }, + }), + ); } function findGeneratedImageId( @@ -327,12 +331,15 @@ async function callDynamicMetadataRoute( body = JSON.stringify(result); } - return new Response(body, { - headers: { - "Content-Type": route.contentType, - "Cache-Control": metadataRouteCacheHeader(route), - }, - }); + // Every branch above serialized `result` to a string — fully materialized. + return markFullyBufferedBody( + new Response(body, { + headers: { + "Content-Type": route.contentType, + "Cache-Control": metadataRouteCacheHeader(route), + }, + }), + ); } function serveStaticMetadataRoute(route: MetadataRuntimeRoute): Response { @@ -348,12 +355,15 @@ function serveStaticMetadataRoute(route: MetadataRuntimeRoute): Response { for (let index = 0; index < binary.length; index++) { bytes[index] = binary.charCodeAt(index); } - return new Response(bytes, { - headers: { - "Content-Type": route.contentType, - "Cache-Control": metadataRouteCacheHeader(route), - }, - }); + // Static file bytes, fully in memory — no producer left. + return markFullyBufferedBody( + new Response(bytes, { + headers: { + "Content-Type": route.contentType, + "Cache-Control": metadataRouteCacheHeader(route), + }, + }), + ); } catch (error) { const reason = error instanceof Error && error.message ? `: ${error.message}` : ""; throw new Error( diff --git a/packages/vinext/src/shims/server.ts b/packages/vinext/src/shims/server.ts index 1896f4f0a2..4b8f340e78 100644 --- a/packages/vinext/src/shims/server.ts +++ b/packages/vinext/src/shims/server.ts @@ -23,6 +23,7 @@ import { getRequestContext, isInsideUnifiedScope, queueAfterCallback, + trackAfterPromise, } from "./unified-request-context.js"; import { assertSafeNavigationUrl } from "./url-safety.js"; import { hasBasePath, stripBasePath } from "../utils/base-path.js"; @@ -1243,9 +1244,12 @@ export function after(task: Promise | (() => T | Promise)): void { if (task == null || typeof (task as PromiseLike).then !== "function") { throw new TypeError("`after()`: Argument must be a promise or a function"); } - const guarded = Promise.resolve(task).catch((error) => { - console.error("[vinext] after() task failed:", error); - }); + const guarded = trackAfterPromise( + requestContext, + Promise.resolve(task).catch((error) => { + console.error("[vinext] after() task failed:", error); + }), + ); getRequestExecutionContext()?.waitUntil(guarded); return; } diff --git a/packages/vinext/src/shims/unified-request-context.ts b/packages/vinext/src/shims/unified-request-context.ts index 1c1243eac5..223175ff90 100644 --- a/packages/vinext/src/shims/unified-request-context.ts +++ b/packages/vinext/src/shims/unified-request-context.ts @@ -63,6 +63,7 @@ export type AfterRequestContext = { callbacks: Array<() => unknown>; responseClosed: boolean; pendingCallbacks: number; + pendingPromises: number; completion: Promise | null; resolveCompletion: (() => void) | null; }; @@ -129,6 +130,7 @@ export function createRequestContext(opts?: Partial): Uni callbacks: [], responseClosed: false, pendingCallbacks: 0, + pendingPromises: 0, completion: null, resolveCompletion: null, }, @@ -194,6 +196,14 @@ export function queueAfterCallback(ctx: UnifiedRequestContext, callback: () => u } } +/** Track promise-form after() work that can register a callback before settling. */ +export function trackAfterPromise(ctx: UnifiedRequestContext, promise: Promise): Promise { + ctx.afterContext.pendingPromises += 1; + return promise.finally(() => { + ctx.afterContext.pendingPromises -= 1; + }); +} + /** Bind a callback to every AsyncLocalStorage context active at registration. */ export function bindRequestContextSnapshot( ctx: UnifiedRequestContext, @@ -226,16 +236,80 @@ export async function closeAfterResponse(ctx: UnifiedRequestContext): Promise 0 || + state.pendingCallbacks > 0 || + state.pendingPromises > 0 || + state.resolveCompletion !== null + ); +} + +type ResponseWithFullyBufferedBodyMetadata = Response & { + __vinextFullyBufferedBody?: boolean; +}; + +/** + * Mark a response whose body vinext constructed from a fully in-memory string + * or byte array, as opposed to a body handed back by user code, which could + * still be producing. With no producer left, no `after()` call can originate + * from this body — the one signal that makes it safe for + * `closeAfterResponseWithBody()` to skip close tracking. + * + * Not set for a metadata route's `result instanceof Response` passthrough (a + * user `icon.tsx`/`opengraph-image.tsx` can return a streaming + * `ImageResponse`) or any handler-returned `new Response(stream)` — those + * bodies can still be producing and must keep close tracking. + */ +export function markFullyBufferedBody(response: Response): Response { + (response as ResponseWithFullyBufferedBodyMetadata).__vinextFullyBufferedBody = true; + return response; +} + +function isFullyBufferedBody(response: Response): boolean { + return (response as ResponseWithFullyBufferedBodyMetadata).__vinextFullyBufferedBody === true; +} + +/** Preserve the internal buffered-body signal when response metadata is rebuilt. */ +export function preserveFullyBufferedBodyMetadata(source: Response, target: Response): Response { + return isFullyBufferedBody(source) ? markFullyBufferedBody(target) : target; +} + +/** + * Wrap a response so deferred `after()` callbacks start on stream completion + * or cancellation. Skipped only when the body is marked fully buffered (see + * `markFullyBufferedBody`) and nothing is currently registered — that lets + * the runtime send it with an accurate `Content-Length` instead of chunked + * transfer encoding. + */ export function closeAfterResponseWithBody( response: Response, ctx: UnifiedRequestContext, ): Response { if (!response.body) { + // Resolve the after-lifecycle now (a no-op if nothing is queued) so a + // callback registered after this call returns doesn't wait forever on a + // responseClosed flag nothing else will set. queueMicrotask(() => void closeAfterResponse(ctx)); return response; } + if (isFullyBufferedBody(response) && !requiresResponseCloseTracking(ctx)) { + return response; + } + const passthrough = new TransformStream(); void response.body.pipeTo(passthrough.writable).then( () => void closeAfterResponse(ctx), diff --git a/tests/after-response-close-worker.test.ts b/tests/after-response-close-worker.test.ts new file mode 100644 index 0000000000..c49bd7b44b --- /dev/null +++ b/tests/after-response-close-worker.test.ts @@ -0,0 +1,239 @@ +/** + * Metadata file convention responses (robots(), sitemap(), manifest(), + * static icons) serialize their body to a string or byte array with no + * producer left, so `closeAfterResponseWithBody()` marks them + * (`markFullyBufferedBody`) and skips the close-tracking wrap when no + * `after()` work is pending — letting the runtime send an accurate + * `Content-Length` instead of chunked transfer encoding. Hand-written Route + * Handler responses are not marked, since they may be `new Response(stream)` + * still producing, and keep close tracking. + * + * Runs inside the actual Cloudflare Workers runtime (via wrangler's workerd): + * `Content-Length` is a wire-level artifact of the runtime's own + * serialization, invisible on an in-memory Headers object and not produced + * the same way by every transport (Vite's Node dev server, for example, + * always streams regardless of this code path). + */ +import fs from "node:fs/promises"; +import path from "node:path"; +import { pathToFileURL } from "node:url"; +import { createBuilder } from "vite"; +import { afterAll, beforeAll, describe, expect, it } from "vite-plus/test"; +import vinext from "../packages/vinext/src/index.js"; +import { APP_FIXTURE_DIR, createIsolatedFixture } from "./helpers.js"; + +const CLOUDFLARE_NODE_MODULES = path.resolve( + import.meta.dirname, + "./fixtures/cf-app-basic/node_modules", +); + +type CloudflarePluginFactory = (options: { + viteEnvironment: { name: string; childEnvironments: string[] }; +}) => import("vite").Plugin; + +async function waitForCondition( + condition: () => boolean | Promise, + options?: { intervalMs?: number; timeoutMs?: number }, +): Promise { + const intervalMs = options?.intervalMs ?? 100; + const deadline = Date.now() + (options?.timeoutMs ?? 3000); + while (!(await condition())) { + if (Date.now() >= deadline) { + throw new Error("Timed out waiting for condition"); + } + await new Promise((resolve) => setTimeout(resolve, intervalMs)); + } +} + +describe("closeAfterResponseWithBody on the Cloudflare Workers runtime", () => { + let root = ""; + let worker: { url: Promise; dispose(): Promise } | undefined; + let baseUrl = ""; + + beforeAll(async () => { + root = await createIsolatedFixture( + APP_FIXTURE_DIR, + "vinext-after-response-close-worker-", + // createIsolatedFixture swaps in a plain workspace node_modules symlink + // (here, cf-app-basic's, for @cloudflare/vite-plugin + wrangler), which + // drops app-basic's own `file:./__test_packages__/*` local packages. + // Exclude the two routes that depend on those — unrelated to this + // regression — so the rest of the fixture still builds. + (src) => + !src.includes(`${path.sep}app${path.sep}context-dedup-test`) && + !src.includes(`${path.sep}app${path.sep}nextjs-compat${path.sep}node-modules-css`), + CLOUDFLARE_NODE_MODULES, + ); + // app-basic has no wrangler config of its own (it's normally only used + // for Vite/Node-mode tests) — write a minimal one so @cloudflare/vite-plugin + // emits dist/server/wrangler.json for unstable_startWorker to read. + await fs.writeFile( + path.join(root, "wrangler.jsonc"), + JSON.stringify({ + name: "vinext-after-response-close-worker-fixture", + compatibility_date: "2026-04-01", + compatibility_flags: ["nodejs_compat"], + main: "vinext/server/fetch-handler", + assets: { not_found_handling: "none", binding: "ASSETS" }, + }), + ); + // Force the metadata response through the middleware header/status rebuild + // that previously discarded the internal fully-buffered marker. + await fs.writeFile( + path.join(root, "middleware.ts"), + `import { NextResponse } from "next/server"; +export function middleware() { + const response = NextResponse.next(); + response.headers.set("x-metadata-middleware", "applied"); + return response; +} +export const config = { matcher: ["/robots.txt"] }; +`, + ); + const cloudflarePluginPath = path.join( + root, + "node_modules/@cloudflare/vite-plugin/dist/index.mjs", + ); + const { cloudflare } = (await import(pathToFileURL(cloudflarePluginPath).href)) as { + cloudflare: CloudflarePluginFactory; + }; + const builder = await createBuilder({ + root, + configFile: false, + plugins: [ + vinext({ appDir: root }), + cloudflare({ viteEnvironment: { name: "rsc", childEnvironments: ["ssr"] } }), + ], + logLevel: "silent", + }); + await builder.buildApp(); + + const wranglerPath = path.join(root, "node_modules/wrangler/wrangler-dist/cli.js"); + const wrangler = (await import(pathToFileURL(wranglerPath).href)) as { + unstable_startWorker(options: { + config: string; + dev: { + remote: false; + persist: false; + logLevel: "none"; + watch: false; + server: { port: 0 }; + }; + }): Promise<{ url: Promise; dispose(): Promise }>; + }; + worker = await wrangler.unstable_startWorker({ + config: path.join(root, "dist/server/wrangler.json"), + dev: { + remote: false, + persist: false, + logLevel: "none", + watch: false, + server: { port: 0 }, + }, + }); + await worker.url; + baseUrl = (await worker.url).origin; + }, 180_000); + + afterAll(async () => { + await worker?.dispose(); + if (root) await fs.rm(root, { recursive: true, force: true }); + }); + + it("preserves Content-Length on a metadata route file convention response (robots.txt)", async () => { + // Accept-Encoding: identity opts out of compression, which would drop + // Content-Length for an unrelated reason — Node's fetch() otherwise + // negotiates gzip/br automatically. + const res = await fetch(`${baseUrl}/robots.txt`, { + headers: { "accept-encoding": "identity" }, + }); + expect(res.status).toBe(200); + expect(res.headers.get("content-encoding")).toBeNull(); + expect(res.headers.get("x-metadata-middleware")).toBe("applied"); + + const bodyText = await res.text(); + expect(bodyText).toContain("Disallow: /private/"); + + const expectedLength = new TextEncoder().encode(bodyText).byteLength; + const contentLength = res.headers.get("content-length"); + expect(contentLength).not.toBeNull(); + expect(Number(contentLength)).toBe(expectedLength); + }); + + it("preserves Content-Length on a metadata sitemap.xml response", async () => { + const res = await fetch(`${baseUrl}/sitemap.xml`, { + headers: { "accept-encoding": "identity" }, + }); + expect(res.status).toBe(200); + expect(res.headers.get("content-encoding")).toBeNull(); + + const bodyText = await res.text(); + expect(bodyText).toContain(" { + const res = await fetch(`${baseUrl}/api/get-only`, { + headers: { "accept-encoding": "identity" }, + }); + expect(res.status).toBe(200); + expect(await res.text()).toBe(JSON.stringify({ method: "GET" })); + expect(res.headers.get("content-length")).toBeNull(); + }); + + it("still defers and runs after() callbacks for a route that registers one", async () => { + const before = (await (await fetch(`${baseUrl}/api/after-test`)).json()) as { + counter: number; + }; + + const postRes = await fetch(`${baseUrl}/api/after-test`, { method: "POST" }); + expect(postRes.status).toBe(200); + expect(await postRes.json()).toEqual({ success: true }); + + // The fixture's after() callback does ~2s of simulated background work + // before incrementing the counter, so poll for it rather than sleeping a + // fixed duration and checking once. + await waitForCondition( + async () => { + const after = (await (await fetch(`${baseUrl}/api/after-test`)).json()) as { + counter: number; + }; + return after.counter === before.counter + 1; + }, + { timeoutMs: 5000 }, + ); + }); + + it("still runs after() registered from within a Route Handler's own still-producing stream, after the handler itself has already returned", async () => { + await fetch(`${baseUrl}/api/late-after-stream?reset=1`); + const before = (await (await fetch(`${baseUrl}/api/late-after-stream?check=1`)).json()) as { + ran: boolean; + }; + expect(before.ran).toBe(false); + + const streamRes = await fetch(`${baseUrl}/api/late-after-stream`); + expect(streamRes.status).toBe(200); + expect(streamRes.headers.get("content-length")).toBeNull(); + expect(await streamRes.text()).toBe("hello-stream"); + + await waitForCondition( + async () => { + const after = (await (await fetch(`${baseUrl}/api/late-after-stream?check=1`)).json()) as { + ran: boolean; + }; + return after.ran === true; + }, + { timeoutMs: 5000 }, + ); + }); + + it("keeps a streaming HTML response chunked", async () => { + const res = await fetch(`${baseUrl}/`); + expect(res.status).toBe(200); + expect(res.headers.get("transfer-encoding")).toBe("chunked"); + }); +}); diff --git a/tests/fixtures/app-basic/app/api/late-after-stream/route.ts b/tests/fixtures/app-basic/app/api/late-after-stream/route.ts new file mode 100644 index 0000000000..e82be242dc --- /dev/null +++ b/tests/fixtures/app-basic/app/api/late-after-stream/route.ts @@ -0,0 +1,32 @@ +import { after } from "next/server"; + +// Module-level state for testing deferred after() timing/ordering, mirroring +// api/after-test's pattern. +let ran = false; + +export async function GET(request: Request) { + const url = new URL(request.url); + if (url.searchParams.get("check") === "1") { + return Response.json({ ran }); + } + if (url.searchParams.get("reset") === "1") { + ran = false; + return Response.json({ resetDone: true }); + } + + // GET() returns immediately with a Response wrapping a ReadableStream + // whose producer — including its after() call — keeps running afterward. + const encoder = new TextEncoder(); + const stream = new ReadableStream({ + async start(controller) { + await new Promise((resolve) => setTimeout(resolve, 200)); + after(() => { + ran = true; + }); + controller.enqueue(encoder.encode("hello-stream")); + controller.close(); + }, + }); + + return new Response(stream); +} diff --git a/tests/shims.test.ts b/tests/shims.test.ts index 82b5a14583..79ed9542a3 100644 --- a/tests/shims.test.ts +++ b/tests/shims.test.ts @@ -5069,6 +5069,306 @@ describe("next/server shim", () => { expect(completed).toBe(true); }); + it("returns a fully-buffered-marked response unchanged when no after() work is registered", async () => { + const { + closeAfterResponseWithBody, + createRequestContext, + markFullyBufferedBody, + runWithRequestContext, + } = await import("../packages/vinext/src/shims/unified-request-context.js"); + const requestContext = createRequestContext(); + const original = markFullyBufferedBody( + new Response("buffered", { + status: 201, + statusText: "Created", + headers: { "x-custom": "value" }, + }), + ); + + const response = await runWithRequestContext(requestContext, () => + closeAfterResponseWithBody(original, requestContext), + ); + + // Same object, not rebuilt through a TransformStream. + expect(response).toBe(original); + expect(response.status).toBe(201); + expect(response.statusText).toBe("Created"); + expect(response.headers.get("x-custom")).toBe("value"); + expect(await response.text()).toBe("buffered"); + }); + + it("still wraps a fully-buffered-marked response when after() is registered", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { + closeAfterResponseWithBody, + createRequestContext, + markFullyBufferedBody, + runWithRequestContext, + } = await import("../packages/vinext/src/shims/unified-request-context.js"); + const requestContext = createRequestContext(); + const original = markFullyBufferedBody(new Response("buffered")); + let called = false; + let response!: Response; + + await runWithRequestContext(requestContext, () => { + after(() => { + called = true; + }); + response = closeAfterResponseWithBody(original, requestContext); + }); + + expect(response).not.toBe(original); + await Promise.resolve(); + await Promise.resolve(); + expect(called).toBe(false); + + await response.text(); + await vi.waitFor(() => expect(called).toBe(true)); + }); + + // Ported from Next.js: packages/next/src/server/after/after-context.test.ts + // https://github.com/vercel/next.js/blob/canary/packages/next/src/server/after/after-context.test.ts + it("runs after() callbacks added from after(promise) after a fully buffered response returns", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { + closeAfterResponseWithBody, + createRequestContext, + markFullyBufferedBody, + runWithRequestContext, + } = await import("../packages/vinext/src/shims/unified-request-context.js"); + const waitUntilCalls: Promise[] = []; + const requestContext = createRequestContext({ + executionContext: { + waitUntil(promise: Promise) { + waitUntilCalls.push(promise); + }, + }, + }); + let releasePromise!: () => void; + const promiseGate = new Promise((resolve) => { + releasePromise = resolve; + }); + let called = false; + let original!: Response; + let response!: Response; + + await runWithRequestContext(requestContext, () => { + after( + (async () => { + await promiseGate; + after(() => { + called = true; + }); + })(), + ); + original = markFullyBufferedBody(new Response("buffered")); + response = closeAfterResponseWithBody(original, requestContext); + }); + + // The outstanding promise keeps close tracking enabled. A callback it + // registers must wait for the body to close, not run as soon as the + // promise settles while the response is still being delivered. + expect(response).not.toBe(original); + releasePromise(); + await vi.waitFor(() => expect(requestContext.afterContext.callbacks).toHaveLength(1)); + expect(called).toBe(false); + + expect(await response.text()).toBe("buffered"); + await vi.waitFor(() => expect(called).toBe(true)); + await Promise.all(waitUntilCalls); + }); + + it("wraps an unmarked response even with no after() work registered yet", async () => { + // No fully-buffered marker (e.g. a hand-written Route Handler returning + // new Response(stream), or a page render whose Suspense-wrapped async + // Server Components can still call after() later) — must stay wrapped. + const { closeAfterResponseWithBody, createRequestContext, runWithRequestContext } = + await import("../packages/vinext/src/shims/unified-request-context.js"); + const requestContext = createRequestContext(); + const original = new Response("streamed"); + + const response = await runWithRequestContext(requestContext, () => + closeAfterResponseWithBody(original, requestContext), + ); + + expect(response).not.toBe(original); + expect(await response.text()).toBe("streamed"); + }); + + it("defers a late after() call registered from an unmarked still-producing stream until it completes", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { closeAfterResponseWithBody, createRequestContext, runWithRequestContext } = + await import("../packages/vinext/src/shims/unified-request-context.js"); + const requestContext = createRequestContext(); + let releaseStream!: () => void; + const streamGate = new Promise((resolve) => { + releaseStream = resolve; + }); + let called = false; + let response!: Response; + + await runWithRequestContext(requestContext, () => { + // No after() call has happened yet when closeAfterResponseWithBody runs. + const body = new ReadableStream({ + async start(controller) { + await streamGate; + after(() => { + called = true; + }); + controller.enqueue(new TextEncoder().encode("streamed")); + controller.close(); + }, + }); + response = closeAfterResponseWithBody(new Response(body), requestContext); + }); + + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + + releaseStream(); + await Promise.resolve(); + await Promise.resolve(); + // after() has just been registered but the stream hasn't closed yet. + expect(called).toBe(false); + + await response.text(); + await vi.waitFor(() => expect(called).toBe(true)); + }); + + it("runs a late after() call from an unmarked still-producing stream when the client cancels instead of consuming it", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { closeAfterResponseWithBody, createRequestContext, runWithRequestContext } = + await import("../packages/vinext/src/shims/unified-request-context.js"); + const requestContext = createRequestContext(); + let releaseStream!: () => void; + const streamGate = new Promise((resolve) => { + releaseStream = resolve; + }); + let called = false; + let response!: Response; + + await runWithRequestContext(requestContext, () => { + const body = new ReadableStream({ + async start(controller) { + await streamGate; + after(() => { + called = true; + }); + controller.enqueue(new TextEncoder().encode("streamed")); + // Deliberately no controller.close() — the client below cancels + // instead of the producer ever finishing on its own. + }, + }); + response = closeAfterResponseWithBody(new Response(body), requestContext); + }); + + releaseStream(); + await Promise.resolve(); + await Promise.resolve(); + await response.body!.cancel(); + + await vi.waitFor(() => expect(called).toBe(true)); + }); + + it("keeps a late after() call from an unmarked still-producing stream wired to ctx.executionContext.waitUntil", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { closeAfterResponseWithBody, createRequestContext, runWithRequestContext } = + await import("../packages/vinext/src/shims/unified-request-context.js"); + const waitUntilCalls: Promise[] = []; + const executionContext = { + waitUntil(promise: Promise) { + waitUntilCalls.push(promise); + }, + }; + const requestContext = createRequestContext({ executionContext }); + let releaseStream!: () => void; + const streamGate = new Promise((resolve) => { + releaseStream = resolve; + }); + let called = false; + let response!: Response; + + await runWithRequestContext(requestContext, () => { + const body = new ReadableStream({ + async start(controller) { + await streamGate; + after(() => { + called = true; + }); + controller.enqueue(new TextEncoder().encode("streamed")); + controller.close(); + }, + }); + response = closeAfterResponseWithBody(new Response(body), requestContext); + }); + + // Nothing registered yet, so no waitUntil call has happened. + expect(waitUntilCalls).toHaveLength(0); + + releaseStream(); + await response.text(); + await vi.waitFor(() => expect(waitUntilCalls).toHaveLength(1)); + await waitUntilCalls[0]; + expect(called).toBe(true); + }); + + it("closeAfterResponseWithBody still resolves queued after() work when the client cancels the body instead of consuming it", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { closeAfterResponseWithBody, createRequestContext, runWithRequestContext } = + await import("../packages/vinext/src/shims/unified-request-context.js"); + const requestContext = createRequestContext(); + let called = false; + + // Deliberately never closes on its own — the test cancels it below, + // simulating a client that disconnects mid-stream. + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("partial")); + }, + }); + + let response!: Response; + await runWithRequestContext(requestContext, () => { + after(() => { + called = true; + }); + response = closeAfterResponseWithBody(new Response(body), requestContext); + }); + + expect(called).toBe(false); + await response.body!.cancel(); + + await vi.waitFor(() => expect(called).toBe(true)); + }); + + it("closeAfterResponseWithBody keeps the after-lifecycle completion wired to ctx.executionContext.waitUntil", async () => { + const { after } = await import("../packages/vinext/src/shims/server.js"); + const { closeAfterResponseWithBody, createRequestContext, runWithRequestContext } = + await import("../packages/vinext/src/shims/unified-request-context.js"); + const waitUntilCalls: Promise[] = []; + const executionContext = { + waitUntil(promise: Promise) { + waitUntilCalls.push(promise); + }, + }; + const requestContext = createRequestContext({ executionContext }); + let called = false; + let response!: Response; + + await runWithRequestContext(requestContext, () => { + after(() => { + called = true; + }); + response = closeAfterResponseWithBody(new Response("streamed"), requestContext); + }); + + expect(waitUntilCalls).toHaveLength(1); + await response.text(); + await waitUntilCalls[0]; + expect(called).toBe(true); + }); + // Next.js uses an unbounded PromiseQueue for after callbacks, so callbacks // registered together start concurrently rather than blocking each other. it("after() starts sibling callbacks concurrently", async () => {