From 68144749b4427c946943a4bd2536001f957e57f8 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Sun, 30 Aug 2026 11:08:28 +0000 Subject: [PATCH] test: speed up fetch-backpressure.test.ts and tighten its assertions The four "server stops writing" tests each drained a 1 GiB body after the pause. They now read until the server has written 1 MiB past the point where it stalled, which proves the resume, and cancel the rest. The exact drain of a body that paused is still covered by the 16 MiB tests for each framing. Fixed waits become conditions: a server exposes `settled()` (its write() returned false or the body is out, then sent() holds still), `wrote(n)` and `closed`. The spawned clients report "stalled" on stdout and wait for "go" on stdin instead of sampling RSS for 60 ms. The small-unread-bodies origin holds the last part until the Response exists instead of 2 ms timers. The released-reader client waits for "go" instead of 50 ms. Assertions: every drained body is checked by md5 against the known digest and every streamed chunk is checked to be all "A". The server's sent() must equal the wire size. Fixture children assert stderr === "". The peer-reset tests assert the error name, code and message. The S3 error and AbortError tests assert the message. The previously sequential describes run concurrently. The two tests that count pooled connections stay serial, at the end of the file. The gzip variant uses stored blocks over a body of "A": the wire still carries about as many bytes as the body, and a debug build no longer spends ~2 s compressing 16 MiB of random bytes on the test's main thread. --- test/js/web/fetch/fetch-backpressure.test.ts | 839 +++++++++++-------- 1 file changed, 483 insertions(+), 356 deletions(-) diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index babca0edde63..1ef0849c556b 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -2,8 +2,7 @@ // HTTP thread from buffering the entire response in memory. import { S3Client } from "bun"; import { describe, expect, test } from "bun:test"; -import { bunEnv, bunExe, isASAN, isDebug, isWindows, tempDir, tls } from "harness"; -import { randomBytes } from "node:crypto"; +import { bunEnv, bunExe, forEachLine, isASAN, isDebug, isWindows, tempDir, tls } from "harness"; import { once } from "node:events"; import { statSync } from "node:fs"; import { stat } from "node:fs/promises"; @@ -19,14 +18,93 @@ import { gzipSync } from "node:zlib"; const CHUNK = 64 * 1024; const COUNT = 256; // 16 MiB const TOTAL = CHUNK * COUNT; +const PAYLOAD = Buffer.alloc(CHUNK, 65); + +function md5(data: Uint8Array | string): string { + return Bun.CryptoHasher.hash("md5", data, "hex"); +} + +// Every body here is "A" repeated. A consumer that received all of a 16 MiB one hashes to this. +const DIGEST = md5(Buffer.alloc(TOTAL, 65)); + +// Whether every byte of `chunk` is "A". A memcmp per 2 MiB window: a byte loop over a body in JS is +// slow in debug builds. +const A_WINDOW = Buffer.alloc(2 * 1024 * 1024, 65); +function isAllA(chunk: Uint8Array): boolean { + const buf = Buffer.from(chunk.buffer, chunk.byteOffset, chunk.byteLength); + for (let off = 0; off < buf.length; off += A_WINDOW.length) { + const end = Math.min(off + A_WINDOW.length, buf.length); + if (buf.compare(A_WINDOW, 0, end - off, off, end) !== 0) return false; + } + return true; +} + +// Longer than any loopback buffer: tcp_rmem[2] + tcp_wmem[2] approach 256 MiB on some hosts. A +// server writes it only until the kernel stops taking bytes, and a test reads only as much of it as +// proves the pause and the resume; no test drains it. +const BIG = 16384; // 1 GiB declared +const BODY = BIG * CHUNK; type Kind = "h1" | "h1-chunked" | "h1-gzip" | "h1-tls" | "h2" | "h3"; -async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () => number } & AsyncDisposable> { +type Server = AsyncDisposable & { + url: string; + // Bytes handed to the socket so far (of the gzip stream, for h1-gzip). + sent: () => number; + // What `sent()` reaches once a whole body is out. + wire: number; + // Resolves with `sent()` once the body has stopped moving: the socket stopped taking writes past + // its first packets (or the whole body is out), and `sent()` then held still for three samples + // 10 ms apart. No event reports "nothing more is coming", so the second half has to sample. + settled: () => Promise; + // `sent()` reached `n`. + wrote: (n: number) => Promise; + // The first response's connection closed (for a body that never ends, that is the client's doing). + closed: Promise; +}; + +// The bookkeeping behind `Server`, shared by every transport. +function progress(wire: number) { let sent = 0; - const payload = Buffer.alloc(CHUNK, 65); + const waiters: [number, () => void][] = []; + const blocked = Promise.withResolvers(); + const finished = Promise.withResolvers(); + const closed = Promise.withResolvers(); + return { + add(n: number) { + sent += n; + for (let i = waiters.length; i--; ) if (sent >= waiters[i][0]) waiters.splice(i, 1)[0][1](); + if (sent >= wire) finished.resolve(); + }, + // A write that did not go through once the socket has taken more than 8 chunks: the kernel is + // full, not merely slow to take the first packets. + block() { + if (sent > 8 * CHUNK) blocked.resolve(); + }, + close: () => closed.resolve(), + api: { + sent: () => sent, + wire, + wrote: (n: number) => (sent >= n ? Promise.resolve() : new Promise(resolve => waiters.push([n, resolve]))), + closed: closed.promise, + async settled() { + await Promise.race([blocked.promise, finished.promise]); + let last = -1; + for (let stable = 0; stable < 3; ) { + await Bun.sleep(10); + stable = sent === last ? stable + 1 : 0; + last = sent; + } + return sent; + }, + }, + }; +} +// Writes `count` chunks of "A" as fast as the socket takes them. +async function serve(kind: Kind, count = COUNT): Promise { if (kind === "h2") { + const p = progress(CHUNK * count); const srv = createSecureServer({ ...tls, allowHTTP1: false }); const sockets = new Set(); srv.on("connection", s => { @@ -36,12 +114,16 @@ async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () srv.on("stream", stream => { stream.respond({ ":status": 200, "content-type": "application/octet-stream" }); stream.on("error", () => {}); + stream.on("close", p.close); let i = 0; const push = () => { while (i < count) { i++; - sent += CHUNK; - if (!stream.write(payload)) return void stream.once("drain", push); + p.add(CHUNK); + if (!stream.write(PAYLOAD)) { + p.block(); + return void stream.once("drain", push); + } } stream.end(); }; @@ -52,7 +134,7 @@ async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () const { port } = srv.address() as import("node:net").AddressInfo; return { url: `https://localhost:${port}/`, - sent: () => sent, + ...p.api, [Symbol.asyncDispose]: async () => { for (const s of sockets) s.destroy(); await new Promise(r => srv.close(r)); @@ -61,6 +143,7 @@ async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () } if (kind === "h3") { + const p = progress(CHUNK * count); const srv = Bun.serve({ port: 0, tls, @@ -71,33 +154,45 @@ async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () return new Response( new ReadableStream({ pull(ctrl) { - if (i++ < count) ctrl.enqueue(payload); - else ctrl.close(); + if (i++ < count) { + p.add(CHUNK); + ctrl.enqueue(PAYLOAD); + // A pull-fed stream reports no failed write; past the client's mark, its pace is + // the client's pause. + p.block(); + } else ctrl.close(); }, }), ); }, }); - return { url: String(srv.url), sent: () => sent, [Symbol.asyncDispose]: () => srv.stop(true) }; + return { url: String(srv.url), ...p.api, [Symbol.asyncDispose]: () => srv.stop(true) }; } // h1 / h1-chunked / h1-gzip / h1-tls - const gz = kind === "h1-gzip" ? gzipSync(randomBytes(CHUNK * count)) : null; - const handler = (req: import("node:http").IncomingMessage, res: import("node:http").ServerResponse) => { + // Stored blocks (level 0): the wire carries about as many bytes as the body, so the transport can + // pause on it, and the decoder still runs over every byte. + const gz = kind === "h1-gzip" ? gzipSync(Buffer.alloc(CHUNK * count, 65), { level: 0 }) : null; + const p = progress(gz ? gz.length : CHUNK * count); + const handler = (_req: import("node:http").IncomingMessage, res: import("node:http").ServerResponse) => { res.on("error", () => {}); + res.on("close", p.close); if (gz) { res.setHeader("content-encoding", "gzip"); res.setHeader("content-length", String(gz.length)); let off = 0; const push = () => { - while (off < gz.length) { + while (off < gz.length && !res.destroyed) { const end = Math.min(off + CHUNK, gz.length); const slice = gz.subarray(off, end); off = end; - sent += slice.length; - if (!res.write(slice)) return void res.once("drain", push); + p.add(slice.length); + if (!res.write(slice)) { + p.block(); + return void res.once("drain", push); + } } - res.end(); + if (!res.destroyed) res.end(); }; return push(); } @@ -105,12 +200,15 @@ async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () res.flushHeaders(); let i = 0; const push = () => { - while (i < count) { + while (i < count && !res.destroyed) { i++; - sent += CHUNK; - if (!res.write(payload)) return void res.once("drain", push); + p.add(CHUNK); + if (!res.write(PAYLOAD)) { + p.block(); + return void res.once("drain", push); + } } - res.end(); + if (!res.destroyed) res.end(); }; push(); }; @@ -120,7 +218,7 @@ async function serve(kind: Kind, count = COUNT): Promise<{ url: string; sent: () const { port } = srv.address() as import("node:net").AddressInfo; return { url: `${kind === "h1-tls" ? "https" : "http"}://127.0.0.1:${port}/`, - sent: () => sent, + ...p.api, [Symbol.asyncDispose]: () => { srv.closeAllConnections(); return new Promise(r => srv.close(() => r(undefined))); @@ -134,78 +232,111 @@ function fetchOpts(kind: Kind): RequestInit { return {}; } -async function spawnClient(url: string, kind: Kind, script: string) { +// Reads a body whose server is stalled at `stalled` until the server has written a MiB more, which +// it can only do once the client takes bytes off the socket again. Then cancels the rest. +async function readUntilResumed(reader: ReadableStreamDefaultReader, server: Server, stalled: number) { + let total = 0; + let foreign = 0; + let ended = false; + while (server.sent() < stalled + 16 * CHUNK) { + const { value, done } = await reader.read(); + if (done) { + ended = true; + break; + } + total += value.byteLength; + if (!isAllA(value)) foreign++; + } + await reader.cancel(); + return { total, foreign, ended }; +} + +// What every client script starts with: `url`, `opts`, a `hasher` for the bytes it receives, +// `stall()`, which tells the test that the client holds its first chunk and then waits until the +// server has stopped writing, and `report(total)`. +const CLIENT = /* js */ ` + const hasher = new Bun.CryptoHasher("md5"); + async function stall() { + process.stdout.write("stalled\\n"); + for await (const line of console) if (line === "go") break; + } + function report(total) { + process.stdout.write(JSON.stringify({ total, digest: hasher.digest("hex") }) + "\\n"); + } +`; + +async function spawnClient(server: Server, kind: Kind, script: string) { await using proc = Bun.spawn({ - cmd: [bunExe(), "-e", `const url=${JSON.stringify(url)};const opts=${JSON.stringify(fetchOpts(kind))};${script}`], + cmd: [ + bunExe(), + "-e", + `const url=${JSON.stringify(server.url)};const opts=${JSON.stringify(fetchOpts(kind))};${CLIENT}${script}`, + ], env: { ...bunEnv, NODE_TLS_REJECT_UNAUTHORIZED: "0", BUN_FEATURE_FLAG_EXPERIMENTAL_HTTP2_CLIENT: "1" }, + stdin: "pipe", stdout: "pipe", stderr: "pipe", }); - const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - if (!stdout) throw new Error(`client exited ${exitCode}: ${stderr}`); - return { ...JSON.parse(stdout), stderr, exitCode }; -} - -const SETTLE_RSS = /* js */ ` - const rss = process.platform === "darwin" && typeof Bun.unsafe.memoryFootprint === "function" ? Bun.unsafe.memoryFootprint : process.memoryUsage.rss; - async function settleRss() { - const before = rss(); - let last = before, stable = 0; - while (stable < 3) { - await Bun.sleep(20); - const now = rss(); - stable = Math.abs(now - last) < (1 << 20) ? stable + 1 : 0; - last = now; - } - return last - before; + const stderr = proc.stderr.text(); + let result = ""; + for await (const line of forEachLine(proc.stdout)) { + if (line === "stalled") { + // The client holds its first chunk and reads no further. Let the server run into the pause + // (or, where loopback buffers take all 16 MiB, finish) before the client goes on. + await server.settled(); + proc.stdin.write("go\n"); + proc.stdin.end(); + } else result = line; } -`; + const exitCode = await proc.exited; + if (!result) throw new Error(`client exited ${exitCode}: ${await stderr}`); + return { ...JSON.parse(result), stderr: await stderr, exitCode }; +} -const STALL_READER = - SETTLE_RSS + - /* js */ ` +const STALL_READER = /* js */ ` const res = await fetch(url, opts); const reader = res.body.getReader(); const first = await reader.read(); - const peak = await settleRss(); + await stall(); + hasher.update(first.value); let total = first.value.byteLength; - for (let r; !(r = await reader.read()).done; ) total += r.value.byteLength; - process.stdout.write(JSON.stringify({ peak, total })); + for (let r; !(r = await reader.read()).done; ) { + hasher.update(r.value); + total += r.value.byteLength; + } + report(total); `; -const STALL_PIPE_TO = - SETTLE_RSS + - /* js */ ` +const STALL_PIPE_TO = /* js */ ` const res = await fetch(url, opts); - let peak = 0, total = 0, first = true; + let total = 0, first = true; await res.body.pipeTo(new WritableStream({ async write(chunk) { + if (first) { first = false; await stall(); } + hasher.update(chunk); total += chunk.byteLength; - if (first) { first = false; peak = await settleRss(); } }, })); - process.stdout.write(JSON.stringify({ peak, total })); + report(total); `; -const STALL_FOR_AWAIT = - SETTLE_RSS + - /* js */ ` +const STALL_FOR_AWAIT = /* js */ ` const res = await fetch(url, opts); - let peak = 0, total = 0, first = true; + let total = 0, first = true; for await (const chunk of res.body) { + if (first) { first = false; await stall(); } + hasher.update(chunk); total += chunk.byteLength; - if (first) { first = false; peak = await settleRss(); } } - process.stdout.write(JSON.stringify({ peak, total })); + report(total); `; -const STALL_NO_CONSUMER = - SETTLE_RSS + - /* js */ ` +const STALL_NO_CONSUMER = /* js */ ` const response = await fetch(url, opts); - const peak = await settleRss(); - const total = (await response.arrayBuffer()).byteLength; - process.stdout.write(JSON.stringify({ peak, total })); + await stall(); + const body = new Uint8Array(await response.arrayBuffer()); + hasher.update(body); + report(body.byteLength); `; for (const kind of ["h1", "h1-chunked", "h1-gzip", "h1-tls", "h2", "h3"] as Kind[]) { @@ -222,42 +353,41 @@ for (const kind of ["h1", "h1-chunked", "h1-gzip", "h1-tls", "h2", "h3"] as Kind ["no consumer", STALL_NO_CONSUMER], ] as const); for (const [name, script] of scripts) { - // Subprocess RSS is too noisy to assert a bound across CI hosts (JIT - // warmup + mimalloc chunks + TLS dylib faulting exceed the 16 MiB - // body on several lanes). These assert the resume path drains the - // full body with no deadlock; the in-process "server stops writing" - // tests below prove the pause. + // How far the server got at the stall is not asserted: loopback buffers on some CI hosts take + // all 16 MiB. These assert that the resume path delivers the whole body intact, that the + // server wrote all of it, and that the client exits cleanly; the in-process "server stops + // writing" tests below prove the pause. test.skipIf(skip)(`stalled ${name} drains the full body`, async () => { await using server = await serve(kind); - const { peak, total, exitCode } = await spawnClient(server.url, kind, script); - expect({ peakMB: peak >> 20, total }).toEqual({ peakMB: expect.any(Number), total: TOTAL }); + const { total, digest, stderr, exitCode } = await spawnClient(server, kind, script); + expect({ total, digest, sent: server.sent(), stderr }).toEqual({ + total: TOTAL, + digest: DIGEST, + sent: server.wire, + stderr: "", + }); expect(exitCode).toBe(0); }); } if (kind === "h1" || kind === "h1-chunked" || kind === "h1-tls") { - test("server stops writing while the reader is stalled, then drains", async () => { - // Body must exceed kernel loopback send+recv autotuning. Some CI - // hosts have tcp_rmem[2]+tcp_wmem[2] approaching 256 MiB, so use - // 1 GiB; the server only actually writes until it blocks. - const big = 16384; - await using server = await serve(kind, big); + test("server stops writing while the reader is stalled, then resumes", async () => { + await using server = await serve(kind, BIG); const res = await fetch(server.url, fetchOpts(kind)); const reader = res.body!.getReader(); const first = await reader.read(); - let last = -1; - let stable = 0; - while (stable < 2) { - await Bun.sleep(10); - const now = server.sent(); - stable = now === last ? stable + 1 : 0; - last = now; - } - expect(server.sent()).toBeLessThan(CHUNK * big); - let total = first.value!.byteLength; - for (let r; !(r = await reader.read()).done; ) total += r.value.byteLength; - expect({ sent: server.sent(), total }).toEqual({ sent: CHUNK * big, total: CHUNK * big }); - }, 60_000); + // One chunk out and no read pending: the client pauses, and the server runs into the + // kernel's buffers. + const stalled = await server.settled(); + expect(stalled).toBeLessThan(BODY); + const { total, foreign, ended } = await readUntilResumed(reader, server, stalled); + expect({ ended, foreign, firstIsA: isAllA(first.value!), took: total > 0 && total <= server.sent() }).toEqual({ + ended: false, + foreign: 0, + firstIsA: true, + took: true, + }); + }); } }); } @@ -266,39 +396,41 @@ describe.concurrent("fetch() receive backpressure — Readable.fromWeb bridge", // `Readable.fromWeb(res.body)` takes the native handle off the ReadableStream // (NativeReadable fast path) and hands chunks to node streams. A stalled // pipe must keep the HTTP-thread socket paused just like `getReader()` does; - // when the pipe resumes, the body must drain to completion. - test("server stops writing while Readable.fromWeb is piped to a stalled Writable, then drains", async () => { - const big = 16384; - await using server = await serve("h1", big); + // when the pipe resumes, the body must keep coming. + test("server stops writing while Readable.fromWeb is piped to a stalled Writable, then resumes", async () => { + await using server = await serve("h1", BIG); const res = await fetch(server.url); let release!: () => void; let got = 0; + let foreign = 0; + let stalled = -1; + const enough = new Error("enough"); const sink = new Writable({ write(chunk, _enc, cb) { got += chunk.length; - if (release) return cb(); - release = cb; + if (!isAllA(chunk)) foreign++; + // The first chunk stalls the pipe until the test releases it. + if (!release) return void (release = cb); + // Resumed: the server wrote a MiB past where the stall left it. + if (server.sent() >= stalled + 16 * CHUNK) return cb(enough); + cb(); }, }); - const readable = Readable.fromWeb(res.body!); - const done = pipeline(readable, sink).then( + const done = pipeline(Readable.fromWeb(res.body!), sink).then( () => null, e => e, ); - let last = -1; - let stable = 0; - while (stable < 2) { - await Bun.sleep(10); - const now = server.sent(); - stable = now === last ? stable + 1 : 0; - last = now; - } - expect(server.sent()).toBeLessThan(CHUNK * big); + stalled = await server.settled(); + expect(stalled).toBeLessThan(BODY); release(); - expect({ err: await done, sent: server.sent(), got }).toEqual({ err: null, sent: CHUNK * big, got: CHUNK * big }); - }, 60_000); + expect({ err: await done, foreign, took: got > 0 && got <= server.sent() }).toEqual({ + err: enough, + foreign: 0, + took: true, + }); + }); // The buffered window between the HTTP-thread recv and `res.write()` is a // chain of native Vecs (FetchTasklet staging + ByteStream overflow) that are @@ -359,7 +491,7 @@ describe.concurrent("fetch() receive backpressure — Readable.fromWeb bridge", // 40 connections × 32 MB each (1.25 GB total). Before the window fix the // per-connection staging/overflow capacity pushed peak RSS to ~247–272 MB // on Linux; with it the same run sits at ~158–187 MB. - expect({ short, peakMB }).toEqual({ short: 0, peakMB: expect.any(Number) }); + expect({ short, peakMB, stderr }).toEqual({ short: 0, peakMB: expect.any(Number), stderr: "" }); expect(peakMB).toBeLessThan(225); expect(exitCode).toBe(0); }, @@ -372,44 +504,50 @@ describe.concurrent("fetch() receive backpressure — Readable.fromWeb bridge", // RSS bound for h2 needs that window lowered, which is a separate change. describe.concurrent("fetch() receive backpressure — buffered consumers are not throttled", () => { - const cases = [ - ["res.arrayBuffer()", async (r: Response) => (await r.arrayBuffer()).byteLength], - ["res.bytes()", async (r: Response) => (await r.bytes()).byteLength], - ["res.text()", async (r: Response) => (await r.text()).length], - ["res.blob()", async (r: Response) => (await r.blob()).size], - ["res.body.bytes()", async (r: Response) => (await r.body!.bytes()).byteLength], - ["res.body.text()", async (r: Response) => (await r.body!.text()).length], - ["res.body.blob()", async (r: Response) => (await r.body!.blob()).size], - [ - "res.body.json() rejects on full body", - async (r: Response) => - r.body!.json().then( - () => 0, - () => TOTAL, - ), - ], + const cases: [string, (r: Response) => Promise][] = [ + ["res.arrayBuffer()", async r => md5(new Uint8Array(await r.arrayBuffer()))], + ["res.bytes()", async r => md5(await r.bytes())], + ["res.text()", async r => md5(await r.text())], + ["res.blob()", async r => md5(await (await r.blob()).bytes())], + ["res.body.bytes()", async r => md5(await r.body!.bytes())], + ["res.body.text()", async r => md5(await r.body!.text())], + ["res.body.blob()", async r => md5(await (await r.body!.blob()).bytes())], [ "Bun.readableStreamToArrayBuffer(res.body)", - async (r: Response) => (await Bun.readableStreamToArrayBuffer(r.body!)).byteLength, + async r => md5(new Uint8Array(await Bun.readableStreamToArrayBuffer(r.body!))), ], - [ - "Bun.readableStreamToBytes(res.body)", - async (r: Response) => (await Bun.readableStreamToBytes(r.body!)).byteLength, - ], - ["Bun.readableStreamToText(res.body)", async (r: Response) => (await Bun.readableStreamToText(r.body!)).length], - ["Bun.readableStreamToBlob(res.body)", async (r: Response) => (await Bun.readableStreamToBlob(r.body!)).size], + ["Bun.readableStreamToBytes(res.body)", async r => md5(await Bun.readableStreamToBytes(r.body!))], + ["Bun.readableStreamToText(res.body)", async r => md5(await Bun.readableStreamToText(r.body!))], + ["Bun.readableStreamToBlob(res.body)", async r => md5(await (await Bun.readableStreamToBlob(r.body!)).bytes())], [ "Bun.readableStreamToArray(res.body)", - async (r: Response) => (await Bun.readableStreamToArray(r.body!)).reduce((n, c) => n + c.byteLength, 0), + async r => { + const hasher = new Bun.CryptoHasher("md5"); + for (const chunk of await Bun.readableStreamToArray(r.body!)) hasher.update(chunk); + return hasher.digest("hex"); + }, ], - ] as const; + ]; for (const [name, consume] of cases) { test(name, async () => { await using server = await serve("h1"); - expect(await consume(await fetch(server.url))).toBe(TOTAL); + expect({ digest: await consume(await fetch(server.url)), sent: server.sent() }).toEqual({ + digest: DIGEST, + sent: TOTAL, + }); }); } + + // 16 MiB of "A" is not JSON. The parse only runs once the whole body is in. + test("res.body.json() rejects on full body", async () => { + await using server = await serve("h1"); + const err = await (await fetch(server.url)).body!.json().then( + () => null, + e => e, + ); + expect({ name: err?.name, sent: server.sent() }).toEqual({ name: "SyntaxError", sent: TOTAL }); + }); }); describe.concurrent("fetch() receive backpressure — streaming consumer shapes", () => { @@ -422,21 +560,26 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" // reader.cancel() aborts the in-flight request (#33227), closing the // connection; the client must recover so a later request still completes. // The abort-vs-drain behavior itself is asserted in regression/issue/33227. - const buf = await (await fetch(server.url, { keepalive: true })).arrayBuffer(); - expect(buf.byteLength).toBe(TOTAL); + const body = await (await fetch(server.url, { keepalive: true })).bytes(); + expect({ length: body.byteLength, digest: md5(body) }).toEqual({ length: TOTAL, digest: DIGEST }); }); test("res.body.tee() both branches drain", async () => { await using server = await serve("h1"); const [a, b] = (await fetch(server.url)).body!.tee(); const sum = async (s: ReadableStream) => { + const hasher = new Bun.CryptoHasher("md5"); let n = 0; - for await (const c of s) n += c.byteLength; - return n; + for await (const c of s) { + n += c.byteLength; + hasher.update(c); + } + return { n, digest: hasher.digest("hex") }; }; - const [na, nb] = await Promise.all([sum(a), sum(b)]); - expect(na).toBe(TOTAL); - expect(nb).toBe(TOTAL); + expect(await Promise.all([sum(a), sum(b)])).toEqual([ + { n: TOTAL, digest: DIGEST }, + { n: TOTAL, digest: DIGEST }, + ]); }); // The peer dying while the transport is receive-paused is only observable @@ -451,14 +594,13 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" // Declares far more than it will send and writes until the kernel stops taking it, which, // with an untouched body, is once the client holds the high-water mark and has paused. const declared = 1 << 30; - const payload = Buffer.alloc(CHUNK, 65); const blocked = Promise.withResolvers(); let sent = 0; const push = (s: import("bun").Socket) => { while (sent < declared) { - const n = s.write(payload); + const n = s.write(PAYLOAD); sent += Math.max(n, 0); - if (n < payload.length) return void (sent > 4 * CHUNK && blocked.resolve(s)); + if (n < PAYLOAD.length) return void (sent > 4 * CHUNK && blocked.resolve(s)); } }; using listener = Bun.listen({ @@ -477,13 +619,24 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" kill(await blocked.promise); const reader = res.body!.getReader(); let total = 0; + let foreign = 0; const err = await (async () => { - for (let r; !(r = await reader.read()).done; ) total += r.value.byteLength; + for (let r; !(r = await reader.read()).done; ) { + total += r.value.byteLength; + if (!isAllA(r.value)) foreign++; + } })().then( () => null, e => e, ); - expect({ code: err?.code, partial: total < declared }).toEqual({ code: "ECONNRESET", partial: true }); + expect({ name: err?.name, code: err?.code, message: err?.message, partial: total < declared, foreign }).toEqual({ + name: "TypeError", + code: "ECONNRESET", + message: + "The socket connection was closed unexpectedly. For more information, pass `verbose: true` in the second argument to fetch()", + partial: true, + foreign: 0, + }); }); } @@ -492,11 +645,17 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" for (let i = 0; i < 2; i++) { const reader = (await fetch(server.url, { keepalive: true })).body!.getReader(); const first = await reader.read(); - await Bun.sleep(20); + // Hold the first chunk until the server has run into the pause (or written everything). + await server.settled(); + const hasher = new Bun.CryptoHasher("md5").update(first.value!); let total = first.value!.byteLength; - for (let r; !(r = await reader.read()).done; ) total += r.value.byteLength; - expect(total).toBe(TOTAL); + for (let r; !(r = await reader.read()).done; ) { + total += r.value.byteLength; + hasher.update(r.value); + } + expect({ total, digest: hasher.digest("hex") }).toEqual({ total: TOTAL, digest: DIGEST }); } + expect(server.sent()).toBe(2 * TOTAL); }); }); @@ -506,69 +665,6 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" // while it did. The transport has to stay paused instead (the bytes have nowhere to go), // the stream has to stay readable, and the paused body must not hold an idle process. -const BIG = 16384; // 1 GiB, as in the "server stops writing" tests above -const BODY = BIG * CHUNK; - -// Writes the body as fast as the socket takes it; `sent()` is how far it got. A paused -// client stops it at the socket buffers (up to a few hundred MiB on some hosts), a client -// that drains lets it write all of BODY. -async function serveUntilBlocked() { - let sent = 0; - let closed = 0; - const payload = Buffer.alloc(CHUNK, 65); - // The kernel stopped taking writes with more than the client's high-water mark outstanding: - // from here only a reader can make room. - const blocked = Promise.withResolvers(); - const firstClosed = Promise.withResolvers(); - const srv = createServer((_req, res) => { - res.on("error", () => {}); - res.on("close", () => { - closed++; - firstClosed.resolve(); - }); - res.flushHeaders(); - let i = 0; - const push = () => { - while (i < BIG && !res.destroyed) { - i++; - sent += CHUNK; - if (!res.write(payload)) { - if (i > 8) blocked.resolve(); - return void res.once("drain", push); - } - } - if (!res.destroyed) res.end(); - }; - push(); - }); - srv.listen(0); - await once(srv, "listening"); - const { port } = srv.address() as import("node:net").AddressInfo; - return { - url: `http://127.0.0.1:${port}/`, - sent: () => sent, - // Resolves with `sent()` once it has not moved for 100ms. - async settled() { - let last = -1; - for (let stable = 0; stable < 5; ) { - await Bun.sleep(20); - stable = sent === last ? stable + 1 : 0; - last = sent; - } - return sent; - }, - async untilClosed() { - while (closed === 0) await Bun.sleep(5); - }, - blocked: blocked.promise, - closed: firstClosed.promise, - [Symbol.asyncDispose]: () => { - srv.closeAllConnections(); - return new Promise(r => srv.close(() => r(undefined))); - }, - }; -} - // Sends the headers and one chunk; the test writes the rest through `response`. async function serveByHand() { let respond!: (res: import("node:http").ServerResponse) => void; @@ -576,7 +672,7 @@ async function serveByHand() { const srv = createServer((_req, res) => { res.on("error", () => {}); res.flushHeaders(); - res.write(Buffer.alloc(CHUNK, 65)); + res.write(PAYLOAD); respond(res); }); srv.listen(0); @@ -592,9 +688,7 @@ async function serveByHand() { }; } -// Sequential on purpose: each test watches one server's write progress, and a sibling -// test draining a body on the same HTTP thread would blur that signal. -describe("fetch() receive backpressure — body stream nothing is reading", () => { +describe.concurrent("fetch() receive backpressure — body stream nothing is reading", () => { const shapes: [string, (res: Response) => Promise>][] = [ ["res.body read by nothing", async res => res.body!], [ @@ -610,73 +704,28 @@ describe("fetch() receive backpressure — body stream nothing is reading", () = for (const [name, shape] of shapes) { test(`${name}: the server blocks, and a later reader resumes the body`, async () => { - await using server = await serveUntilBlocked(); + await using server = await serve("h1-chunked", BIG); const body = await shape(await fetch(server.url)); - expect(await server.settled()).toBeLessThan(BODY); - - const reader = body.getReader(); - let got = 0; - while (got < 16 * CHUNK) got += (await reader.read()).value!.byteLength; + const stalled = await server.settled(); + expect(stalled).toBeLessThan(BODY); - await reader.cancel(); - await server.untilClosed(); - }, 60_000); + const { total, foreign, ended } = await readUntilResumed(body.getReader(), server, stalled); + expect({ ended, foreign, took: total > 0 && total <= server.sent() }).toEqual({ + ended: false, + foreign: 0, + took: true, + }); + // cancel() closed the connection. + await server.closed; + }); } - // The other half of the rule: a small body nobody reads is still taken off the socket, so - // the keep-alive connection goes back to the pool instead of staying pinned under it. - test("small unread bodies complete and their connection is reused", async () => { - // The body trails the headers in several packets, so it reaches a stream nothing reads - // chunk by chunk; the client has to keep taking it (it is under the high-water mark) for - // the response to finish and the connection to go back to the pool. A client that parks - // on the first unread chunk needs a new connection for every request here. - const PART = 32 * 1024; - const PARTS = 4; - let connections = 0; - let finished = Promise.resolve(); - const srv = createServer(async (_req, res) => { - finished = new Promise(r => res.on("finish", () => r())); - res.setHeader("content-length", String(PART * PARTS)); - res.flushHeaders(); - for (let i = 0; i < PARTS; i++) { - await new Promise(r => setTimeout(r, 2)); - res.write(Buffer.alloc(PART, 65)); - } - res.end(); - }).on("connection", () => connections++); - srv.listen(0); - await once(srv, "listening"); - try { - const url = `http://127.0.0.1:${(srv.address() as import("node:net").AddressInfo).port}/`; - const N = 10; - for (let i = 0; i < N; i++) { - const res = await fetch(url); - expect(res.status).toBe(200); - void res.body; - await finished; - } - // Not exactly 1: a request can start before the previous body's last packet was taken. - expect(connections).toBeLessThan(N / 2); - } finally { - srv.closeAllConnections(); - await new Promise(r => srv.close(() => r(undefined))); - } - }); - // The other ways back to a parked body: a buffered read and a native sink both have to // pick the transport up again and see the whole body. for (const [name, drain] of [ - ["res.text()", (res: Response) => res.text().then(t => t.length)], - [ - "HTMLRewriter.transform(res)", - (res: Response) => - new HTMLRewriter() - .on("x", {}) - .transform(res) - .arrayBuffer() - .then(b => b.byteLength), - ], + ["res.text()", (res: Response) => res.text().then(md5)], + ["HTMLRewriter.transform(res)", (res: Response) => new HTMLRewriter().on("x", {}).transform(res).bytes().then(md5)], ] as const) { test(`res.body read by nothing, then ${name} drains the whole body`, async () => { await using server = await serve("h1"); @@ -684,13 +733,8 @@ describe("fetch() receive backpressure — body stream nothing is reading", () = void res.body; // Let the client park (16 MiB body, 256 KiB mark): wait until bytes stop moving. How far // the server got is not asserted; loopback buffers on some hosts can take the whole body. - let last = -1; - for (let stable = 0; stable < 3; ) { - await Bun.sleep(20); - stable = server.sent() === last ? stable + 1 : 0; - last = server.sent(); - } - expect(await drain(res)).toBe(TOTAL); + await server.settled(); + expect({ digest: await drain(res), sent: server.sent() }).toEqual({ digest: DIGEST, sent: TOTAL }); }); } }); @@ -708,11 +752,10 @@ type Framing = "content-length" | "chunked" | "close-delimited"; // the last CHUNK of every body is held back until `finishHeld()`, which also stops holding. // The origin never ends a body by closing, so every close it sees is the client's. async function rawOrigin(framing: Framing, length: number, holdTail = false) { - const payload = Buffer.alloc(CHUNK, 65); const frame = framing === "chunked" - ? Buffer.concat([Buffer.from(`${CHUNK.toString(16)}\r\n`), payload, Buffer.from("\r\n")]) - : payload; + ? Buffer.concat([Buffer.from(`${CHUNK.toString(16)}\r\n`), PAYLOAD, Buffer.from("\r\n")]) + : PAYLOAD; const tail = framing === "chunked" ? Buffer.concat([frame, Buffer.from("0\r\n\r\n")]) : frame; const head = framing === "content-length" @@ -802,12 +845,11 @@ async function collectUntil(event: Promise): Promise { return result; } -// Sequential on purpose, as above: these watch connections and closes on their own origin. -describe("fetch() receive backpressure — a Response whose body nothing touches", () => { +describe.concurrent("fetch() receive backpressure — a Response whose body nothing touches", () => { const N = 4; test("a long body, Response held untouched: the process is not held and the body stays where it is", async () => { - await using server = await serveUntilBlocked(); + await using server = await serve("h1-chunked", BIG); await using proc = Bun.spawn({ cmd: [bunExe(), "-e", `globalThis.keep = await fetch(${JSON.stringify(server.url)});`], env: bunEnv, @@ -831,27 +873,6 @@ describe("fetch() receive backpressure — a Response whose body nothing touches }); } - // Not close-delimited: such a body ends with its connection, so there is nothing to reuse. - for (const framing of ["content-length", "chunked"] as Framing[]) { - test(`a short ${framing} body, Response still held: it is received, and its connection is reused`, async () => { - // Each body's tail is held back until its Response exists, so every body is still underway - // when fetch() resolves, as it is over a real network. - await using origin = await rawOrigin(framing, 2 * CHUNK, true); - const responses: Response[] = []; - for (let i = 0; i < N; i++) responses.push(await fetch(origin.url)); - // Every body is underway, so no connection was free for the next request. - expect(origin.connections()).toBe(N); - - origin.finishHeld(); - // The held bodies complete on their own and give their connections back. One request after - // another from here needs at most one more connection (the first can leave before the tails - // were taken); before, every one of them did, since each held body pinned its connection. - for (let i = 0; i < N; i++) expect((await (await fetch(origin.url)).arrayBuffer()).byteLength).toBe(2 * CHUNK); - expect(origin.connections() - N).toBeLessThanOrEqual(1); - expect({ closed: origin.closed(), held: responses.length }).toEqual({ closed: 0, held: N }); - }); - } - // The boundary of the abort above: a consumer that waits for the whole body (`.text()` through // a promise, `Bun.write()` through a native callback) may be all that is left of a Response. // Its body still has to arrive. @@ -887,10 +908,10 @@ describe("fetch() receive backpressure — a Response whose body nothing touches // S3 downloads go through the same HTTP client with their own body producer // (S3DownloadStreamWrapper). The same rule applies: a reader that stalls pauses the transport, // an unread stream does not hold the process, and a collected one aborts the download. -describe("S3 receive backpressure", () => { +describe.concurrent("S3 receive backpressure", () => { // A GET-only fake bucket: every object is BODY bytes written as fast as the socket takes them. async function fakeBucket() { - const server = await serveUntilBlocked(); + const server = await serve("h1-chunked", BIG); const s3 = new S3Client({ accessKeyId: "test", secretAccessKey: "test", endpoint: server.url, bucket: "b" }); return Object.assign(server, { s3 }); } @@ -898,10 +919,16 @@ describe("S3 receive backpressure", () => { test("a reader that stalls and comes back drains the body; cancel() closes the connection", async () => { await using bucket = await fakeBucket(); const reader = bucket.s3.file("big").stream().getReader(); - let got = (await reader.read()).value!.byteLength; - await bucket.blocked; - while (got < 64 * CHUNK) got += (await reader.read()).value!.byteLength; - await reader.cancel(); + const first = await reader.read(); + const stalled = await bucket.settled(); + expect(stalled).toBeLessThan(BODY); + const { total, foreign, ended } = await readUntilResumed(reader, bucket, stalled); + expect({ ended, foreign, firstIsA: isAllA(first.value!), took: total > 0 && total <= bucket.sent() }).toEqual({ + ended: false, + foreign: 0, + firstIsA: true, + took: true, + }); await bucket.closed; }); @@ -911,7 +938,11 @@ describe("S3 receive backpressure", () => { const s3 = new S3Client({ accessKeyId: "test", secretAccessKey: "test", endpoint: server.url, bucket: "b" }); const dest = join(String(dir), "out.bin"); expect(await Bun.write(dest, s3.file("k"))).toBe(TOTAL); - expect(statSync(dest).size).toBe(TOTAL); + expect({ size: statSync(dest).size, digest: md5(await Bun.file(dest).bytes()), sent: server.sent() }).toEqual({ + size: TOTAL, + digest: DIGEST, + sent: TOTAL, + }); }); // A paused, unread stream releases the loop: the process exits with most of the body unsent. @@ -943,7 +974,7 @@ describe("S3 receive backpressure", () => { await r.read(); r.releaseLock(); })(); - await bucket.blocked; + await bucket.settled(); await collectUntil(bucket.closed); expect(bucket.sent()).toBeLessThan(BODY); }); @@ -957,15 +988,27 @@ describe("S3 receive backpressure", () => { new Response(`AccessDenied${message}`, { status: 403 }), }); const s3 = new S3Client({ accessKeyId: "t", secretAccessKey: "t", endpoint: bucket.url.href, bucket: "b" }); - await expect(s3.file("denied").stream().getReader().read()).rejects.toThrow( - expect.objectContaining({ code: "AccessDenied" }), - ); + const err = await s3 + .file("denied") + .stream() + .getReader() + .read() + .then( + () => null, + e => e, + ); + expect({ name: err?.name, code: err?.code, message: err?.message }).toEqual({ + name: "S3Error", + code: "AccessDenied", + message, + }); }); // The other direction: a fetch body uploaded to S3. The multipart sink's queue back-pressures // the fetch, and both `Bun.write(s3file, res)` and `s3file.writer()` resolve with the bytes sent. async function fakeUploadBucket(holdParts?: Promise) { let uploaded = 0; + let foreign = 0; let parts = 0; let completed = 0; const firstPart = Promise.withResolvers(); @@ -976,10 +1019,11 @@ describe("S3 receive backpressure", () => { if (req.method === "POST" && url.searchParams.has("uploads")) return new Response("u"); if (req.method === "PUT") { - const body = await req.arrayBuffer(); + const body = await req.bytes(); firstPart.resolve(); await holdParts; uploaded += body.byteLength; + if (!isAllA(body)) foreign++; return new Response("", { headers: { etag: `"e${++parts}"` } }); } if (req.method === "POST" && url.searchParams.has("uploadId")) { @@ -998,13 +1042,14 @@ describe("S3 receive backpressure", () => { s3, firstPart: firstPart.promise, uploaded: () => uploaded, + foreign: () => foreign, completed: () => completed, }); } test("fetch → Bun.write(s3file, res) is paced by the part uploads, and an aborted source commits nothing", async () => { const hold = Promise.withResolvers(); - await using origin = await serveUntilBlocked(); + await using origin = await serve("h1-chunked", BIG); await using bucket = await fakeUploadBucket(hold.promise); const abort = new AbortController(); // partSize 5 MiB × queueSize 1: with the first part held, the sink fills and the origin has @@ -1014,20 +1059,30 @@ describe("S3 receive backpressure", () => { await fetch(origin.url, { signal: abort.signal }), ); await bucket.firstPart; - await origin.blocked; - expect(origin.sent()).toBeLessThan(BODY); + expect(await origin.settled()).toBeLessThan(BODY); // Aborting the source fails the upload: nothing is committed. abort.abort(); hold.resolve(); - await expect(written).rejects.toThrow(expect.objectContaining({ name: "AbortError" })); - expect(bucket.completed()).toBe(0); + const err = await written.then( + () => null, + e => e, + ); + expect({ name: err?.name, message: err?.message, completed: bucket.completed() }).toEqual({ + name: "AbortError", + message: "The operation was aborted.", + completed: 0, + }); }); test("Bun.write(s3file, res) resolves with the byte count", async () => { await using origin = await serve("h1"); await using bucket = await fakeUploadBucket(); expect(await Bun.write(bucket.s3.file("up"), await fetch(origin.url))).toBe(TOTAL); - expect(bucket.uploaded()).toBe(TOTAL); + expect({ uploaded: bucket.uploaded(), foreign: bucket.foreign(), completed: bucket.completed() }).toEqual({ + uploaded: TOTAL, + foreign: 0, + completed: 1, + }); }); test("s3file.writer().end() resolves with the byte count, also once the writer is collected", async () => { @@ -1085,7 +1140,7 @@ describe.concurrent("fetch() receive backpressure — a body nothing waits for d for (const [holding, script] of idleClients) { test(`a process holding ${holding} exits on its own`, async () => { - await using server = await serveUntilBlocked(); + await using server = await serve("h1-chunked", BIG); await using proc = Bun.spawn({ cmd: [bunExe(), "-e", `const res = await fetch(${JSON.stringify(server.url)}); ${script} ${diagnose}`], env: bunEnv, @@ -1093,18 +1148,18 @@ describe.concurrent("fetch() receive backpressure — a body nothing waits for d stderr: "pipe", }); const stderr = proc.stderr.text(); - const exited = proc.exited.then(exitCode => ({ exitCode })); - const deadline = performance.now() + 10_000; - let outcome: { exitCode: number } | { stillAlive: true; serverSentKiB: number } | undefined; - while (outcome === undefined) { - outcome = await Promise.race([exited, Bun.sleep(20).then(() => undefined)]); - // Took the whole body: it is draining. Took little and is still around: it is - // holding the paused body. Neither exited. - if (outcome === undefined && (server.sent() >= BODY || performance.now() >= deadline)) { - outcome = { stillAlive: true, serverSentKiB: server.sent() >> 10 }; - proc.kill(); - } - } + // Took the whole body: it is draining. Still around after 10 s with little taken: it is + // holding the paused body. Neither exited. + let deadline!: ReturnType; + const outcome = await Promise.race([ + proc.exited.then(exitCode => ({ exitCode })), + server.wrote(BODY).then(() => ({ stillAlive: true, drained: true })), + new Promise<{ stillAlive: true; serverSentKiB: number }>(resolve => { + deadline = setTimeout(() => resolve({ stillAlive: true, serverSentKiB: server.sent() >> 10 }), 10_000); + }), + ]); + clearTimeout(deadline); + if ("stillAlive" in outcome) proc.kill(); expect({ outcome, stderr: await stderr }).toEqual({ outcome: { exitCode: 0 }, stderr: "" }); }, 20_000); } @@ -1131,9 +1186,9 @@ describe.concurrent("fetch() receive backpressure — a body nothing waits for d let total = (await reader.read()).value.byteLength; reader.releaseLock(); console.log("released"); - // Lets the chunk the test writes on "released" arrive while nothing reads the body. - // A slow machine only narrows what this exercises; the assertions do not depend on it. - await Bun.sleep(50); + // The test writes a chunk on "released" and says "go" once the socket took it: that chunk + // arrives while nothing reads the body. + for await (const line of console) if (line === "go") break; reader = res.body.getReader(); const next = reader.read(); console.log("waiting"); @@ -1150,22 +1205,21 @@ describe.concurrent("fetch() receive backpressure — a body nothing waits for d await using proc = Bun.spawn({ cmd: [bunExe(), "-e", `const res = await fetch(${JSON.stringify(server.url)}); ${script}`], env: bunEnv, + stdin: "pipe", stdout: "pipe", stderr: "pipe", }); const stderr = proc.stderr.text(); const upstream = await server.response; const lines: string[] = []; - let pending = ""; - for await (const text of proc.stdout.pipeThrough(new TextDecoderStream())) { - pending += text; - for (let nl; (nl = pending.indexOf("\n")) !== -1; ) { - const line = pending.slice(0, nl); - pending = pending.slice(nl + 1); - lines.push(line); - if (line === "released") upstream.write(Buffer.alloc(CHUNK, 66)); - if (line === "waiting") upstream.end(Buffer.alloc(CHUNK, 67)); - } + for await (const line of forEachLine(proc.stdout)) { + lines.push(line); + if (line === "released") + upstream.write(Buffer.alloc(CHUNK, 66), () => { + proc.stdin.write("go\n"); + proc.stdin.end(); + }); + if (line === "waiting") upstream.end(Buffer.alloc(CHUNK, 67)); } expect({ lines: lines.slice(-2), stderr: await stderr, exitCode: await proc.exited }).toEqual({ lines: ["waiting", `total ${chunks * CHUNK}`], @@ -1175,3 +1229,76 @@ describe.concurrent("fetch() receive backpressure — a body nothing waits for d }); } }); + +// Serial on purpose: both count the connections a pooled client opens, and a connection is back in +// the pool only once the HTTP thread has taken the last of its body. A sibling test keeping that +// thread busy would turn a late return into an extra connection. +describe.serial("fetch() receive backpressure — an unread body hands its connection back", () => { + // A small body nobody reads is still taken off the socket, so the keep-alive connection goes + // back to the pool instead of staying pinned under it. + test("small unread bodies complete and their connection is reused", async () => { + // The body trails the headers in several packets; the last one is held until the Response + // exists and its body is a stream nothing reads, so it reaches such a stream chunk by chunk. + // The client has to keep taking it (it is under the high-water mark) for the response to + // finish and the connection to go back to the pool. A client that parks on the first unread + // chunk needs a new connection for every request here. + const PART = 32 * 1024; + const PARTS = 4; + let connections = 0; + let lastPart = Promise.withResolvers(); + let finished = Promise.withResolvers(); + const srv = createServer(async (_req, res) => { + res.on("finish", () => finished.resolve()); + res.setHeader("content-length", String(PART * PARTS)); + res.flushHeaders(); + for (let i = 0; i < PARTS - 1; i++) await new Promise(r => res.write(Buffer.alloc(PART, 65), r)); + await lastPart.promise; + res.end(Buffer.alloc(PART, 65)); + }).on("connection", () => connections++); + srv.listen(0); + await once(srv, "listening"); + try { + const url = `http://127.0.0.1:${(srv.address() as import("node:net").AddressInfo).port}/`; + const N = 10; + for (let i = 0; i < N; i++) { + lastPart = Promise.withResolvers(); + finished = Promise.withResolvers(); + const res = await fetch(url); + expect(res.status).toBe(200); + void res.body; + lastPart.resolve(); + await finished.promise; + } + // Not exactly 1: a request can start before the previous body's last packet was taken. + expect(connections).toBeLessThan(N / 2); + } finally { + srv.closeAllConnections(); + await new Promise(r => srv.close(() => r(undefined))); + } + }); + + // Not close-delimited: such a body ends with its connection, so there is nothing to reuse. + for (const framing of ["content-length", "chunked"] as Framing[]) { + test(`a short ${framing} body, Response still held: it is received, and its connection is reused`, async () => { + const N = 4; + // Each body's tail is held back until its Response exists, so every body is still underway + // when fetch() resolves, as it is over a real network. + await using origin = await rawOrigin(framing, 2 * CHUNK, true); + const responses: Response[] = []; + for (let i = 0; i < N; i++) responses.push(await fetch(origin.url)); + // Every body is underway, so no connection was free for the next request. + expect(origin.connections()).toBe(N); + + origin.finishHeld(); + // The held bodies complete on their own and give their connections back. One request after + // another from here needs at most one more connection (the first can leave before the tails + // were taken); before, every one of them did, since each held body pinned its connection. + for (let i = 0; i < N; i++) { + const body = await (await fetch(origin.url)).bytes(); + expect({ length: body.byteLength, allA: isAllA(body) }).toEqual({ length: 2 * CHUNK, allA: true }); + } + expect(origin.connections() - N).toBeLessThanOrEqual(1); + expect({ closed: origin.closed(), held: responses.length }).toEqual({ closed: 0, held: N }); + }); + } +});