From 532ad123114a3033ffa0bce9150cc8f09c61e373 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 24 Sep 2026 17:16:05 +0000 Subject: [PATCH 1/2] Bun.serve: do not end a Response body stream that another reader holds A locked stream belongs to its lock holder: a reader the handler took, or the server's own sink for another Response around the same stream. Two places ended such a stream. cancel_unread_body cancels the body stream of a Response the server does not transmit: a HEAD response, a 101/204/205/304 response, and a handler result the server drops. It now cancels only a stream that is not locked. ReadableStream.prototype.cancel rejects on a locked stream for the same reason. GET refuses a locked stream body with ERR_STREAM_CANNOT_PIPE. The refusal left the stream attached to the Response, so the teardown of the refused request found it and ended it. The refusal now detaches the stream. It also no longer calls unprotect() on the stream, which no protect() matched. --- src/runtime/server/RequestContext.rs | 14 +- .../serve-pending-promise-abort-leak.test.ts | 205 ++++++++++++++++++ .../js/bun/http/serve-reused-response.test.ts | 118 ++++++++++ 3 files changed, 333 insertions(+), 4 deletions(-) diff --git a/src/runtime/server/RequestContext.rs b/src/runtime/server/RequestContext.rs index eee6f4f59b60..1208bca9ee01 100644 --- a/src/runtime/server/RequestContext.rs +++ b/src/runtime/server/RequestContext.rs @@ -736,13 +736,16 @@ where Ok(JSValue::UNDEFINED) } - /// Cancel the body stream of a Response the server will not transmit. + /// Cancel the body stream of a Response the server will not transmit, unless a reader holds it. fn cancel_unread_body(response: &Response, global_this: &JSGlobalObject) { if let Some(stream) = response.get_body_readable_stream() { let _keep = jsc::EnsureStillAlive(stream.value); response.detach_readable_stream(global_this); - // Not `cancel()`: it skips a stream with no reader, which an unattached body is. - crate::dispatch::fold(stream.cancel_with_reason(global_this, JSValue::UNDEFINED)); + // A locked stream belongs to its reader, which can be the sink of another Response. + if !stream.is_locked(global_this) { + // Not `cancel()`: it skips a stream with no reader, which an unattached body is. + crate::dispatch::fold(stream.cancel_with_reason(global_this, JSValue::UNDEFINED)); + } } *response.get_body_value() = Body::Value::Used; } @@ -3154,7 +3157,10 @@ where ), ..Default::default() }; - stream.value.unprotect(); + // Teardown must not find the stream: it belongs to its reader. + if let Some(response) = this.response_mut() { + response.detach_readable_stream(global_this); + } let js_err = err.to_error_instance(global_this); this.run_error_handler(js_err); return; diff --git a/test/js/bun/http/serve-pending-promise-abort-leak.test.ts b/test/js/bun/http/serve-pending-promise-abort-leak.test.ts index 6782ec905f07..fb37f3b0fc65 100644 --- a/test/js/bun/http/serve-pending-promise-abort-leak.test.ts +++ b/test/js/bun/http/serve-pending-promise-abort-leak.test.ts @@ -964,3 +964,208 @@ test.each([204, 304])("a %d Response with a ReadableStream body cancels the stre expect(cancelReasons).toStrictEqual([undefined]); await stopAndAssertDrained(server); }); + +// A locked stream belongs to its reader. The server does not transmit these +// bodies, but it is not the one reading them either, so it must not cancel +// them: ReadableStream.prototype.cancel() rejects on a locked stream for the +// same reason. +const heldReaderRows = [204, 205, 304].flatMap(status => ["GET", "HEAD"].map(method => [status, method] as const)); +test.each(heldReaderRows)( + "a %d Response (%s) does not cancel a body stream the handler still reads", + async (status, method) => { + let cancels = 0; + let controller!: ReadableStreamDefaultController; + let reader!: ReadableStreamDefaultReader; + using server = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch() { + const response = new Response( + new ReadableStream({ + start(c) { + controller = c; + c.enqueue(new TextEncoder().encode("part1;")); + }, + cancel() { + cancels++; + }, + }), + { status }, + ); + reader = response.body!.getReader(); + return response; + }, + }); + + const res = await fetch(server.url, { method }); + expect(res.status).toBe(status); + expect(await res.text()).toBe(""); + + // The response is complete. The reader the handler took still owns the stream. + expect(cancels).toBe(0); + controller.enqueue(new TextEncoder().encode("part2;")); + controller.close(); + let read = ""; + const decoder = new TextDecoder(); + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + read += decoder.decode(chunk.value, { stream: true }); + } + expect({ read, cancels }).toEqual({ read: "part1;part2;", cancels: 0 }); + await stopAndAssertDrained(server); + }, +); + +// Two Responses can adopt one stream while nothing has read it. While the first is on the wire, +// the server holds the stream's lock: through a reader, through the sink of a direct stream, or +// through the native pipe of a fetch() body. A second request that sends no body must leave it. +type Producing = AsyncDisposable & { stream: ReadableStream; finish(): void; cancels(): number }; +const producing: Record Promise> = { + "a ReadableStream": async () => { + let cancels = 0; + let controller!: ReadableStreamDefaultController; + const stream = new ReadableStream({ + start(c) { + controller = c; + c.enqueue(new TextEncoder().encode("part1;")); + }, + cancel() { + cancels++; + }, + }); + return { + stream, + finish() { + controller.enqueue(new TextEncoder().encode("part2;")); + controller.close(); + }, + cancels: () => cancels, + async [Symbol.asyncDispose]() {}, + }; + }, + 'a type: "direct" ReadableStream': async () => { + let cancels = 0; + const { promise: gate, resolve: finish } = Promise.withResolvers(); + const stream = new ReadableStream({ + type: "direct", + async pull(controller) { + controller.write("part1;"); + await controller.flush(); + await gate; + controller.write("part2;"); + controller.close(); + }, + cancel() { + cancels++; + }, + }); + return { stream, finish, cancels: () => cancels, async [Symbol.asyncDispose]() {} }; + }, + "a fetch() body": async () => { + let cancels = 0; + const { promise: gate, resolve: finish } = Promise.withResolvers(); + const upstream = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch: () => + new Response( + new ReadableStream({ + async start(c) { + c.enqueue(new TextEncoder().encode("part1;")); + await gate; + c.enqueue(new TextEncoder().encode("part2;")); + c.close(); + }, + // Runs when the fetch() behind `body` drops its connection. + cancel() { + cancels++; + }, + }), + ), + }); + const { body } = await fetch(upstream.url); + return { stream: body!, finish, cancels: () => cancels, [Symbol.asyncDispose]: () => upstream.stop(true) }; + }, +}; + +async function readToEnd(reader: ReadableStreamDefaultReader, decoder: TextDecoder) { + let text = ""; + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + text += decoder.decode(chunk.value, { stream: true }); + } + return text; +} + +test.each(Object.keys(producing))( + "a HEAD for another Response around %s that is being sent leaves that stream alone", + async kind => { + await using held = await producing[kind](); + const responses = { "/sent": new Response(held.stream), "/head": new Response(held.stream) }; + using server = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch: req => responses[new URL(req.url).pathname as keyof typeof responses], + }); + + const inFlight = await fetch(new URL("/sent", server.url)); + const body = inFlight.body!.getReader(); + const decoder = new TextDecoder(); + let read = decoder.decode((await body.read()).value, { stream: true }); + + const head = await fetch(new URL("/head", server.url), { method: "HEAD" }); + expect(await head.text()).toBe(""); + expect(held.cancels()).toBe(0); + + held.finish(); + read += await readToEnd(body, decoder); + expect({ read, cancels: held.cancels() }).toEqual({ read: "part1;part2;", cancels: 0 }); + await stopAndAssertDrained(server); + }, +); + +test.each(Object.keys(producing))( + "a Promise that settles after the client aborted leaves %s alone while another response sends it", + async kind => { + const { promise: handlerEntered, resolve: signalHandler } = Promise.withResolvers(); + const { promise: abortObserved, resolve: signalAbort } = Promise.withResolvers(); + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + let handlerResult: Promise | undefined; + await using held = await producing[kind](); + // Both adopt the stream while nothing has read it. + const sent = new Response(held.stream); + const late = new Response(held.stream); + + using server = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch(req) { + if (new URL(req.url).pathname === "/sent") return sent; + req.signal.addEventListener("abort", () => signalAbort(), { once: true }); + signalHandler(); + handlerResult = gate.then(() => late); + return handlerResult; + }, + }); + + const inFlight = await fetch(new URL("/sent", server.url)); + const body = inFlight.body!.getReader(); + const decoder = new TextDecoder(); + let read = decoder.decode((await body.read()).value, { stream: true }); + + const ac = new AbortController(); + const aborted = fetch(new URL("/late", server.url), { signal: ac.signal }).catch(() => {}); + await handlerEntered; + ac.abort(); + await aborted; + await abortObserved; + + // The server drops this Response. Its stream is the one the first request still sends. + openGate(); + await handlerResult!; + expect(held.cancels()).toBe(0); + + held.finish(); + read += await readToEnd(body, decoder); + expect({ read, cancels: held.cancels() }).toEqual({ read: "part1;part2;", cancels: 0 }); + await stopAndAssertDrained(server); + }, +); diff --git a/test/js/bun/http/serve-reused-response.test.ts b/test/js/bun/http/serve-reused-response.test.ts index 4c4772e76e6b..f870cb3eb003 100644 --- a/test/js/bun/http/serve-reused-response.test.ts +++ b/test/js/bun/http/serve-reused-response.test.ts @@ -224,4 +224,122 @@ describe("returning a Response with an already-used body", () => { // The error is reported like any other unhandled error thrown from the fetch handler. expect(exitCode).toBe(1); }); + + // The server refuses a stream body that another reader holds. The stream stays with that reader: + // the teardown of the refused request must not end it. The default 500 reports the refusal as an + // unhandled error, so each case runs in its own process. + describe.concurrent("a fetch() body that another reader holds", () => { + const streamUsed = "Stream already used, please create a new one"; + async function run(script: string) { + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", script], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + return { ...JSON.parse(stdout), reports: stderr.split(streamUsed).length - 1, exitCode }; + } + + const whole = "chunk1;chunk2;chunk3;chunk4;chunk5;chunk6;"; + it.each([ + ["GET", { status: 500, read: whole, reports: 1, exitCode: 1 }], + // HEAD sends no body, so the server has nothing to refuse. + ["HEAD", { status: 200, read: whole, reports: 0, exitCode: 0 }], + ] as const)("%s: the handler's reader still reads the whole body", async (method, expected) => { + const result = await run(` + const release = Promise.withResolvers(); + const upstream = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch() { + let sent = 0; + return new Response( + new ReadableStream({ + async pull(controller) { + if (sent === 1) await release.promise; + controller.enqueue(new TextEncoder().encode("chunk" + ++sent + ";")); + if (sent === 6) controller.close(); + }, + }), + ); + }, + }); + let reader; + let read = ""; + const decoder = new TextDecoder(); + const server = Bun.serve({ + port: 0, + idleTimeout: 0, + development: false, + // The handler reads the upstream body itself and still returns the upstream Response. + async fetch() { + const upstreamResponse = await fetch(upstream.url); + reader = upstreamResponse.body.getReader(); + read += decoder.decode((await reader.read()).value, { stream: true }); + return upstreamResponse; + }, + }); + const response = await fetch(server.url, { method: ${JSON.stringify(method)} }); + await response.text(); + // The request is over and torn down. The upstream sends the rest only now. + release.resolve(); + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + read += decoder.decode(chunk.value, { stream: true }); + } + console.log(JSON.stringify({ status: response.status, read })); + await server.stop(true); + await upstream.stop(true); + `); + expect(result).toEqual(expected); + }); + + it("GET: another client's download of the same body completes", async () => { + const result = await run(` + const release = Promise.withResolvers(); + const upstream = Bun.serve({ + port: 0, + idleTimeout: 0, + fetch: () => + new Response( + new ReadableStream({ + async start(controller) { + controller.enqueue(new TextEncoder().encode("part1;")); + await release.promise; + controller.enqueue(new TextEncoder().encode("part2;")); + controller.close(); + }, + }), + ), + }); + // Both adopt the body while nothing has read it. + const { body } = await fetch(upstream.url); + const responses = { "/sent": new Response(body), "/refused": new Response(body) }; + const server = Bun.serve({ + port: 0, + idleTimeout: 0, + development: false, + fetch: req => responses[new URL(req.url).pathname], + }); + const inFlight = await fetch(new URL("/sent", server.url)); + const reader = inFlight.body.getReader(); + const decoder = new TextDecoder(); + let read = decoder.decode((await reader.read()).value, { stream: true }); + const refused = await fetch(new URL("/refused", server.url)); + await refused.text(); + release.resolve(); + try { + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + read += decoder.decode(chunk.value, { stream: true }); + } + } catch (error) { + read += "<" + error.code + ">"; + } + console.log(JSON.stringify({ status: refused.status, read })); + await server.stop(true); + await upstream.stop(true); + `); + expect(result).toEqual({ status: 500, read: "part1;part2;", reports: 1, exitCode: 1 }); + }); + }); }); From b995e28d43081c5c17c7c1a3b9756c4a35cab6e9 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 25 Sep 2026 11:10:25 +0000 Subject: [PATCH 2/2] Bun.serve: HEAD sends a locked ReadableStream body to error(), as GET does GET refuses a Response whose body stream a reader already holds: it calls error() with ERR_STREAM_CANNOT_PIPE. The HEAD renderer had no such check and answered the handler's status with transfer-encoding: chunked. can_send_body_stream is the one place both renderers ask whether a stream body can be sent. refuse_body_stream builds the error() argument and sets the state a refusal leaves: the body is used and the Response lets go of the stream without a cancel. The HEAD arm asks before it writes any header. It hands the stream it resolved to cancel_unlocked_body, so a HEAD with a sendable stream resolves the stream once and asks for the lock once. The shared parts are free functions that are not inlined, so the eight RequestContext monomorphizations share one copy. --- src/runtime/server/RequestContext.rs | 82 ++++-- .../serve-pending-promise-abort-leak.test.ts | 7 + .../js/bun/http/serve-reused-response.test.ts | 269 +++++++++++++++++- 3 files changed, 327 insertions(+), 31 deletions(-) diff --git a/src/runtime/server/RequestContext.rs b/src/runtime/server/RequestContext.rs index 1208bca9ee01..d3c4127b6e1c 100644 --- a/src/runtime/server/RequestContext.rs +++ b/src/runtime/server/RequestContext.rs @@ -384,6 +384,46 @@ fn release_body_stream(response: &mut Response, global_this: &JSGlobalObject) { } } +/// Whether the server can send this stream body. GET and HEAD both ask here, before any header is written. +#[inline] +fn can_send_body_stream(stream: &ReadableStream, global_this: &JSGlobalObject) -> bool { + !stream.is_locked(global_this) +} + +/// Refuses a stream body: detaches it for its reader, and returns the `error()` argument. +#[cold] +#[inline(never)] +fn refuse_body_stream(response: Option<&Response>, global_this: &JSGlobalObject) -> JSValue { + bun_core::scoped_log!(ReadableStream, "was locked but it shouldn't be"); + if let Some(response) = response { + response.detach_readable_stream(global_this); + *response.get_body_value() = Body::Value::Used; + } + jsc::SystemError { + code: BunString::static_(<&'static str>::from(jsc::ErrorCode::ERR_STREAM_CANNOT_PIPE)), + message: BunString::static_("Stream already used, please create a new one"), + ..Default::default() + } + .to_error_instance(global_this) +} + +/// Cancel the body stream of a Response the server will not transmit. `stream` is not locked. +#[inline(never)] +fn cancel_unlocked_body( + response: &Response, + stream: Option, + global_this: &JSGlobalObject, +) { + if let Some(stream) = stream { + debug_assert!(!stream.is_locked(global_this)); + let _keep = jsc::EnsureStillAlive(stream.value); + response.detach_readable_stream(global_this); + // Not `cancel()`: it skips a stream with no reader, which an unattached body is. + crate::dispatch::fold(stream.cancel_with_reason(global_this, JSValue::UNDEFINED)); + } + *response.get_body_value() = Body::Value::Used; +} + // ─── sibling-subtree shims ─────────────────────────────────────────────────── // These forward to methods that exist in webcore/ but are currently inside // impl blocks that fail to compile (codegen gc-slot stubs, opaque AbortSignal). @@ -737,17 +777,16 @@ where } /// Cancel the body stream of a Response the server will not transmit, unless a reader holds it. + #[inline(never)] fn cancel_unread_body(response: &Response, global_this: &JSGlobalObject) { - if let Some(stream) = response.get_body_readable_stream() { - let _keep = jsc::EnsureStillAlive(stream.value); - response.detach_readable_stream(global_this); + match response.get_body_readable_stream() { // A locked stream belongs to its reader, which can be the sink of another Response. - if !stream.is_locked(global_this) { - // Not `cancel()`: it skips a stream with no reader, which an unattached body is. - crate::dispatch::fold(stream.cancel_with_reason(global_this, JSValue::UNDEFINED)); + Some(stream) if stream.is_locked(global_this) => { + response.detach_readable_stream(global_this); + *response.get_body_value() = Body::Value::Used; } + stream => cancel_unlocked_body(response, stream, global_this), } - *response.get_body_value() = Body::Value::Used; } /// [`Self::cancel_unread_body`] for a rooted handler result: a `Response` or a settled promise of one. @@ -2663,6 +2702,14 @@ where this.end_without_body(this.should_close_connection()); } Body::Value::Locked(_) => { + let stream = response.get_body_readable_stream(); + if let Some(stream) = &stream + && !can_send_body_stream(stream, global_this) + { + let js_err = refuse_body_stream(Some(&*response), global_this); + this.run_error_handler(js_err); + return; + } this.render_metadata(); if !MUX { // SAFETY: FFI handle @@ -2670,7 +2717,7 @@ where } // HEAD never transmits the body. if let Some(response) = this.response_mut() { - Self::cancel_unread_body(response, global_this); + cancel_unlocked_body(response, stream, global_this); } this.end_without_body(this.should_close_connection()); } @@ -3146,22 +3193,9 @@ where if let Some(stream) = readable_stream { *value = Body::Value::Used; - if stream.is_locked(global_this) { - stream_log!("was locked but it shouldn't be"); - let err = jsc::SystemError { - code: BunString::static_(<&'static str>::from( - jsc::ErrorCode::ERR_STREAM_CANNOT_PIPE, - )), - message: BunString::static_( - "Stream already used, please create a new one", - ), - ..Default::default() - }; - // Teardown must not find the stream: it belongs to its reader. - if let Some(response) = this.response_mut() { - response.detach_readable_stream(global_this); - } - let js_err = err.to_error_instance(global_this); + if !can_send_body_stream(&stream, global_this) { + let js_err = + refuse_body_stream(this.response_mut().map(|r| &*r), global_this); this.run_error_handler(js_err); return; } diff --git a/test/js/bun/http/serve-pending-promise-abort-leak.test.ts b/test/js/bun/http/serve-pending-promise-abort-leak.test.ts index fb37f3b0fc65..3fcad553d1b0 100644 --- a/test/js/bun/http/serve-pending-promise-abort-leak.test.ts +++ b/test/js/bun/http/serve-pending-promise-abort-leak.test.ts @@ -1100,10 +1100,15 @@ test.each(Object.keys(producing))( async kind => { await using held = await producing[kind](); const responses = { "/sent": new Response(held.stream), "/head": new Response(held.stream) }; + const errors: unknown[] = []; using server = Bun.serve({ port: 0, idleTimeout: 0, fetch: req => responses[new URL(req.url).pathname as keyof typeof responses], + error(err: any) { + errors.push(err.code); + return new Response("handled", { status: 500 }); + }, }); const inFlight = await fetch(new URL("/sent", server.url)); @@ -1111,8 +1116,10 @@ test.each(Object.keys(producing))( const decoder = new TextDecoder(); let read = decoder.decode((await body.read()).value, { stream: true }); + // The server cannot send the stream for this Response, so HEAD reports that as GET does. const head = await fetch(new URL("/head", server.url), { method: "HEAD" }); expect(await head.text()).toBe(""); + expect({ status: head.status, errors }).toEqual({ status: 500, errors: ["ERR_STREAM_CANNOT_PIPE"] }); expect(held.cancels()).toBe(0); held.finish(); diff --git a/test/js/bun/http/serve-reused-response.test.ts b/test/js/bun/http/serve-reused-response.test.ts index f870cb3eb003..0a73ff3e3c4a 100644 --- a/test/js/bun/http/serve-reused-response.test.ts +++ b/test/js/bun/http/serve-reused-response.test.ts @@ -1,6 +1,8 @@ import { serve } from "bun"; import { describe, expect, it, jest } from "bun:test"; import { bunEnv, bunExe } from "harness"; +import { once } from "node:events"; +import http2 from "node:http2"; // A fetch handler that returns a Response whose body has already been used // (most often the same Response object returned for every request) must invoke @@ -12,6 +14,11 @@ describe("returning a Response with an already-used body", () => { message: "Response body already used. A Response body can only be sent once; create a new Response for each request.", }; + const streamUsedError = { + code: "ERR_STREAM_CANNOT_PIPE", + name: "Error", + message: "Stream already used, please create a new one", + }; const reusedBodies = { string: () => new Response("cached-route-body"), @@ -114,15 +121,212 @@ describe("returning a Response with an already-used body", () => { expect(await third.text()).toBe("handled"); expect(third.status).toBe(500); - const streamUsedError = { - code: "ERR_STREAM_CANNOT_PIPE", - name: "Error", - message: "Stream already used, please create a new one", - }; expect(errors).toEqual([streamUsedError, streamUsedError]); }, ); + // The server cannot send a stream that another reader holds. GET and HEAD report that through + // error() before any header is written (RFC 9110 §9.3.2), and both leave the stream to its reader. + describe("a stream body that a reader already holds", () => { + const encoder = new TextEncoder(); + const decoder = new TextDecoder(); + async function readToEnd(reader: ReadableStreamDefaultReader) { + let text = ""; + for (let chunk = await reader.read(); !chunk.done; chunk = await reader.read()) { + text += decoder.decode(chunk.value, { stream: true }); + } + return text; + } + + const sources = { + start: (cancel: () => void) => + new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("body")); + controller.close(); + }, + cancel, + }), + pull: (cancel: () => void) => + new ReadableStream({ + pull(controller) { + controller.enqueue(encoder.encode("body")); + controller.close(); + }, + cancel, + }), + direct: (cancel: () => void) => + new ReadableStream({ + type: "direct", + pull(controller) { + controller.write("body"); + controller.close(); + }, + cancel, + }), + }; + + // `rest()` is what the holder of the lock reads from the stream once the request is over. + type Held = { response: Response; rest(): Promise }; + const viaBodyReader = (response: Response): Held => { + const reader = response.body!.getReader(); + return { response, rest: () => readToEnd(reader) }; + }; + const shapes: Record void) => Held | Promise> = { + "body.getReader() on a start() source": cancel => viaBodyReader(new Response(sources.start(cancel))), + "body.getReader() on a pull() source": cancel => viaBodyReader(new Response(sources.pull(cancel))), + 'body.getReader() on a type: "direct" source': cancel => viaBodyReader(new Response(sources.direct(cancel))), + "body.getReader() and a Content-Length header": cancel => + viaBodyReader(new Response(sources.start(cancel), { headers: { "Content-Length": "4" } })), + "stream.getReader() without touching .body": cancel => { + const stream = sources.start(cancel); + const response = new Response(stream); + const reader = stream.getReader(); + return { response, rest: () => readToEnd(reader) }; + }, + "body.tee()": cancel => { + const response = new Response(sources.start(cancel)); + const [branch] = response.body!.tee(); + return { response, rest: () => readToEnd(branch.getReader()) }; + }, + "await response.text()": async cancel => { + const response = new Response(sources.start(cancel)); + const text = await response.text(); + return { response, rest: async () => text }; + }, + "another Response around the same stream was read": async cancel => { + const stream = sources.start(cancel); + const [first, response] = [new Response(stream), new Response(stream)]; + const text = await first.text(); + return { response, rest: async () => text }; + }, + }; + + const deliveries: Record () => Response | Promise> = { + "a Response": response => () => response, + "a fulfilled promise": response => () => Promise.resolve(response), + "a pending promise": response => async () => { + await new Promise(resolve => setImmediate(resolve)); + return response; + }, + }; + + async function request(url: URL | string, method: string) { + const response = await fetch(url, { method }); + return { + status: response.status, + contentLength: response.headers.get("Content-Length"), + transferEncoding: response.headers.get("Transfer-Encoding"), + body: await response.text(), + }; + } + const handled = { status: 500, contentLength: "7", transferEncoding: null, body: "handled" }; + + // Each case pins what GET does, and asserts that HEAD does the same with no body. + describe.each(Object.keys(shapes))("%s", shape => { + it.each(Object.keys(deliveries))("returned as %s: HEAD reports what GET reports", async delivery => { + async function outcome(method: string) { + let cancels = 0; + const held = await shapes[shape](() => void cancels++); + const errors: unknown[] = []; + await using server = serve({ + port: 0, + fetch: deliveries[delivery](held.response), + error(err: any) { + errors.push({ code: err.code, name: err.constructor.name, message: err.message }); + return new Response("handled", { status: 500 }); + }, + }); + const response = await request(server.url, method); + // `rest` and `cancels`: the server neither ended the stream nor cancelled it under its reader. + return { ...response, errors, rest: await held.rest(), cancels }; + } + + const get = await outcome("GET"); + expect(get).toEqual({ ...handled, errors: [streamUsedError], rest: "body", cancels: 0 }); + expect(await outcome("HEAD")).toEqual({ ...get, body: "" }); + }); + }); + + const routeShapes: Record Response) => Bun.Serve.Routes> = { + "an any-method route": handler => ({ "/": handler }), + "a HEAD route": handler => ({ "/": { GET: handler, HEAD: handler } }), + "a HEAD derived from a GET route": handler => ({ "/": { GET: handler } }), + }; + it.each(Object.keys(routeShapes))("%s: HEAD reports what GET reports", async routeShape => { + const errors: unknown[] = []; + await using server = serve({ + port: 0, + routes: routeShapes[routeShape](() => viaBodyReader(new Response(sources.start(() => {}))).response), + error(err: any) { + errors.push({ code: err.code, name: err.constructor.name, message: err.message }); + return new Response("handled", { status: 500 }); + }, + }); + + const get = await request(server.url, "GET"); + expect(get).toEqual(handled); + expect(await request(server.url, "HEAD")).toEqual({ ...get, body: "" }); + expect(errors).toEqual([streamUsedError, streamUsedError]); + }); + + it("over HTTP/2, HEAD reports what GET reports", async () => { + const errors: unknown[] = []; + await using server = serve({ + port: 0, + http2: true, + fetch: () => viaBodyReader(new Response(sources.start(() => {}))).response, + error(err: any) { + errors.push({ code: err.code, name: err.constructor.name, message: err.message }); + return new Response("handled", { status: 500 }); + }, + }); + + const session = http2.connect(`http://127.0.0.1:${server.port}`); + try { + async function h2(method: string) { + const stream = session.request({ ":path": "/", ":method": method }); + const [headers] = await once(stream, "response"); + let body = ""; + for await (const chunk of stream.setEncoding("utf8")) body += chunk; + return { status: headers[":status"], contentLength: headers["content-length"], body }; + } + const get = await h2("GET"); + expect(get).toEqual({ status: 500, contentLength: "7", body: "handled" }); + expect(await h2("HEAD")).toEqual({ ...get, body: "" }); + expect(errors).toEqual([streamUsedError, streamUsedError]); + } finally { + session.close(); + } + }); + + // The refusal uses the body up, as a sent body is: the next request gets the already-used error. + it.each([ + ["GET", "GET"], + ["HEAD", "HEAD"], + ["HEAD", "GET"], + ["GET", "HEAD"], + ])("the same Response returned for %s and then %s", async (first, second) => { + const { response } = viaBodyReader(new Response(sources.start(() => {}))); + const errors: unknown[] = []; + await using server = serve({ + port: 0, + fetch: () => response, + error(err: any) { + errors.push({ code: err.code, name: err.constructor.name, message: err.message }); + return new Response("handled", { status: 500 }); + }, + }); + + const bodyFor = (method: string) => (method === "HEAD" ? "" : "handled"); + expect([await request(server.url, first), await request(server.url, second)]).toEqual([ + { ...handled, body: bodyFor(first) }, + { ...handled, body: bodyFor(second) }, + ]); + expect(errors).toEqual([streamUsedError, alreadyUsedError]); + }); + }); + // The second input leaves a Content-Length on the used Response: it must not frame a HEAD 200. const leftoverContentLength = { headers: { "Content-Length": "24" } }; it.each([ @@ -244,8 +448,7 @@ describe("returning a Response with an already-used body", () => { const whole = "chunk1;chunk2;chunk3;chunk4;chunk5;chunk6;"; it.each([ ["GET", { status: 500, read: whole, reports: 1, exitCode: 1 }], - // HEAD sends no body, so the server has nothing to refuse. - ["HEAD", { status: 200, read: whole, reports: 0, exitCode: 0 }], + ["HEAD", { status: 500, read: whole, reports: 1, exitCode: 1 }], ] as const)("%s: the handler's reader still reads the whole body", async (method, expected) => { const result = await run(` const release = Promise.withResolvers(); @@ -294,6 +497,58 @@ describe("returning a Response with an already-used body", () => { expect(result).toEqual(expected); }); + // error() has run already, so the refusal of its own Response gets the default 500. + it.each(["GET", "HEAD"])( + "%s: an error() Response whose stream a reader holds gets the default 500", + async method => { + const result = await run(` + const held = () => { + const response = new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("body")); + controller.close(); + }, + }), + { status: 418 }, + ); + response.body.getReader(); + return response; + }; + const errorHandlers = [ + () => held(), + () => Promise.resolve(held()), + async () => { + await new Promise(resolve => setImmediate(resolve)); + return held(); + }, + ]; + const responses = []; + for (const error of errorHandlers) { + const server = Bun.serve({ + port: 0, + development: false, + fetch() { + throw new Error("boom"); + }, + error, + }); + const response = await fetch(server.url, { method: ${JSON.stringify(method)} }); + responses.push([ + response.status, + response.headers.get("Content-Length"), + response.headers.get("Transfer-Encoding"), + await response.text(), + ]); + await server.stop(true); + } + console.log(JSON.stringify({ responses })); + `); + const response = [500, "21", null, method === "HEAD" ? "" : "Something went wrong!"]; + expect(result).toEqual({ responses: [response, response, response], reports: 3, exitCode: 1 }); + }, + ); + it("GET: another client's download of the same body completes", async () => { const result = await run(` const release = Promise.withResolvers();