From ceb32be1ae2946d2e2fdf6bdf8887503d401c52a Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 17 Sep 2026 22:32:35 +0000 Subject: [PATCH 01/10] http: a held response body outlives its socket instead of being decoded whole at EOF The decode budget from #43123 stopped at the end of the transport. `finalize_body_on_eof` decodes everything that is held, and a lot is held whenever the end arrives while the reader is behind: - Through a CONNECT tunnel, whose socket is never paused, all of the body and then the origin's close reach the client at once. A 256 MB-logical zstd or gzip body (270 KB on the wire) grew the client by 507 to 551 MB whether it ended by Content-Length, chunked or close; 1 GiB-logical by 1278 MB. - On a direct TLS connection the socket is paused, which defers a FIN but not a close_notify that came in with the body: 3 of 30 runs grew by 288 MB for a 64 MB-logical body. Now the end of the transport marks the body complete and decodes one budgeted pass like any other read. If input is still held, the client registers in `socketless_bodies` by `async_http_id` and stays alive without a socket. The consumer's pulls (`drain_queued_receive_resumes`) decode the rest a pass at a time through `send_progress_update_without_socket`, the socket-free update h2 and h3 already use, and an abort (`drain_queued_shutdowns`) fails it. The entry is removed by `unregister_abort_tracker`, which every terminal path runs before the client is freed. A tunnel whose inner TLS stream closes first keeps draining through its outer socket until the proxy closes that too. Also: the CONNECT test from #43123 never reached its proxy where the ambient NO_PROXY lists 127.0.0.1. The client env now drops the proxy variables and the test asserts that the proxy saw the CONNECT. --- src/http/HTTPThread.rs | 6 + src/http/InternalState.rs | 4 + src/http/ProxyTunnel.rs | 13 ++- src/http/lib.rs | 112 ++++++++++++++++--- test/js/web/fetch/fetch-backpressure.test.ts | 57 ++++++++-- 5 files changed, 160 insertions(+), 32 deletions(-) diff --git a/src/http/HTTPThread.rs b/src/http/HTTPThread.rs index d3b0f526cb6b..85bc795aea4f 100644 --- a/src/http/HTTPThread.rs +++ b/src/http/HTTPThread.rs @@ -590,6 +590,10 @@ impl HttpThread { if self.abort_pending_h2_waiter(http.async_http_id) { continue; } + if let Some(client) = crate::socketless_body(http.async_http_id) { + client.abort_socketless_body(); + continue; + } // Or it's on an HTTP/3 session, which has no TCP socket to // register in the tracker. if h3::ClientContext::abort_by_http_id(http.async_http_id) { @@ -788,6 +792,8 @@ impl HttpThread { } } } + } else if let Some(client) = crate::socketless_body(id) { + client.drain_socketless_body(); } else { h3::ClientContext::resume_receive_by_http_id(id); } diff --git a/src/http/InternalState.rs b/src/http/InternalState.rs index 1a083f5e38d2..0079b55342b6 100644 --- a/src/http/InternalState.rs +++ b/src/http/InternalState.rs @@ -84,6 +84,9 @@ pub struct InternalStateFlags { pub(crate) body_compressed: bool, /// Held input or buffered decoder output remains for `HTTPClient::drain_response_body`. pub(crate) decompress_output_pending: bool, + /// The socket is gone and the consumer still pulls held input: the client is in + /// `socketless_bodies`, which is how a resume or an abort finds it. + pub(crate) body_outlived_socket: bool, } impl InternalStateFlags { @@ -100,6 +103,7 @@ impl InternalStateFlags { receive_paused: false, body_compressed: false, decompress_output_pending: false, + body_outlived_socket: false, } } } diff --git a/src/http/ProxyTunnel.rs b/src/http/ProxyTunnel.rs index ac7bfaa877ca..3c79edef6995 100644 --- a/src/http/ProxyTunnel.rs +++ b/src/http/ProxyTunnel.rs @@ -506,10 +506,17 @@ fn on_close(ctx: *mut HTTPClient) { && !this.state.flags.is_redirect_pending; let mut fail_err: Option = None; if in_progress && this.state.is_body_complete_on_close() { - match this.state.finalize_body_on_eof() { + match this.finish_body_on_close() { Ok(()) => { - // `this` dead (NLL); reborrow via `client_from_ctx` inside. - progress_update_for_proxy_socket(ctx, proxy_nn); + // A held body with nothing decoded has no update yet: the consumer's pull + // drains it through the outer socket, or through `socketless_bodies` once + // the proxy closes that too. + let held = + this.state.has_pending_compressed() && this.state.decoded_body.list.is_empty(); + if !held { + // `this` dead (NLL); reborrow via `client_from_ctx` inside. + progress_update_for_proxy_socket(ctx, proxy_nn); + } crate::http_thread().schedule_proxy_deref(keepalive); return; } diff --git a/src/http/lib.rs b/src/http/lib.rs index b436832e317a..14f5171a2e9b 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -923,6 +923,13 @@ pub(crate) static SOCKET_ASYNC_HTTP_ABORT_TRACKER: bun_core::RacyCell< Option>, > = bun_core::RacyCell::new(None); +/// h1 clients whose socket closed while their consumer still pulls a held body, by +/// `async_http_id`. HTTP-thread-only, like the abort tracker. An entry is removed by +/// `unregister_abort_tracker`, which every terminal path runs before the client is freed. +pub(crate) static SOCKETLESS_BODIES: bun_core::RacyCell< + Option>>>, +> = bun_core::RacyCell::new(None); + // ═══════════════════════════════════════════════════════════════════════ // Prelude: imports, constants, helper fns, and bridge impls the // `impl HTTPClient` state machine needs. Kept separate from the head/tail @@ -1158,6 +1165,21 @@ fn abort_tracker() -> &'static mut ArrayHashMap { unsafe { (*SOCKET_ASYNC_HTTP_ABORT_TRACKER.get()).get_or_insert_with(ArrayHashMap::new) } } +/// Same contract as [`abort_tracker`]. +#[inline] +fn socketless_bodies() -> &'static mut ArrayHashMap>> { + // SAFETY: HTTP-thread only; every call site is a per-statement reborrow. + unsafe { (*SOCKETLESS_BODIES.get()).get_or_insert_with(ArrayHashMap::new) } +} + +/// The client behind `async_http_id` if its body outlived its socket. +pub(crate) fn socketless_body<'b>(async_http_id: u32) -> Option<&'b mut HTTPClient<'static>> { + socketless_bodies() + .get(&async_http_id) + .copied() + .map(HTTPClient::from_erased_backref) +} + /// Remove every abort-tracker entry whose stored socket is `socket`. /// /// Backstop for the per-client `unregister_abort_tracker()` calls: when @@ -1796,6 +1818,9 @@ impl<'a> HTTPClient<'a> { // SAFETY: HTTP-thread only; per-statement reborrow. let _ = abort_tracker().swap_remove(&self.async_http_id); } + if core::mem::take(&mut self.state.flags.body_outlived_socket) { + let _ = socketless_bodies().swap_remove(&self.async_http_id); + } } /// Runs once per request: for a new connection via [`Self::on_connect`], @@ -2144,10 +2169,14 @@ impl<'a> HTTPClient<'a> { return; } if in_progress && self.state.is_body_complete_on_close() { - if let Err(err) = self.state.finalize_body_on_eof() { + if let Err(err) = self.finish_body_on_close() { self.fail(err); return; } + if self.state.has_pending_compressed() { + self.outlive_socket(); + return; + } let ctx = self.get_ssl_ctx::(); self.progress_update::(ctx, socket); return; @@ -4148,45 +4177,86 @@ impl<'a> HTTPClient<'a> { } pub(crate) fn drain_response_body(&mut self, socket: HttpSocket) { - if self.pump_held_body::(socket) { + if self.pump_held_body_or_close::(socket) { let ctx = self.get_ssl_ctx::(); self.send_progress_update_without_stage_check::(ctx, socket); } } + fn pump_held_body_or_close(&mut self, socket: HttpSocket) -> bool { + match self.pump_held_body() { + Ok(has_update) => has_update, + Err(err) => { + self.close_and_fail::(err, socket); + false + } + } + } + /// Decodes the next piece of a held body. Returns whether there is an update to send. - fn pump_held_body(&mut self, socket: HttpSocket) -> bool { + fn pump_held_body(&mut self) -> crate::Result { // Find out if we should not send any update. match self.state.stage { - Stage::Done | Stage::Fail => return false, + Stage::Done | Stage::Fail => return Ok(false), _ => {} } if self.state.fail.is_some() { // If there's any error at all, do not drain. - return false; + return Ok(false); } // If there's a pending redirect, then don't bother to send a response body // as that wouldn't make sense and I want to defensively avoid edgecases // from that. if self.state.flags.is_redirect_pending { - return false; + return Ok(false); } // A consumer that paused again gets another resume when it unpauses. let pumped = self.state.has_pending_compressed() && !self.signals.is_receive_paused(); if pumped { let is_final = self.state.is_done(); - if let Err(err) = self.process_received_body(is_final) { - self.close_and_fail::(err, socket); - return false; - } + self.process_received_body(is_final)?; } // A pump that ends the body has to say so even with no bytes (a stream trailer alone). let ended = pumped && self.state.is_done() && !self.state.has_pending_compressed(); - !self.state.decoded_body.list.is_empty() || ended + Ok(!self.state.decoded_body.list.is_empty() || ended) + } + + /// The transport ended, and with it the body. What a consumer's budget holds stays held. + pub(crate) fn finish_body_on_close(&mut self) -> crate::Result<()> { + self.state.flags.received_last_chunk = true; + self.process_received_body(true).map(drop) + } + + /// Keeps this client reachable by id once its socket is gone, then delivers what it can. + fn outlive_socket(&mut self) { + self.state.flags.body_outlived_socket = true; + let _ = socketless_bodies().put(self.async_http_id, self.as_erased_ptr()); + if !self.state.decoded_body.list.is_empty() && self.send_progress_update_without_socket() { + self.drain_socketless_body(); + } + } + + /// A consumer's pull (`drain_queued_receive_resumes`) for a body that outlived its socket. + pub(crate) fn drain_socketless_body(&mut self) { + loop { + match self.pump_held_body() { + Ok(true) => {} + Ok(false) => return, + Err(err) => return self.fail(err), + } + if !self.send_progress_update_without_socket() { + return; + } + } + } + + /// An abort (`drain_queued_shutdowns`) for a body that outlived its socket. + pub(crate) fn abort_socketless_body(&mut self) { + self.fail(crate::Error::Aborted); } fn send_progress_update_without_stage_check( @@ -4195,12 +4265,13 @@ impl<'a> HTTPClient<'a> { socket: HttpSocket, ) { if self.flags.protocol != Protocol::Http1_1 { - return self.send_progress_update_multiplexed(); + self.send_progress_update_without_socket(); + return; } // A loop, not a call back into `drain_response_body`: a consumer that never pauses // (`BufferAll`, or an S3 error body that is collected whole) takes one pass per turn. while self.send_one_progress_update::(ctx, socket) - && self.pump_held_body::(socket) + && self.pump_held_body_or_close::(socket) {} } @@ -4336,9 +4407,13 @@ impl<'a> HTTPClient<'a> { /// `send_progress_update_without_stage_check` minus the per-request TCP socket /// release/close. Used by HTTP/2 and HTTP/3, whose session owns the - /// transport, so there is no `ctx`/`socket` to hand back to the pool here. - fn send_progress_update_multiplexed(&mut self) { - debug_assert!(self.flags.protocol != Protocol::Http1_1); + /// transport, and by an h1 body that outlived its socket, so there is no + /// `ctx`/`socket` to hand back to the pool here. + /// Returns whether a held body is left that its consumer will not ask for. + fn send_progress_update_without_socket(&mut self) -> bool { + debug_assert!( + self.flags.protocol != Protocol::Http1_1 || self.state.flags.body_outlived_socket + ); let callback = self.result_callback; let mut result = self.to_result(); @@ -4359,7 +4434,7 @@ impl<'a> HTTPClient<'a> { if is_done { result.body_owned = decoded_body.list; callback.run(parent, result); - return; + return false; } result.body = decoded_body.list.as_slice(); callback.run(parent, result); @@ -4367,6 +4442,7 @@ impl<'a> HTTPClient<'a> { decoded_body.list.clear(); self.state.decoded_body = decoded_body; } + self.state.has_pending_compressed() && !self.signals.is_receive_paused() } /// `do_redirect` minus the per-request socket release/close. The session @@ -4418,7 +4494,7 @@ impl<'a> HTTPClient<'a> { } return; } - self.send_progress_update_multiplexed(); + self.send_progress_update_without_socket(); } pub(crate) fn do_redirect_h3(&mut self) { diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index f697b83759b3..346d0fff29a5 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -559,14 +559,26 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }; } - // Close-delimited, and the origin never closes: for the client this body does not end. - async function serveBomb(enc: Enc, secure: boolean) { + // How the origin ends the body. "never": close-delimited and the origin never closes, so for + // the client this body does not end. The others send the whole body and close at once, so the + // end of the transport reaches a client that still holds nearly all of the body undecoded. + type Ending = "never" | "content-length" | "chunked" | "close-delimited"; + async function serveBomb(enc: Enc, secure: boolean, ending: Ending = "never") { const bomb = bombFor(enc); + const framing = + ending === "content-length" + ? `Content-Length: ${bomb.length}\r\n` + : ending === "chunked" + ? "Transfer-Encoding: chunked\r\n" + : ""; const handler = (s: import("node:net").Socket) => { s.on("error", () => {}); s.once("data", () => { - s.write(`HTTP/1.1 200 OK\r\nContent-Encoding: ${enc}\r\nConnection: close\r\n\r\n`); + s.write(`HTTP/1.1 200 OK\r\nContent-Encoding: ${enc}\r\n${framing}Connection: close\r\n\r\n`); + if (ending === "chunked") s.write(`${bomb.length.toString(16)}\r\n`); s.write(bomb); + if (ending === "chunked") s.write("\r\n0\r\n\r\n"); + if (ending !== "never") s.end(); }); }; const server = await listening(secure ? createTlsServer(tls, handler) : createTcpServer(handler)); @@ -576,13 +588,15 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not // A CONNECT proxy that pipes both ways. bun does not pause a tunnelled socket, so the origin's // bytes keep arriving while the reader is paused. async function serveConnectProxy() { + let connects = 0; const server = await listening( createTcpServer(client => { let upstream: import("node:net").Socket | undefined; client.on("error", () => upstream?.destroy()); client.on("close", () => upstream?.destroy()); client.once("data", head => { - const [, target] = head.toString("latin1").split(" "); + const [method, target] = head.toString("latin1").split(" "); + if (method === "CONNECT") connects++; const colon = target.lastIndexOf(":"); upstream = connect(Number(target.slice(colon + 1)), target.slice(0, colon), () => { client.write("HTTP/1.1 200 Connection Established\r\n\r\n"); @@ -594,7 +608,7 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }); }), ); - return { ...server, url: `http://127.0.0.1:${server.port}` }; + return { ...server, url: `http://127.0.0.1:${server.port}`, connects: () => connects }; } // Takes one chunk, lets the client's memory settle, takes a few more, and reports the largest @@ -623,10 +637,17 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not process.stdout.write(JSON.stringify({ got, peak, zeros })); `; + // An ambient NO_PROXY that lists 127.0.0.1 makes fetch() ignore its `proxy` option. + const clientEnv = { ...bunEnv }; + for (const key of ["NO_PROXY", "HTTP_PROXY", "HTTPS_PROXY", "ALL_PROXY"]) { + delete clientEnv[key]; + delete clientEnv[key.toLowerCase()]; + } + async function runClient(url: string, opts: object, script: string) { await using proc = Bun.spawn({ cmd: [bunExe(), "-e", `const url=${JSON.stringify(url)};const opts=${JSON.stringify(opts)};${script}`], - env: bunEnv, + env: clientEnv, stdout: "pipe", stderr: "pipe", }); @@ -662,11 +683,25 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }, ); - test("zstd through a CONNECT proxy: a reader that takes a little holds a little", async () => { - await using server = await serveBomb("zstd", true); - await using proxy = await serveConnectProxy(); - const opts = { proxy: proxy.url, tls: { rejectUnauthorized: false } }; - expectBounded(await runClient(server.url, opts, READ_A_LITTLE)); + // bun does not pause a tunnelled socket, so all of the body and then the origin's close reach + // the client at once. The end of the transport must not decode what the reader has not asked + // for: the client outlives its socket until the reader has pulled the rest, or cancels. + test.each(["never", "content-length", "chunked", "close-delimited"] as Ending[])( + "zstd through a CONNECT proxy, body ending %s: a reader that takes a little holds a little", + async ending => { + await using server = await serveBomb("zstd", true, ending); + await using proxy = await serveConnectProxy(); + const opts = { proxy: proxy.url, tls: { rejectUnauthorized: false } }; + expectBounded(await runClient(server.url, opts, READ_A_LITTLE)); + expect(proxy.connects()).toBe(1); + }, + ); + + // Without a tunnel the socket is paused, which defers a FIN but not a TLS close_notify that + // came in with the body. + test.each([false, true])("gzip, origin closes after the body (tls: %p)", async secure => { + await using server = await serveBomb("gzip", secure, "content-length"); + expectBounded(await runClient(server.url, { tls: { rejectUnauthorized: false } }, READ_A_LITTLE)); }); // A live stream: two flushed messages in one packet, then the origin goes quiet with the frame From 0760886f28d2990a03932b76edfb72e6f5720916 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 00:10:43 +0000 Subject: [PATCH 02/10] http: a finished transport hands the undecoded rest of a body to its consumer When every wire byte of a compressed body has arrived and the consumer's decode budget has left part of it undecoded, the request now ends as usual. The final result carries a HeldBody: the undecoded input and the decoder. FetchTasklet owns it and decodes one budget per reader pull in a work-pool job. No client outlives its socket, so the socketless registry is gone, and so are the special cases for a body that is complete but held. --- src/http/AsyncHTTP.rs | 3 + src/http/Decompressor.rs | 6 +- src/http/HTTPThread.rs | 6 - src/http/InternalState.rs | 52 ++++- src/http/ProxyTunnel.rs | 11 +- src/http/Signals.rs | 8 + src/http/lib.rs | 191 ++++++----------- src/runtime/webcore/fetch/FetchTasklet.rs | 160 ++++++++++++++- test/js/web/fetch/fetch-backpressure.test.ts | 203 ++++++++++++++++--- 9 files changed, 454 insertions(+), 186 deletions(-) diff --git a/src/http/AsyncHTTP.rs b/src/http/AsyncHTTP.rs index 31776836f9c5..7c87254e960e 100644 --- a/src/http/AsyncHTTP.rs +++ b/src/http/AsyncHTTP.rs @@ -253,6 +253,8 @@ pub struct Options<'a> { pub verbose: Option, pub disable_keepalive: Option, pub disable_decompression: Option, + /// The consumer takes `HTTPClientResult::held_body`: a compressed body is decoded as it reads. + pub takes_held_body: bool, pub max_redirects: Option, pub reject_unauthorized: Option, pub tls_props: Option, @@ -476,6 +478,7 @@ impl<'a> AsyncHTTP<'a> { if let Some(val) = options.disable_decompression { this.client.flags.disable_decompression = val; } + this.client.flags.takes_held_body = options.takes_held_body; if let Some(val) = options.max_redirects { this.client.remaining_redirect_count = (val.min(126) + 1) as i8; } diff --git a/src/http/Decompressor.rs b/src/http/Decompressor.rs index fac59577b7d5..1804ef3c6376 100644 --- a/src/http/Decompressor.rs +++ b/src/http/Decompressor.rs @@ -1,4 +1,3 @@ -use bun_core::MutableString; use bun_http_types::Encoding::Encoding; // The streaming decoders below own only their C-side state and take @@ -59,14 +58,14 @@ impl Decompressor { } /// Feed one body chunk `buffer` through the decoder, appending the - /// decompressed output to `body_out_str` until it holds `max_output` bytes. Creates the + /// decompressed output to `out` until it holds `max_output` bytes. Creates the /// decoder on first call. Returns the input bytes consumed. Returns `ShortRead` when more /// input is needed and the stream is not yet done. pub(crate) fn decompress_chunk( &mut self, encoding: Encoding, buffer: &[u8], - body_out_str: &mut MutableString, + out: &mut Vec, max_output: usize, is_done: bool, ) -> crate::Result { @@ -76,7 +75,6 @@ impl Decompressor { if matches!(self, Decompressor::None) { self.init(encoding, buffer)?; } - let out = &mut body_out_str.list; match self { Decompressor::Zlib(reader) => Ok(reader.decompress(buffer, out, max_output, is_done)?), Decompressor::Brotli(reader) => { diff --git a/src/http/HTTPThread.rs b/src/http/HTTPThread.rs index 85bc795aea4f..d3b0f526cb6b 100644 --- a/src/http/HTTPThread.rs +++ b/src/http/HTTPThread.rs @@ -590,10 +590,6 @@ impl HttpThread { if self.abort_pending_h2_waiter(http.async_http_id) { continue; } - if let Some(client) = crate::socketless_body(http.async_http_id) { - client.abort_socketless_body(); - continue; - } // Or it's on an HTTP/3 session, which has no TCP socket to // register in the tracker. if h3::ClientContext::abort_by_http_id(http.async_http_id) { @@ -792,8 +788,6 @@ impl HttpThread { } } } - } else if let Some(client) = crate::socketless_body(id) { - client.drain_socketless_body(); } else { h3::ClientContext::resume_receive_by_http_id(id); } diff --git a/src/http/InternalState.rs b/src/http/InternalState.rs index 0079b55342b6..882dec23602b 100644 --- a/src/http/InternalState.rs +++ b/src/http/InternalState.rs @@ -84,9 +84,6 @@ pub struct InternalStateFlags { pub(crate) body_compressed: bool, /// Held input or buffered decoder output remains for `HTTPClient::drain_response_body`. pub(crate) decompress_output_pending: bool, - /// The socket is gone and the consumer still pulls held input: the client is in - /// `socketless_bodies`, which is how a resume or an abort finds it. - pub(crate) body_outlived_socket: bool, } impl InternalStateFlags { @@ -103,7 +100,6 @@ impl InternalStateFlags { receive_paused: false, body_compressed: false, decompress_output_pending: false, - body_outlived_socket: false, } } } @@ -227,16 +223,30 @@ impl<'a> InternalState<'a> { /// than failing it: chunked decoder already in the trailers state, or a /// close-delimited response (no Content-Length, no Transfer-Encoding). pub(crate) fn is_body_complete_on_close(&self) -> bool { - // Every byte arrived; only the decode is outstanding. - if self.flags.decompress_output_pending && self.is_done() { - return true; - } if self.is_chunked_encoding() { return bun_picohttp::phr_decode_chunked_is_in_trailers(&self.chunked_decoder) != 0; } self.content_length.is_none() && self.response_stage == HTTPStage::Body } + /// Drops input that nothing will decode. + pub(crate) fn discard_held_input(&mut self) { + self.flags.decompress_output_pending = false; + self.compressed_body.list.clear(); + self.compressed_body_consumed = 0; + } + + /// The rest of a body that is complete on the wire, for `HTTPClientResult::held_body`. + pub(crate) fn take_held_body(&mut self) -> HeldBody { + self.flags.decompress_output_pending = false; + HeldBody { + encoding: self.encoding, + decompressor: core::mem::take(&mut self.decompressor), + input: core::mem::take(&mut self.compressed_body.list), + consumed: core::mem::take(&mut self.compressed_body_consumed), + } + } + /// Mark the body complete and drive `process_body_buffer` one last time /// with `is_final_chunk = true` so a compressed stream that never reached /// stream-end is rejected. Call from every site that flips @@ -398,7 +408,7 @@ impl<'a> InternalState<'a> { match self.decompressor.decompress_chunk( self.encoding, buffer, - &mut self.decoded_body, + &mut self.decoded_body.list, max_output, is_done, ) { @@ -481,6 +491,30 @@ impl<'a> InternalState<'a> { } } +/// The undecoded rest of a body whose transport finished, and its decoder. The consumer owns it. +pub struct HeldBody { + encoding: Encoding, + decompressor: Decompressor, + input: Vec, + consumed: usize, +} + +// SAFETY: owns its bytes and its decoder's C state, which no thread has a claim on. +unsafe impl Send for HeldBody {} + +impl HeldBody { + /// Appends decoded bytes to `out` until it holds `max_output`. `Ok(true)`: the body has ended. + pub fn decode(&mut self, out: &mut Vec, max_output: usize) -> Result { + let input = &self.input[self.consumed..]; + self.consumed += + self.decompressor + .decompress_chunk(self.encoding, input, out, max_output, true)?; + let ended = out.len() < max_output + || (self.consumed == self.input.len() && !self.decompressor.is_mid_stream()); + Ok(ended) + } +} + #[derive(Clone, Copy, PartialEq, Eq)] pub enum HTTPStage { Pending, diff --git a/src/http/ProxyTunnel.rs b/src/http/ProxyTunnel.rs index 3c79edef6995..f1f3d2239827 100644 --- a/src/http/ProxyTunnel.rs +++ b/src/http/ProxyTunnel.rs @@ -508,15 +508,8 @@ fn on_close(ctx: *mut HTTPClient) { if in_progress && this.state.is_body_complete_on_close() { match this.finish_body_on_close() { Ok(()) => { - // A held body with nothing decoded has no update yet: the consumer's pull - // drains it through the outer socket, or through `socketless_bodies` once - // the proxy closes that too. - let held = - this.state.has_pending_compressed() && this.state.decoded_body.list.is_empty(); - if !held { - // `this` dead (NLL); reborrow via `client_from_ctx` inside. - progress_update_for_proxy_socket(ctx, proxy_nn); - } + // `this` dead (NLL); reborrow via `client_from_ctx` inside. + progress_update_for_proxy_socket(ctx, proxy_nn); crate::http_thread().schedule_proxy_deref(keepalive); return; } diff --git a/src/http/Signals.rs b/src/http/Signals.rs index 868a028bb1f0..75d065f1a21e 100644 --- a/src/http/Signals.rs +++ b/src/http/Signals.rs @@ -85,6 +85,14 @@ impl Signals { .is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Paused as u8) } + /// Nothing will read the body, and its consumer is shutting the transport down. + #[inline] + pub(crate) fn is_body_abandoned(self) -> bool { + self.body_receive_mode + .map(bun_ptr::BackRef::from) + .is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Abandoned as u8) + } + /// `Flowing` or `Paused`: a consumer takes the body piece by piece. #[inline] pub(crate) fn is_demand_driven(self) -> bool { diff --git a/src/http/lib.rs b/src/http/lib.rs index 14f5171a2e9b..73cce6febd73 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -59,7 +59,7 @@ pub use http_context::{HTTPContext, HTTPSocket, PeerVerification}; pub use http_request_body::HTTPRequestBody; pub use http_thread::HttpThread as HTTPThread; pub use http_thread::shutdown_for_exit; -pub use internal_state::InternalState; +pub use internal_state::{HeldBody, InternalState}; pub use proxy_tunnel::ProxyTunnel; pub use send_file::SendFile; pub use signals::Signals; @@ -222,6 +222,8 @@ pub struct Flags { pub forced_protocol: Option, pub(crate) h3_retried: bool, pub is_node_http_client: bool, + /// `Options::takes_held_body`. + pub(crate) takes_held_body: bool, } impl Default for Flags { @@ -245,6 +247,7 @@ impl Default for Flags { forced_protocol: None, h3_retried: false, is_node_http_client: false, + takes_held_body: false, } } } @@ -487,6 +490,8 @@ pub struct HTTPClientResult<'a> { /// Boxed: it is large and rare, and every result is moved and dropped /// several times per request. pub proxy_connect_response: Option>, + /// Final result only: what the consumer's budget left undecoded (`Options::takes_held_body`). + pub held_body: Option>, } /// Keep-alive pool partition of the fetch session a request belongs to. @@ -576,6 +581,7 @@ impl<'a> HTTPClientResult<'a> { certificate_info: self.certificate_info, connect_errno: self.connect_errno, proxy_connect_response: self.proxy_connect_response, + held_body: self.held_body, } } } @@ -923,13 +929,6 @@ pub(crate) static SOCKET_ASYNC_HTTP_ABORT_TRACKER: bun_core::RacyCell< Option>, > = bun_core::RacyCell::new(None); -/// h1 clients whose socket closed while their consumer still pulls a held body, by -/// `async_http_id`. HTTP-thread-only, like the abort tracker. An entry is removed by -/// `unregister_abort_tracker`, which every terminal path runs before the client is freed. -pub(crate) static SOCKETLESS_BODIES: bun_core::RacyCell< - Option>>>, -> = bun_core::RacyCell::new(None); - // ═══════════════════════════════════════════════════════════════════════ // Prelude: imports, constants, helper fns, and bridge impls the // `impl HTTPClient` state machine needs. Kept separate from the head/tail @@ -1165,21 +1164,6 @@ fn abort_tracker() -> &'static mut ArrayHashMap { unsafe { (*SOCKET_ASYNC_HTTP_ABORT_TRACKER.get()).get_or_insert_with(ArrayHashMap::new) } } -/// Same contract as [`abort_tracker`]. -#[inline] -fn socketless_bodies() -> &'static mut ArrayHashMap>> { - // SAFETY: HTTP-thread only; every call site is a per-statement reborrow. - unsafe { (*SOCKETLESS_BODIES.get()).get_or_insert_with(ArrayHashMap::new) } -} - -/// The client behind `async_http_id` if its body outlived its socket. -pub(crate) fn socketless_body<'b>(async_http_id: u32) -> Option<&'b mut HTTPClient<'static>> { - socketless_bodies() - .get(&async_http_id) - .copied() - .map(HTTPClient::from_erased_backref) -} - /// Remove every abort-tracker entry whose stored socket is `socket`. /// /// Backstop for the per-client `unregister_abort_tracker()` calls: when @@ -1818,9 +1802,6 @@ impl<'a> HTTPClient<'a> { // SAFETY: HTTP-thread only; per-statement reborrow. let _ = abort_tracker().swap_remove(&self.async_http_id); } - if core::mem::take(&mut self.state.flags.body_outlived_socket) { - let _ = socketless_bodies().swap_remove(&self.async_http_id); - } } /// Runs once per request: for a new connection via [`Self::on_connect`], @@ -2173,10 +2154,6 @@ impl<'a> HTTPClient<'a> { self.fail(err); return; } - if self.state.has_pending_compressed() { - self.outlive_socket(); - return; - } let ctx = self.get_ssl_ctx::(); self.progress_update::(ctx, socket); return; @@ -2216,11 +2193,6 @@ impl<'a> HTTPClient<'a> { if self.flags.disable_timeout { return; } - // A fully received body that waits on its consumer expects nothing from the socket. - if self.state.has_pending_compressed() && self.state.is_done() { - socket.set_timeout(0); - return; - } bun_core::scoped_log!(fetch, "Timeout {}\n", BStr::new(self.url.href)); // Terminate (mark dead + close) BEFORE failing, matching // `close_and_fail`: `fail()` dispatches the final result, which frees @@ -4116,10 +4088,18 @@ impl<'a> HTTPClient<'a> { socket.set_timeout(self.effective_idle_timeout_seconds()); } - /// Output budget of one decode pass. h1 only: h2/h3 detach before held input could drain. + /// A compressed body is decoded one budget at a time, as its consumer reads. h1 only. + #[inline] + fn decodes_on_demand(&self) -> bool { + self.flags.takes_held_body + && self.flags.protocol == Protocol::Http1_1 + && self.signals.is_demand_driven() + } + + /// Output budget of one decode pass. #[inline] fn decompress_output_cap(&self) -> usize { - if self.flags.protocol == Protocol::Http1_1 && self.signals.is_demand_driven() { + if self.decodes_on_demand() { signals::BODY_HIGH_WATER_MARK } else { usize::MAX @@ -4129,13 +4109,17 @@ impl<'a> HTTPClient<'a> { /// Decodes what has arrived under the consumer's budget. Returns whether to report bytes. fn process_received_body(&mut self, is_final_chunk: bool) -> crate::Result { let max_output = self.decompress_output_cap(); - // Nothing is decoded for a paused consumer (a tunnelled socket keeps reading anyway). - if max_output != usize::MAX - && self.state.encoding.is_compressed() - && self.signals.is_receive_paused() - { - self.state.flags.decompress_output_pending = true; - return Ok(false); + if self.state.encoding.is_compressed() { + // Nothing will read it, and its transport is being shut down. + if self.signals.is_body_abandoned() { + self.state.discard_held_input(); + return Ok(false); + } + // Nothing is decoded for a paused consumer (a tunnelled socket keeps reading anyway). + if max_output != usize::MAX && self.signals.is_receive_paused() { + self.state.flags.decompress_output_pending = true; + return Ok(false); + } } // `process_body_buffer` takes `&mut self.state`, so the bytes move out first. let buffer = core::mem::take(&mut self.state.get_body_buffer().list); @@ -4176,87 +4160,49 @@ impl<'a> HTTPClient<'a> { self.set_timeout(&socket); } + /// The transport ended, and with it a close-delimited body. + pub(crate) fn finish_body_on_close(&mut self) -> crate::Result<()> { + self.state.flags.received_last_chunk = true; + self.process_received_body(true).map(drop) + } + pub(crate) fn drain_response_body(&mut self, socket: HttpSocket) { - if self.pump_held_body_or_close::(socket) { + if self.pump_held_body::(socket) { let ctx = self.get_ssl_ctx::(); self.send_progress_update_without_stage_check::(ctx, socket); } } - fn pump_held_body_or_close(&mut self, socket: HttpSocket) -> bool { - match self.pump_held_body() { - Ok(has_update) => has_update, - Err(err) => { - self.close_and_fail::(err, socket); - false - } - } - } - /// Decodes the next piece of a held body. Returns whether there is an update to send. - fn pump_held_body(&mut self) -> crate::Result { + fn pump_held_body(&mut self, socket: HttpSocket) -> bool { // Find out if we should not send any update. match self.state.stage { - Stage::Done | Stage::Fail => return Ok(false), + Stage::Done | Stage::Fail => return false, _ => {} } if self.state.fail.is_some() { // If there's any error at all, do not drain. - return Ok(false); + return false; } // If there's a pending redirect, then don't bother to send a response body // as that wouldn't make sense and I want to defensively avoid edgecases // from that. if self.state.flags.is_redirect_pending { - return Ok(false); + return false; } // A consumer that paused again gets another resume when it unpauses. - let pumped = self.state.has_pending_compressed() && !self.signals.is_receive_paused(); - if pumped { - let is_final = self.state.is_done(); - self.process_received_body(is_final)?; - } - - // A pump that ends the body has to say so even with no bytes (a stream trailer alone). - let ended = pumped && self.state.is_done() && !self.state.has_pending_compressed(); - Ok(!self.state.decoded_body.list.is_empty() || ended) - } - - /// The transport ended, and with it the body. What a consumer's budget holds stays held. - pub(crate) fn finish_body_on_close(&mut self) -> crate::Result<()> { - self.state.flags.received_last_chunk = true; - self.process_received_body(true).map(drop) - } - - /// Keeps this client reachable by id once its socket is gone, then delivers what it can. - fn outlive_socket(&mut self) { - self.state.flags.body_outlived_socket = true; - let _ = socketless_bodies().put(self.async_http_id, self.as_erased_ptr()); - if !self.state.decoded_body.list.is_empty() && self.send_progress_update_without_socket() { - self.drain_socketless_body(); - } - } - - /// A consumer's pull (`drain_queued_receive_resumes`) for a body that outlived its socket. - pub(crate) fn drain_socketless_body(&mut self) { - loop { - match self.pump_held_body() { - Ok(true) => {} - Ok(false) => return, - Err(err) => return self.fail(err), - } - if !self.send_progress_update_without_socket() { - return; + if self.state.has_pending_compressed() && !self.signals.is_receive_paused() { + // Not the final chunk: a body that is complete on the wire has ended the request. + if let Err(err) = self.process_received_body(false) { + self.close_and_fail::(err, socket); + return false; } } - } - /// An abort (`drain_queued_shutdowns`) for a body that outlived its socket. - pub(crate) fn abort_socketless_body(&mut self) { - self.fail(crate::Error::Aborted); + !self.state.decoded_body.list.is_empty() } fn send_progress_update_without_stage_check( @@ -4265,13 +4211,12 @@ impl<'a> HTTPClient<'a> { socket: HttpSocket, ) { if self.flags.protocol != Protocol::Http1_1 { - self.send_progress_update_without_socket(); - return; + return self.send_progress_update_multiplexed(); } // A loop, not a call back into `drain_response_body`: a consumer that never pauses // (`BufferAll`, or an S3 error body that is collected whole) takes one pass per turn. while self.send_one_progress_update::(ctx, socket) - && self.pump_held_body_or_close::(socket) + && self.pump_held_body::(socket) {} } @@ -4407,13 +4352,9 @@ impl<'a> HTTPClient<'a> { /// `send_progress_update_without_stage_check` minus the per-request TCP socket /// release/close. Used by HTTP/2 and HTTP/3, whose session owns the - /// transport, and by an h1 body that outlived its socket, so there is no - /// `ctx`/`socket` to hand back to the pool here. - /// Returns whether a held body is left that its consumer will not ask for. - fn send_progress_update_without_socket(&mut self) -> bool { - debug_assert!( - self.flags.protocol != Protocol::Http1_1 || self.state.flags.body_outlived_socket - ); + /// transport, so there is no `ctx`/`socket` to hand back to the pool here. + fn send_progress_update_multiplexed(&mut self) { + debug_assert!(self.flags.protocol != Protocol::Http1_1); let callback = self.result_callback; let mut result = self.to_result(); @@ -4434,7 +4375,7 @@ impl<'a> HTTPClient<'a> { if is_done { result.body_owned = decoded_body.list; callback.run(parent, result); - return false; + return; } result.body = decoded_body.list.as_slice(); callback.run(parent, result); @@ -4442,7 +4383,6 @@ impl<'a> HTTPClient<'a> { decoded_body.list.clear(); self.state.decoded_body = decoded_body; } - self.state.has_pending_compressed() && !self.signals.is_receive_paused() } /// `do_redirect` minus the per-request socket release/close. The session @@ -4494,7 +4434,7 @@ impl<'a> HTTPClient<'a> { } return; } - self.send_progress_update_without_socket(); + self.send_progress_update_multiplexed(); } pub(crate) fn do_redirect_h3(&mut self) { @@ -4612,6 +4552,13 @@ impl<'a> HTTPClient<'a> { None }; let certificate_info = self.state.certificate_info.take(); + // check if we are reporting cert errors, do not have a fail state and we are not done + let has_more = + certificate_info.is_some() || (self.state.fail.is_none() && !self.state.is_done()); + // The transport has all of the body. What is not decoded yet is the consumer's to decode. + let held_body = + (!has_more && self.state.fail.is_none() && self.state.has_pending_compressed()) + .then(|| Box::new(self.state.take_held_body())); if certificate_info.is_none() { if let Some(metadata) = self.state.cloned_metadata.take() { // transfer ownership of the metadata here @@ -4625,8 +4572,8 @@ impl<'a> HTTPClient<'a> { dns_hostname: self.state.dns_hostname.take(), connect_errno: self.state.connect_errno, proxy_connect_response: None, - has_more: self.state.fail.is_none() - && (!self.state.is_done() || self.state.has_pending_compressed()), + has_more, + held_body, body_size, certificate_info: None, can_stream: (self.state.request_stage == RequestStage::Body @@ -4646,10 +4593,8 @@ impl<'a> HTTPClient<'a> { dns_hostname: self.state.dns_hostname.take(), connect_errno: self.state.connect_errno, proxy_connect_response, - // check if we are reporting cert errors, do not have a fail state and we are not done - has_more: certificate_info.is_some() - || (self.state.fail.is_none() - && (!self.state.is_done() || self.state.has_pending_compressed())), + has_more, + held_body, body_size, certificate_info, // we can stream the request_body at this stage @@ -4677,7 +4622,7 @@ impl<'a> HTTPClient<'a> { && let Some(len) = content_length && incoming_data.len() >= len // The single-packet path decodes the whole body with no output budget. - && !(self.state.encoding.is_compressed() && self.signals.is_demand_driven()) + && !(self.state.encoding.is_compressed() && self.decodes_on_demand()) { self.handle_response_body_from_single_packet(&incoming_data[0..len])?; Ok(true) @@ -4769,7 +4714,6 @@ impl<'a> HTTPClient<'a> { // Close-delimited bodies still need per-packet decompression, but // a non-streaming consumer must not see per-packet progress: the // terminal callback (on close) is the first to carry metadata. - let is_done = is_done && !self.state.has_pending_compressed(); return Ok(is_done || (processed && is_streaming)); } Ok(false) @@ -4782,7 +4726,7 @@ impl<'a> HTTPClient<'a> { let small_len = 16 * 1024usize; if incoming_data.len() <= small_len && self.state.get_body_buffer().list.is_empty() - && !(self.state.encoding.is_compressed() && self.signals.is_demand_driven()) + && !(self.state.encoding.is_compressed() && self.decodes_on_demand()) { self.handle_response_body_chunked_encoding_from_single_packet(incoming_data) } else { @@ -4854,12 +4798,11 @@ impl<'a> HTTPClient<'a> { // Done _ => { self.state.flags.received_last_chunk = true; - let processed = self.process_received_body(true)?; + self.process_received_body(true)?; self.report_progress(buffer_len); - // A held body ends when `drain_response_body` has pumped it dry, not here. - return Ok(processed || !self.state.has_pending_compressed()); + return Ok(true); } } } diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 2eac893af7f4..19170ab1d851 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -117,6 +117,10 @@ pub struct FetchTasklet { /// buffer used to stream response to JS pub(crate) scheduled_response_buffer: MutableString, + /// `HTTPClientResult::held_body`, for `drive_held_body`. `None` while a `HeldBodyPass` has it. + pub(crate) held_body: JsCell>>, + /// A held body arrived: the HTTP thread is done with this fetch, the response body is not. + pub(crate) is_transport_done: bool, /// response weak ref we need this to track the response JS lifetime pub(crate) response: jsc::Weak, /// native response ref if we still need it when JS is discarted @@ -904,6 +908,10 @@ impl FetchTasklet { bun_output::scoped_log!(FetchTasklet, "onProgressUpdate"); self.mutex.lock(); self.has_schedule_callback.store(false, Ordering::Relaxed); + if let Some(held_body) = self.result.held_body.take() { + self.held_body.set(Some(held_body)); + self.is_transport_done = true; + } let is_done = !self.result.has_more; let vm = self.global_this.bun_vm(); @@ -920,7 +928,8 @@ impl FetchTasklet { } } self.mutex.unlock(); - if is_done { + // Nothing will read a held body either, and it is all that is left of this fetch. + if is_done || self.held_body.take().is_some() { // SAFETY: `self` is the live heap tasklet; we hold a ref. FetchTasklet::deref(std::ptr::from_mut(self)); } @@ -942,6 +951,10 @@ impl FetchTasklet { .with_mut(|poll_ref| poll_ref.unref(bun_io::js_vm_ctx())); // SAFETY: `this` is the live heap tasklet; we hold a ref. FetchTasklet::deref(std::ptr::from_mut(this)); + } else if this.is_transport_done { + // As above: nothing takes the rest of an upload. The response body goes on. + this.cancel_request_body_sink(JSValue::UNDEFINED); + this.drive_held_body(); } }; @@ -950,7 +963,7 @@ impl FetchTasklet { if let Err(err) = self.start_request_stream() { // The VM is being stopped: leave like the `!script_allowed()` gate above does. self.mutex.unlock(); - if is_done { + if is_done || self.held_body.take().is_some() { // SAFETY: `self` is the live heap tasklet; we hold a ref. FetchTasklet::deref(std::ptr::from_mut(self)); } @@ -1789,6 +1802,76 @@ impl FetchTasklet { if let Some(http_) = self.http.as_ref() { http::http_thread().schedule_receive_resume(http_.async_http_id); } + self.drive_held_body(); + } + + /// Produces a held body. Called wherever the HTTP thread would hear from the consumer. + fn drive_held_body(&self) { + if self.held_body.get().is_none() { + return; + } + let mode = self.signal_store.body_receive_mode(); + let unread = self.is_body_unread(mode); + if mode == BodyReceiveMode::Paused && !unread { + return; + } + let held = self.held_body.take(); + let cx = self + .global_this + .js_thread(self.global_this.bun_vm().context_of(self.context)); + self.ref_(); + jsc::Job::::schedule( + &cx, + HeldBodyPass { + // `None`: the pass only takes the end of the body to a turn that can report it. + held: held.filter(|_| !unread), + max_output: if mode == BodyReceiveMode::BufferAll { + usize::MAX + } else { + BODY_HIGH_WATER_MARK + }, + out: Vec::new(), + ended: Ok(false), + }, + HeldBodyPassOwner(std::ptr::from_ref(self).cast_mut()), + ); + } + + /// Nothing will read the rest of the body: its fetch was aborted, or its response is gone. + fn is_body_unread(&self, mode: BodyReceiveMode) -> bool { + mode == BodyReceiveMode::Abandoned || self.signal_store.aborted.load(Ordering::Relaxed) + } + + /// A `HeldBodyPass` is back: what `callback` and the task it posts do for the HTTP thread. + fn on_held_body_pass(&mut self, pass: HeldBodyPass) -> JsResult<()> { + let unread = self.is_body_unread(self.signal_store.body_receive_mode()); + self.mutex.lock(); + match (pass.held, pass.ended) { + (Some(held), Ok(ended)) if !unread => { + self.result.has_more = !ended; + if !ended { + self.held_body.set(Some(held)); + } + let scheduled = &mut self.scheduled_response_buffer.list; + if scheduled.is_empty() { + *scheduled = pass.out; + } else { + scheduled.extend_from_slice(&pass.out); + } + if !ended && scheduled.len() >= BODY_HIGH_WATER_MARK { + self.signal_store.pause_receive(); + } + } + (_, ended) => { + self.result.has_more = false; + self.result.fail = Some(match ended { + Err(err) if !unread => err, + _ => http::Error::Aborted, + }); + } + } + self.mutex.unlock(); + self.on_progress_update() } fn to_body_value(&mut self) -> BodyValue { @@ -1949,6 +2032,8 @@ impl FetchTasklet { request_body: fetch_options.body, request_body_streaming_buffer: None, scheduled_response_buffer: MutableString::default(), + held_body: JsCell::new(None), + is_transport_done: false, response: jsc::Weak::default(), native_response: JsCell::new(None), response_stream: Default::default(), @@ -2097,6 +2182,7 @@ impl FetchTasklet { idle_timeout_seconds: fetch_options.idle_timeout_seconds, disable_keepalive: Some(fetch_options.disable_keepalive), disable_decompression: Some(fetch_options.disable_decompression), + takes_held_body: true, max_redirects: fetch_options.max_redirects, reject_unauthorized: Some(fetch_options.reject_unauthorized), verbose: Some(fetch_options.verbose), @@ -2240,7 +2326,7 @@ impl FetchTasklet { data: RequestBodyChunk<'_>, high_water_mark: usize, ) -> Writable { - if self.signal_aborted() { + if self.signal_aborted() || self.is_transport_done { return Writable::Done; } // An empty chunk is a no-op on every framing path. It must not reach @@ -2304,8 +2390,10 @@ impl FetchTasklet { pub(crate) fn write_end_request(&mut self, err: Option) { bun_output::scoped_log!(FetchTasklet, "writeEndRequest hasError? {}", err.is_some()); let this_ptr = std::ptr::from_mut(self); + // An upload that ends after its response is complete on the wire changes nothing. + let is_over = self.is_transport_done || self.signal_store.aborted.load(Ordering::Relaxed); if let Some(js_error) = err { - if self.signal_store.aborted.load(Ordering::Relaxed) || self.abort_reason.has() { + if is_over || self.abort_reason.has() { // SAFETY: `this_ptr` derived from live `&mut self`; we hold a ref. FetchTasklet::deref(this_ptr); return; @@ -2315,7 +2403,7 @@ impl FetchTasklet { } self.abort_task(); } else { - if self.signal_store.aborted.load(Ordering::Relaxed) { + if is_over { // SAFETY: `this_ptr` derived from live `&mut self`; we hold a ref. FetchTasklet::deref(this_ptr); return; @@ -2357,6 +2445,7 @@ impl FetchTasklet { if let Some(http_) = self.http.as_deref() { http::http_thread().schedule_shutdown(http_); } + self.drive_held_body(); true } @@ -2384,7 +2473,10 @@ impl FetchTasklet { let global_this = self.global_this; self.abort_reason.set(&global_this, reason); } - self.abort_task(); + // An abort would take the rest of the response body with it. + if !self.is_transport_done { + self.abort_task(); + } if let Some(sink) = self.sink_mut() { sink.pending.result = Writable::Done; sink.pending.run(); @@ -2395,7 +2487,7 @@ impl FetchTasklet { } if is_native { // No pump promise exists to balance the `+1` from - // `start_request_stream`; `aborted` is set above so + // `start_request_stream`; the fetch is over (see `write_end_request`), so // `write_end_request(Some(_))` is just the balancing deref. self.write_end_request(Some(reason)); } @@ -2483,6 +2575,16 @@ impl FetchTasklet { // SAFETY: lifetime erasure for non-body fields; `body` is stored as // `&'static []` so no borrow escapes. task_ref.result = unsafe { result.detach_lifetime() }; + // Read once: the JS thread can abandon the body at any point of this callback. + let abandoned = task_ref.signal_store.body_receive_mode() == BodyReceiveMode::Abandoned; + // The transport is done, the body is not: `on_progress_update` takes the rest from here. + if task_ref.result.held_body.is_some() { + if abandoned { + task_ref.result.held_body = None; + } else { + task_ref.result.has_more = true; + } + } // can_stream is a one-shot signal to start the request body stream; don't let a // later coalesced result clobber it before the JS thread sees it. task_ref.result.can_stream = task_ref.result.can_stream || prev_can_stream; @@ -2508,7 +2610,7 @@ impl FetchTasklet { let success = task_ref.result.is_success(); - if task_ref.signal_store.body_receive_mode() == BodyReceiveMode::Abandoned { + if abandoned { if task_ref.scheduled_response_buffer.list.capacity() > 0 { task_ref.scheduled_response_buffer = MutableString::default(); } @@ -2680,6 +2782,48 @@ impl FetchTasklet { } } +/// One decode pass over a held body, on the work pool. +struct HeldBodyPass { + held: Option>, + max_output: usize, + out: Vec, + ended: Result, +} + +/// The tasklet a `HeldBodyPass` reports to, and a ref on it. +struct HeldBodyPassOwner(*mut FetchTasklet); + +// SAFETY: a ref on the tasklet, which is the JS thread's; used and dropped there. +unsafe impl bun_jsc::job::JsAffine for HeldBodyPassOwner {} + +impl Drop for HeldBodyPassOwner { + /// Released unrun: nothing reports the end of the body, so the ref it would drop goes too. + fn drop(&mut self) { + FetchTasklet::deref(self.0); + FetchTasklet::deref(self.0); + } +} + +impl jsc::JobContext for HeldBodyPass { + type OffThread = Self; + type Js = HeldBodyPassOwner; + + fn run(pass: &mut Self, done: jsc::Completion) -> Option> { + if let Some(held) = pass.held.as_mut() { + pass.ended = held.decode(&mut pass.out, pass.max_output); + } + Some(done) + } + + fn then(pass: Self, owner: HeldBodyPassOwner, _: &jsc::JsThread<'_>) -> JsResult<()> { + let tasklet = core::mem::ManuallyDrop::new(owner).0; + // SAFETY: the pass holds a ref; JS thread. + let result = unsafe { (*tasklet).on_held_body_pass(pass) }; + FetchTasklet::deref(tasklet); + result + } +} + pub struct FetchOptions { pub method: Method, pub(crate) headers: Headers, diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 346d0fff29a5..3dd4bc7e7087 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -685,7 +685,7 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not // bun does not pause a tunnelled socket, so all of the body and then the origin's close reach // the client at once. The end of the transport must not decode what the reader has not asked - // for: the client outlives its socket until the reader has pulled the rest, or cancels. + // for: the request ends there, and the Response keeps the rest of the body undecoded. test.each(["never", "content-length", "chunked", "close-delimited"] as Ending[])( "zstd through a CONNECT proxy, body ending %s: a reader that takes a little holds a little", async ending => { @@ -697,11 +697,107 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }, ); - // Without a tunnel the socket is paused, which defers a FIN but not a TLS close_notify that - // came in with the body. - test.each([false, true])("gzip, origin closes after the body (tls: %p)", async secure => { - await using server = await serveBomb("gzip", secure, "content-length"); - expectBounded(await runClient(server.url, { tls: { rejectUnauthorized: false } }, READ_A_LITTLE)); + // The origin sends the whole of a body at once and keeps the connection: the request ends with + // most of the body undecoded, and the Response keeps that part. + async function serveWhole(body: Buffer) { + const head = `HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: ${body.length}\r\n\r\n`; + const server = await listening( + createTcpServer(s => { + s.on("error", () => {}); + // Answers each request head. What else arrives is an upload it does not wait for. + s.on("data", data => { + if (data.includes(" HTTP/1.1\r\n")) s.write(Buffer.concat([Buffer.from(head), body])); + }); + }), + ); + return { ...server, url: `http://127.0.0.1:${server.port}/` }; + } + + // The reader is eight budgets into the body when it reaches the damage, long after the + // request ended. The error has to reach it there: no hang, and no body that just stops. + test("gzip: damage in the part the Response kept undecoded rejects the read", async () => { + const member = gzipSync(Buffer.alloc(8 * MARK), { level: 9 }); + const damaged = Buffer.from(member); + damaged.fill(0xff, 20, 60); + await using server = await serveWhole(Buffer.concat([member, damaged])); + const script = /* js */ ` + const res = await fetch(url, opts); + let got = 0, code; + try { + for await (const chunk of res.body) got += chunk.byteLength; + } catch (e) { + code = e.code; + } + process.stdout.write(JSON.stringify({ got, code })); + `; + const { got, code, exitCode } = await runClient(server.url, {}, script); + expect({ got, code }).toEqual({ got: 8 * MARK, code: "ZlibError" }); + expect(exitCode).toBe(0); + }); + + test("an abort reaches the part of a body that the Response kept undecoded", async () => { + await using server = await serveWhole(gzipSync(Buffer.alloc(64 * MARK))); + + const reading = new AbortController(); + const reader = (await fetch(server.url, { signal: reading.signal })).body!.getReader(); + const first = await reader.read(); + reading.abort(); + const rest = (async () => { + while (!(await reader.read()).done); + return "ended"; + })().catch(e => e.name); + + const untouched = new AbortController(); + const res = await fetch(server.url, { signal: untouched.signal }); + untouched.abort(); + + expect({ + first: first.value!.byteLength > 0, + rest: await rest, + untouched: await res.text().catch(e => e.name), + }).toEqual({ first: true, rest: "AbortError", untouched: "AbortError" }); + }); + + // The request is over when its response is complete on the wire, so an upload that is still + // underway ends there, as it does for a body that is not compressed. The Response is not over. + test("an upload still underway ends with the request, not with the part the Response kept undecoded", async () => { + await using server = await serveWhole(gzipSync(Buffer.alloc(16 * MARK, 65))); + const cancelled = Promise.withResolvers(); + const upload = new ReadableStream({ + pull: controller => controller.enqueue(new Uint8Array(1024)), + cancel: () => cancelled.resolve(), + }); + const res = await fetch(server.url, { method: "POST", body: upload, duplex: "half" } as RequestInit); + await cancelled.promise; + const bytes = await res.bytes(); + expect({ length: bytes.byteLength, allA: isAllA(bytes) }).toEqual({ length: 16 * MARK, allA: true }); + }); + + // One Response keeps its part at rest; the other has a reader, so a decode pass is out on + // another thread when the worker goes. + test("a worker can go away while its Responses keep parts of bodies undecoded", async () => { + await using server = await serveWhole(gzipSync(Buffer.alloc(256 * MARK))); + const worker = /* js */ ` + const { parentPort, workerData } = require("node:worker_threads"); + (async () => { + globalThis.atRest = await fetch(workerData.url); + const reader = (await fetch(workerData.url)).body.getReader(); + await reader.read(); + parentPort.postMessage("ready"); + while (!(await reader.read()).done); + })(); + `; + const script = /* js */ ` + const { Worker } = require("node:worker_threads"); + const worker = new Worker(${JSON.stringify(worker)}, { eval: true, workerData: { url } }); + const ready = await new Promise(resolve => worker.once("message", resolve)); + await worker.terminate(); + process.stdout.write(JSON.stringify({ ready })); + `; + const { ready, stderr, exitCode } = await runClient(server.url, {}, script); + expect(stderr).toBe(""); + expect(ready).toBe("ready"); + expect(exitCode).toBe(0); }); // A live stream: two flushed messages in one packet, then the origin goes quiet with the frame @@ -818,28 +914,27 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not // Content-Length bodies end with the origin's FIN, which reaches a client that still holds // most of the body undecoded. Chunked ones stay open. - async function serveBody(kind: Kind, chunked: boolean) { + async function serveBody(kind: Kind, chunked: boolean, secure = false) { const body = bodyFor(kind); const enc = kind === "br-hq" ? "br" : kind; - const server = await listening( - createTcpServer(s => { - s.on("error", () => {}); - s.once("data", () => { - if (chunked) { - s.write(`HTTP/1.1 200 OK\r\nContent-Encoding: ${enc}\r\nTransfer-Encoding: chunked\r\n\r\n`); - s.write(`${body.length.toString(16)}\r\n`); - s.write(body); - s.write("\r\n0\r\n\r\n"); - } else { - s.write( - `HTTP/1.1 200 OK\r\nContent-Encoding: ${enc}\r\nContent-Length: ${body.length}\r\nConnection: close\r\n\r\n`, - ); - s.end(body); - } - }); - }), - ); - return { ...server, url: `http://127.0.0.1:${server.port}/` }; + const handler = (s: import("node:net").Socket) => { + s.on("error", () => {}); + s.once("data", () => { + if (chunked) { + s.write(`HTTP/1.1 200 OK\r\nContent-Encoding: ${enc}\r\nTransfer-Encoding: chunked\r\n\r\n`); + s.write(`${body.length.toString(16)}\r\n`); + s.write(body); + s.write("\r\n0\r\n\r\n"); + } else { + s.write( + `HTTP/1.1 200 OK\r\nContent-Encoding: ${enc}\r\nContent-Length: ${body.length}\r\nConnection: close\r\n\r\n`, + ); + s.end(body); + } + }); + }; + const server = await listening(secure ? createTlsServer(tls, handler) : createTcpServer(handler)); + return { ...server, url: `${secure ? "https" : "http"}://127.0.0.1:${server.port}/` }; } const STREAM = /* js */ ` @@ -885,6 +980,22 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not expect(exitCode).toBe(0); }); }); + + // Through a tunnel the whole body and the origin's close are in before the reader's second + // pull, so nearly all of it comes from the part that the Response kept undecoded. + test.each([ + ["a streaming reader", STREAM, { several: true }], + ["res.bytes()", BUFFER, {}], + ])("zstd through a CONNECT proxy, origin closes after the body: %s", async (_, script, seen) => { + await using server = await serveBody("zstd", false, true); + await using proxy = await serveConnectProxy(); + const opts = { proxy: proxy.url, tls: { rejectUnauthorized: false } }; + const { stderr, exitCode, ...result } = await runClient(server.url, opts, script); + expect(stderr).toBe(""); + expect(result).toEqual({ total: SIZE, digest, ...seen }); + expect(proxy.connects()).toBe(1); + expect(exitCode).toBe(0); + }); }); }); @@ -1662,6 +1773,46 @@ describe.serial("fetch() receive backpressure — an unread body hands its conne } }); + // A compressed body is decoded as it is read, so all of it can be in while most of it is still + // undecoded. That ends the request: the Response keeps the rest, the connection goes back. + test("a compressed body that is complete on the wire, Response still held: its connection is reused", async () => { + const SIZE = 4 * MARK; + const body = gzipSync(Buffer.alloc(SIZE, 65)); + let connections = 0; + const sockets = new Set(); + const srv = createTcpServer(socket => { + connections++; + sockets.add(socket); + socket.on("error", () => {}); + socket.on("data", () => { + socket.write( + Buffer.concat([ + Buffer.from(`HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: ${body.length}\r\n\r\n`), + body, + ]), + ); + }); + }); + srv.listen(0, "127.0.0.1"); + await once(srv, "listening"); + try { + const url = `http://127.0.0.1:${(srv.address() as import("node:net").AddressInfo).port}/`; + const N = 8; + const responses: Response[] = []; + for (let i = 0; i < N; i++) responses.push(await fetch(url)); + // Not exactly 1: a request can start before the previous body's last packet was taken. + // A body that pins its connection until it is read needs N of them. + expect(connections).toBeLessThan(N / 2); + for (const res of responses) { + const bytes = await res.bytes(); + expect({ length: bytes.byteLength, allA: isAllA(bytes) }).toEqual({ length: SIZE, allA: true }); + } + } finally { + for (const socket of sockets) socket.destroy(); + 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 () => { From f8ef451e090b5163ba2cf7f2b1fe8655bc313ad3 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 02:08:22 +0000 Subject: [PATCH 03/10] http: keep the decode budget of a consumer that takes no HeldBody while the wire is live The S3 download stream reads under the same backpressure signals as fetch, and takes no HeldBody. It keeps its per-pass budget until the last chunk. The last chunk is decoded in full for it, as a close did before. --- src/http/AsyncHTTP.rs | 2 +- src/http/lib.rs | 13 ++++++------- test/js/web/fetch/fetch-backpressure.test.ts | 12 ++++++++++-- 3 files changed, 17 insertions(+), 10 deletions(-) diff --git a/src/http/AsyncHTTP.rs b/src/http/AsyncHTTP.rs index 7c87254e960e..8d3e886c4e51 100644 --- a/src/http/AsyncHTTP.rs +++ b/src/http/AsyncHTTP.rs @@ -253,7 +253,7 @@ pub struct Options<'a> { pub verbose: Option, pub disable_keepalive: Option, pub disable_decompression: Option, - /// The consumer takes `HTTPClientResult::held_body`: a compressed body is decoded as it reads. + /// The consumer takes `HTTPClientResult::held_body`. Others get it decoded with the last chunk. pub takes_held_body: bool, pub max_redirects: Option, pub reject_unauthorized: Option, diff --git a/src/http/lib.rs b/src/http/lib.rs index 73cce6febd73..b9b01f120fe9 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -4091,15 +4091,13 @@ impl<'a> HTTPClient<'a> { /// A compressed body is decoded one budget at a time, as its consumer reads. h1 only. #[inline] fn decodes_on_demand(&self) -> bool { - self.flags.takes_held_body - && self.flags.protocol == Protocol::Http1_1 - && self.signals.is_demand_driven() + self.flags.protocol == Protocol::Http1_1 && self.signals.is_demand_driven() } - /// Output budget of one decode pass. + /// Output budget of one decode pass. None for the last of a body that has no `HeldBody` taker. #[inline] - fn decompress_output_cap(&self) -> usize { - if self.decodes_on_demand() { + fn decompress_output_cap(&self, is_final_chunk: bool) -> usize { + if self.decodes_on_demand() && (self.flags.takes_held_body || !is_final_chunk) { signals::BODY_HIGH_WATER_MARK } else { usize::MAX @@ -4108,7 +4106,7 @@ impl<'a> HTTPClient<'a> { /// Decodes what has arrived under the consumer's budget. Returns whether to report bytes. fn process_received_body(&mut self, is_final_chunk: bool) -> crate::Result { - let max_output = self.decompress_output_cap(); + let max_output = self.decompress_output_cap(is_final_chunk); if self.state.encoding.is_compressed() { // Nothing will read it, and its transport is being shut down. if self.signals.is_body_abandoned() { @@ -4559,6 +4557,7 @@ impl<'a> HTTPClient<'a> { let held_body = (!has_more && self.state.fail.is_none() && self.state.has_pending_compressed()) .then(|| Box::new(self.state.take_held_body())); + debug_assert!(held_body.is_none() || self.flags.takes_held_body); if certificate_info.is_none() { if let Some(metadata) = self.state.cloned_metadata.take() { // transfer ownership of the metadata here diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 3dd4bc7e7087..682ce97db642 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -617,8 +617,10 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not const base = process.memoryUsage.rss(); let peak = 0; const sample = () => (peak = Math.max(peak, process.memoryUsage.rss() - base)); - const res = await fetch(url, opts); - const reader = res.body.getReader(); + const body = opts.s3 + ? new Bun.S3Client({ accessKeyId: "t", secretAccessKey: "t", endpoint: url, bucket: "b" }).file("k").stream() + : (await fetch(url, opts)).body; + const reader = body.getReader(); const first = await reader.read(); for (let last = sample(), stable = 0; stable < 3; ) { await Bun.sleep(20); @@ -683,6 +685,12 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }, ); + // The S3 client reads a bucket's compressed object through the same decoder. + test("gzip from S3: a reader that takes a little holds a little", async () => { + await using server = await serveBomb("gzip", false); + expectBounded(await runClient(server.url, { s3: true }, READ_A_LITTLE)); + }); + // bun does not pause a tunnelled socket, so all of the body and then the origin's close reach // the client at once. The end of the transport must not decode what the reader has not asked // for: the request ends there, and the Response keeps the rest of the body undecoded. From 2a8a1315bfc8e7d8dae4b04205278d7237e592b4 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 03:08:32 +0000 Subject: [PATCH 04/10] fetch: end an unread held body in a task that runs in a stopped context too A decode pass whose context stopped is released unrun, and the end of an unread body was such a pass. The tasklet was then freed with the Response body still locked on it. The end of the body is now the tasklet's own task, which runs whatever became of the context, as the HTTP thread's last result does. A pass that is released unrun queues the same task. --- src/runtime/webcore/fetch/FetchTasklet.rs | 85 ++++++++++++++------ test/js/web/fetch/fetch-backpressure.test.ts | 29 +++++++ 2 files changed, 90 insertions(+), 24 deletions(-) diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 19170ab1d851..7dee3a93fb80 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -117,8 +117,8 @@ pub struct FetchTasklet { /// buffer used to stream response to JS pub(crate) scheduled_response_buffer: MutableString, - /// `HTTPClientResult::held_body`, for `drive_held_body`. `None` while a `HeldBodyPass` has it. - pub(crate) held_body: JsCell>>, + /// `HTTPClientResult::held_body`, from when `on_progress_update` has it. JS thread only. + held_body: JsCell, /// A held body arrived: the HTTP thread is done with this fetch, the response body is not. pub(crate) is_transport_done: bool, /// response weak ref we need this to track the response JS lifetime @@ -909,9 +909,14 @@ impl FetchTasklet { self.mutex.lock(); self.has_schedule_callback.store(false, Ordering::Relaxed); if let Some(held_body) = self.result.held_body.take() { - self.held_body.set(Some(held_body)); + self.held_body.set(HeldBodyState::AtRest(held_body)); self.is_transport_done = true; } + if matches!(self.held_body.get(), HeldBodyState::Unread) { + self.held_body.set(HeldBodyState::None); + self.result.has_more = false; + self.result.fail = Some(http::Error::Aborted); + } let is_done = !self.result.has_more; let vm = self.global_this.bun_vm(); @@ -929,7 +934,7 @@ impl FetchTasklet { } self.mutex.unlock(); // Nothing will read a held body either, and it is all that is left of this fetch. - if is_done || self.held_body.take().is_some() { + if is_done || self.drop_held_body_at_rest() { // SAFETY: `self` is the live heap tasklet; we hold a ref. FetchTasklet::deref(std::ptr::from_mut(self)); } @@ -963,7 +968,7 @@ impl FetchTasklet { if let Err(err) = self.start_request_stream() { // The VM is being stopped: leave like the `!script_allowed()` gate above does. self.mutex.unlock(); - if is_done || self.held_body.take().is_some() { + if is_done || self.drop_held_body_at_rest() { // SAFETY: `self` is the live heap tasklet; we hold a ref. FetchTasklet::deref(std::ptr::from_mut(self)); } @@ -1807,15 +1812,19 @@ impl FetchTasklet { /// Produces a held body. Called wherever the HTTP thread would hear from the consumer. fn drive_held_body(&self) { - if self.held_body.get().is_none() { + if !matches!(self.held_body.get(), HeldBodyState::AtRest(_)) { return; } let mode = self.signal_store.body_receive_mode(); - let unread = self.is_body_unread(mode); - if mode == BodyReceiveMode::Paused && !unread { + if self.is_body_unread(mode) { + return self.end_unread_held_body(); + } + if mode == BodyReceiveMode::Paused { return; } - let held = self.held_body.take(); + let HeldBodyState::AtRest(held) = self.held_body.replace(HeldBodyState::InPass) else { + return; + }; let cx = self .global_this .js_thread(self.global_this.bun_vm().context_of(self.context)); @@ -1823,8 +1832,7 @@ impl FetchTasklet { jsc::Job::::schedule( &cx, HeldBodyPass { - // `None`: the pass only takes the end of the body to a turn that can report it. - held: held.filter(|_| !unread), + held, max_output: if mode == BodyReceiveMode::BufferAll { usize::MAX } else { @@ -1842,16 +1850,34 @@ impl FetchTasklet { mode == BodyReceiveMode::Abandoned || self.signal_store.aborted.load(Ordering::Relaxed) } + /// Queues the end of the body. This task runs in a stopped context too: the body gets settled. + fn end_unread_held_body(&self) { + self.held_body.set(HeldBodyState::Unread); + let task = Task::init(std::ptr::from_ref(self).cast_mut()); + // SAFETY: JS thread; the loop is this VM's. + unsafe { (*self.global_this.bun_vm().event_loop()).enqueue_task(task) }; + } + + fn drop_held_body_at_rest(&self) -> bool { + let at_rest = matches!(self.held_body.get(), HeldBodyState::AtRest(_)); + if at_rest { + self.held_body.set(HeldBodyState::None); + } + at_rest + } + /// A `HeldBodyPass` is back: what `callback` and the task it posts do for the HTTP thread. fn on_held_body_pass(&mut self, pass: HeldBodyPass) -> JsResult<()> { let unread = self.is_body_unread(self.signal_store.body_receive_mode()); self.mutex.lock(); - match (pass.held, pass.ended) { - (Some(held), Ok(ended)) if !unread => { + match pass.ended { + Ok(ended) if !unread => { self.result.has_more = !ended; - if !ended { - self.held_body.set(Some(held)); - } + self.held_body.set(if ended { + HeldBodyState::None + } else { + HeldBodyState::AtRest(pass.held) + }); let scheduled = &mut self.scheduled_response_buffer.list; if scheduled.is_empty() { *scheduled = pass.out; @@ -1862,7 +1888,8 @@ impl FetchTasklet { self.signal_store.pause_receive(); } } - (_, ended) => { + ended => { + self.held_body.set(HeldBodyState::None); self.result.has_more = false; self.result.fail = Some(match ended { Err(err) if !unread => err, @@ -2032,7 +2059,7 @@ impl FetchTasklet { request_body: fetch_options.body, request_body_streaming_buffer: None, scheduled_response_buffer: MutableString::default(), - held_body: JsCell::new(None), + held_body: JsCell::new(HeldBodyState::None), is_transport_done: false, response: jsc::Weak::default(), native_response: JsCell::new(None), @@ -2782,9 +2809,20 @@ impl FetchTasklet { } } +/// Where the rest of a body is once its transport is done. +enum HeldBodyState { + /// The HTTP thread produces the body, or the body has ended. + None, + AtRest(Box), + /// A `HeldBodyPass` has it. + InPass, + /// Nothing will read it, and the `on_progress_update` that ends the body is queued. + Unread, +} + /// One decode pass over a held body, on the work pool. struct HeldBodyPass { - held: Option>, + held: Box, max_output: usize, out: Vec, ended: Result, @@ -2797,9 +2835,10 @@ struct HeldBodyPassOwner(*mut FetchTasklet); unsafe impl bun_jsc::job::JsAffine for HeldBodyPassOwner {} impl Drop for HeldBodyPassOwner { - /// Released unrun: nothing reports the end of the body, so the ref it would drop goes too. + /// Released unrun: the context that fetched stopped, or the VM is going. fn drop(&mut self) { - FetchTasklet::deref(self.0); + // SAFETY: the pass holds a ref; JS thread. + unsafe { (*self.0).end_unread_held_body() }; FetchTasklet::deref(self.0); } } @@ -2809,9 +2848,7 @@ impl jsc::JobContext for HeldBodyPass { type Js = HeldBodyPassOwner; fn run(pass: &mut Self, done: jsc::Completion) -> Option> { - if let Some(held) = pass.held.as_mut() { - pass.ended = held.decode(&mut pass.out, pass.max_output); - } + pass.ended = pass.held.decode(&mut pass.out, pass.max_output); Some(done) } diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 682ce97db642..46e93fa6106e 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -781,6 +781,35 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not expect({ length: bytes.byteLength, allA: isAllA(bytes) }).toEqual({ length: 16 * MARK, allA: true }); }); + // A host may keep a Response that a Bun.ModuleGraph's code fetched. dispose() aborts the fetch, + // and the body the host still holds has to end with it before what produced it is freed. + test("a Response outlives the Bun.ModuleGraph that fetched it: the part it kept undecoded ends", async () => { + await using server = await serveWhole(gzipSync(Buffer.alloc(64 * MARK))); + using dir = tempDir("held-body-graph", { + "app.ts": `export const get = (url: string) => fetch(url);`, + "host.ts": ` + const graph = new Bun.ModuleGraph({}); + const app = await graph.import(import.meta.dir + "/app.ts"); + const res: Response = await graph.run(() => app.get(process.argv[2])); + graph.dispose(); + // dispose() queued the task that ends the body. It runs before the loop gets here. + for (let i = 0; i < 3; i++) await new Promise(resolve => setImmediate(resolve)); + process.stdout.write(await res.text().catch(e => e.name)); + `, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "host.ts", server.url], + env: clientEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + expect(stdout).toBe("AbortError"); + expect(exitCode).toBe(0); + }); + // One Response keeps its part at rest; the other has a reader, so a decode pass is out on // another thread when the worker goes. test("a worker can go away while its Responses keep parts of bodies undecoded", async () => { From 31132d642ee6f00751222678a4c6d8fe917e2b95 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 03:48:24 +0000 Subject: [PATCH 05/10] fetch: reach a held body's deferred work through the tasklet's allocation pointer A pass and the task that ends an unread body form &mut and can drop the last ref, so their pointer must not come from a &self. The append of a pass goes through handle_oom, like the HTTP thread's. --- src/runtime/webcore/fetch/FetchTasklet.rs | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 7dee3a93fb80..f7e4799f6785 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -119,6 +119,8 @@ pub struct FetchTasklet { pub(crate) scheduled_response_buffer: MutableString, /// `HTTPClientResult::held_body`, from when `on_progress_update` has it. JS thread only. held_body: JsCell, + /// What `heap::into_raw` returned for this tasklet, for work that outlives a `&self` caller. + allocation: *mut FetchTasklet, /// A held body arrived: the HTTP thread is done with this fetch, the response body is not. pub(crate) is_transport_done: bool, /// response weak ref we need this to track the response JS lifetime @@ -1841,7 +1843,7 @@ impl FetchTasklet { out: Vec::new(), ended: Ok(false), }, - HeldBodyPassOwner(std::ptr::from_ref(self).cast_mut()), + HeldBodyPassOwner(self.allocation), ); } @@ -1853,7 +1855,7 @@ impl FetchTasklet { /// Queues the end of the body. This task runs in a stopped context too: the body gets settled. fn end_unread_held_body(&self) { self.held_body.set(HeldBodyState::Unread); - let task = Task::init(std::ptr::from_ref(self).cast_mut()); + let task = Task::init(self.allocation); // SAFETY: JS thread; the loop is this VM's. unsafe { (*self.global_this.bun_vm().event_loop()).enqueue_task(task) }; } @@ -1878,13 +1880,13 @@ impl FetchTasklet { } else { HeldBodyState::AtRest(pass.held) }); - let scheduled = &mut self.scheduled_response_buffer.list; - if scheduled.is_empty() { - *scheduled = pass.out; + let scheduled = &mut self.scheduled_response_buffer; + if scheduled.list.is_empty() { + scheduled.list = pass.out; } else { - scheduled.extend_from_slice(&pass.out); + bun_core::handle_oom(scheduled.write(&pass.out)); } - if !ended && scheduled.len() >= BODY_HIGH_WATER_MARK { + if !ended && scheduled.list.len() >= BODY_HIGH_WATER_MARK { self.signal_store.pause_receive(); } } @@ -2060,6 +2062,7 @@ impl FetchTasklet { request_body_streaming_buffer: None, scheduled_response_buffer: MutableString::default(), held_body: JsCell::new(HeldBodyState::None), + allocation: core::ptr::null_mut(), is_transport_done: false, response: jsc::Weak::default(), native_response: JsCell::new(None), @@ -2150,6 +2153,7 @@ impl FetchTasklet { let fetch_tasklet_ptr = bun_core::heap::into_raw(fetch_tasklet); // SAFETY: just allocated; exclusive access until returned let fetch_tasklet = unsafe { &mut *fetch_tasklet_ptr }; + fetch_tasklet.allocation = fetch_tasklet_ptr; // This task gets queued on the HTTP thread. // `AsyncHTTP::init` takes several `&'static [u8]` borrows From 6b5d6da73428eb6f76e351d1fd141f9f8bc1759f Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Wed, 23 Sep 2026 20:12:31 -0700 Subject: [PATCH 06/10] fetch: move a large held-body pass into the response buffer instead of copying it When .text() takes the rest of a held body, the decoded output can be far larger than the bytes already buffered. Copying it onto the buffer doubled peak memory. The buffered bytes now go in front of the output, and the output becomes the buffer. The collected-Response test now runs its client without proxy variables, and its GC loop is bounded so that a tasklet that is never freed fails the count. --- src/runtime/webcore/fetch/FetchTasklet.rs | 4 ++++ test/js/web/fetch/fetch-backpressure.test.ts | 5 +++-- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 9efbcc365458..6584cec46029 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -1926,6 +1926,10 @@ impl FetchTasklet { let scheduled = &mut self.scheduled_response_buffer; if scheduled.list.is_empty() { scheduled.list = pass.out; + } else if scheduled.list.len() < pass.out.len() { + let mut out = pass.out; + out.splice(0..0, scheduled.list.drain(..)); + scheduled.list = out; } else { bun_core::handle_oom(scheduled.write(&pass.out)); } diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 68a8eb0f291b..f7f994217fef 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -1287,7 +1287,8 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not (async () => { for await (const line of console) if (line === "done") done = true; })(); - while (!done) { + // Bounded, so that a tasklet that is never freed fails the count below. + for (let turn = 0; turn < 200 && !done; turn++) { Bun.gc(true); await new Promise(resolve => setImmediate(resolve)); } @@ -1295,7 +1296,7 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not `; await using proc = Bun.spawn({ cmd: [bunExe(), "-e", script], - env: { ...bunEnv, BUN_DEBUG_FetchTasklet: "1", BUN_DEBUG_HTTPInternalState: "1" }, + env: { ...clientEnv, BUN_DEBUG_FetchTasklet: "1", BUN_DEBUG_HTTPInternalState: "1" }, stdin: "pipe", stdout: "pipe", stderr: "pipe", From 9cbd14609a70a18c362d1b6152a36f659a69c4e6 Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Wed, 23 Sep 2026 20:20:12 -0700 Subject: [PATCH 07/10] http: log each zlib pass over a held body The debug-log tests in fetch-backpressure.test.ts count decode passes. A zlib pass over a held body wrote no log line, so a fallback from the one libdeflate call to zlib was invisible to them. --- src/http/InternalState.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/http/InternalState.rs b/src/http/InternalState.rs index 3d422bf05f33..e33fe6ad8d1a 100644 --- a/src/http/InternalState.rs +++ b/src/http/InternalState.rs @@ -526,6 +526,7 @@ impl HeldBody { return Ok(true); } let input = &self.input[self.consumed..]; + log!("Decompressing {} bytes\n", input.len()); self.consumed += self.decompressor .decompress_chunk(self.encoding, input, out, max_output, true)?; From b3498578646856eeab41595cb6046d2304996eb7 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 24 Sep 2026 04:01:57 +0000 Subject: [PATCH 08/10] test: a body that is read shows that the decode-pass scan sees a pass The collected-Response case asserts that no pass is logged while the bodies are held. After they are freed, the client now reads one more body, and the test expects exactly one pass from it. --- test/js/web/fetch/fetch-backpressure.test.ts | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index f7f994217fef..b16c3ded59fc 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -1278,9 +1278,10 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }); await using server = await listening(srv); const script = /* js */ ` + const url = ${JSON.stringify(`http://127.0.0.1:${server.port}/`)}; // Its own frame, so that nothing on the caller's still refers to the response afterwards. async function abandon() { - if ((await fetch(${JSON.stringify(`http://127.0.0.1:${server.port}/`)})).status !== 200) throw new Error("status"); + if ((await fetch(url)).status !== 200) throw new Error("status"); } for (let i = 0; i < ${N}; i++) await abandon(); let done = false; @@ -1292,6 +1293,8 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not Bun.gc(true); await new Promise(resolve => setImmediate(resolve)); } + // A body that is read logs its pass, so the scan below is known to see one. + await (await fetch(url)).bytes(); process.exit(0); `; await using proc = Bun.spawn({ @@ -1302,19 +1305,25 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not stderr: "pipe", }); let freed = 0; + let passesWhileHeld = -1; const passes: string[] = []; // Output.scoped writes to whichever stream it chose at init; scan both. const scan = async (stream: ReadableStream) => { for await (const line of forEachLine(stream)) { if (/Decompressing \d+ bytes/.test(line)) passes.push(line); if (/\[FetchTasklet\] deinit/i.test(line) && ++freed === N) { + passesWhileHeld = passes.length; proc.stdin.write("done\n"); proc.stdin.end(); } } }; await Promise.all([scan(proc.stdout), scan(proc.stderr)]); - expect({ freed, passes }).toEqual({ freed: N, passes: [] }); + expect({ freed: Math.min(freed, N), passesWhileHeld, passesOnceRead: passes.length }).toEqual({ + freed: N, + passesWhileHeld: 0, + passesOnceRead: 1, + }); expect(await proc.exited).toBe(0); }, ); From 985938ded6066bf95b3b13f9f046a05a1ad52c7f Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Wed, 23 Sep 2026 21:09:14 -0700 Subject: [PATCH 09/10] http: end a body on the HTTP thread when its undecoded rest is short A held body is the undecoded rest of a response whose transport is done. Its consumer decodes it in work-pool jobs, and the first job starts pool threads. That is not worth it for a short rest. - With the last of a body, the HTTP thread decodes up to 64 KB more than the budget. A body that ends inside that leaves no held body. - A paused consumer gets the same 64 KB with the last of a body. - A zlib pass over a held body is tagged in the debug log, so a test can tell it from a pass on the HTTP thread. --- src/http/InternalState.rs | 2 +- src/http/lib.rs | 30 +++++++++++----- test/js/web/fetch/fetch-backpressure.test.ts | 38 ++++++++++++++++++++ 3 files changed, 60 insertions(+), 10 deletions(-) diff --git a/src/http/InternalState.rs b/src/http/InternalState.rs index e33fe6ad8d1a..67b3b212b8fc 100644 --- a/src/http/InternalState.rs +++ b/src/http/InternalState.rs @@ -526,7 +526,7 @@ impl HeldBody { return Ok(true); } let input = &self.input[self.consumed..]; - log!("Decompressing {} bytes\n", input.len()); + log!("Decompressing {} bytes of a held body\n", input.len()); self.consumed += self.decompressor .decompress_chunk(self.encoding, input, out, max_output, true)?; diff --git a/src/http/lib.rs b/src/http/lib.rs index c8716601c7eb..da7ea630abbe 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -4119,13 +4119,22 @@ impl<'a> HTTPClient<'a> { /// Output budget of one decode pass. None for the last of a body that has no `HeldBody` taker. #[inline] fn decompress_output_cap(&self, is_final_chunk: bool) -> usize { - if self.decodes_on_demand() && (self.flags.takes_held_body || !is_final_chunk) { - signals::BODY_HIGH_WATER_MARK + if !self.decodes_on_demand() { + return usize::MAX; + } + if !is_final_chunk { + return signals::BODY_HIGH_WATER_MARK; + } + if self.flags.takes_held_body { + signals::BODY_HIGH_WATER_MARK + Self::HELD_BODY_MIN } else { usize::MAX } } + /// A rest that decodes to less than this ends on the HTTP thread: it is not worth a work-pool job. + const HELD_BODY_MIN: usize = 64 * 1024; + /// Decodes what has arrived under the consumer's budget. Returns whether to report bytes. fn process_received_body(&mut self, is_final_chunk: bool) -> crate::Result { let mut max_output = self.decompress_output_cap(is_final_chunk); @@ -4136,19 +4145,22 @@ impl<'a> HTTPClient<'a> { return Ok(false); } if max_output != usize::MAX { - // Nothing is decoded for a paused consumer (a tunnelled socket keeps reading anyway). - if self.signals.is_receive_paused() { - self.state.flags.decompress_output_pending = true; - return Ok(false); - } - // A body that one libdeflate call can inflate goes whole to its consumer. + let paused = self.signals.is_receive_paused(); if is_final_chunk && self.state.wants_exact_size_inflate() { - if self.signals.is_body_unclaimed() { + // A body that one libdeflate call can inflate goes whole to its consumer. + if paused || self.signals.is_body_unclaimed() { self.state.flags.decompress_output_pending = true; return Ok(false); } // A consumer attached after the cap was read. max_output = self.decompress_output_cap(is_final_chunk); + } else if paused { + // Nothing is decoded for a paused consumer (a tunnelled socket keeps reading anyway). + if !is_final_chunk { + self.state.flags.decompress_output_pending = true; + return Ok(false); + } + max_output = Self::HELD_BODY_MIN; } } } diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index b16c3ded59fc..9ccda868289b 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -1035,6 +1035,44 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }); }); + // A pass over a held body is a work-pool job, and it logs "... of a held body". A rest that + // is short is not worth one: the HTTP thread decodes it with the last of the body. + test.skipIf(!isDebug).each([ + ["a short rest ends on the HTTP thread", MARK + 32 * 1024, false], + ["a long rest goes to its consumer", 4 * MARK, true], + ])("%s", async (_, size, held) => { + const body = brotliCompressSync(Buffer.alloc(size, "alpha beta gamma delta lorem ipsum ")); + const reply = Buffer.concat([ + Buffer.from(`HTTP/1.1 200 OK\r\nContent-Encoding: br\r\nContent-Length: ${body.length}\r\n\r\n`), + body, + ]); + const srv = createTcpServer(socket => { + socket.on("error", () => {}); + socket.on("data", () => socket.write(reply)); + }); + await using server = await listening(srv); + const script = /* js */ ` + const res = await fetch(${JSON.stringify(`http://127.0.0.1:${server.port}/`)}); + let total = 0; + for await (const chunk of res.body) total += chunk.byteLength; + console.log("RESULT", total); + `; + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", script], + env: { ...clientEnv, BUN_DEBUG_HTTPInternalState: "1" }, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + // Output.scoped writes to whichever stream it chose at init; scan both. + const output = stdout + stderr; + expect({ + result: output.match(/RESULT .*/g), + held: /Decompressing \d+ bytes of a held body/.test(output), + }).toEqual({ result: [`RESULT ${size}`], held }); + expect(exitCode).toBe(0); + }); + // A gzip body that arrives whole, and whose trailer says it inflates to between 512 KB and // 32 MB, is worth one exact-size libdeflate call instead of zlib passes. A reader's budget rules // that call out, and such a body is usually here before anyone has said how the Response is From 7d93c6e0e826f7b1ebd0a7baeb1b5c668b03cb39 Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Thu, 24 Sep 2026 13:30:47 -0700 Subject: [PATCH 10/10] http: update comments that described the hold in the HTTP client Three comments still said that the HTTP client keeps a complete body and that the consumer resumes it. The Response keeps the body now. No code change. --- src/http/lib.rs | 2 +- test/js/web/fetch/fetch-backpressure.test.ts | 6 ++---- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/src/http/lib.rs b/src/http/lib.rs index da7ea630abbe..4fc1f5698c0f 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -4116,7 +4116,7 @@ impl<'a> HTTPClient<'a> { self.flags.protocol == Protocol::Http1_1 && self.signals.is_demand_driven() } - /// Output budget of one decode pass. None for the last of a body that has no `HeldBody` taker. + /// Output budget of one decode pass. The last pass has `HELD_BODY_MIN` more, or no budget without a `HeldBody` taker. #[inline] fn decompress_output_cap(&self, is_final_chunk: bool) -> usize { if !self.decodes_on_demand() { diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 9ccda868289b..476fa9da9c3c 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -1189,7 +1189,7 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not expect(await consume(res)).toBe(digest); }); - // A tunnelled socket is never paused. The consumer's resume reaches the client all the same. + // The same through a tunnel: the Response has the body, so the consumer needs no socket. test.each([ ["res.bytes()", `[await (await fetch(url, opts)).bytes()]`], ["a streaming reader", `await Array.fromAsync((await fetch(url, opts)).body)`], @@ -1233,9 +1233,7 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not expect(await consume(res)).toBe(digest); }); - // A tunnelled socket is never paused, so the close reaches the client while it still holds - // the body. `server.closed` settles once the client has closed its side, which is after it - // handled the end of the body. Only then does the consumer attach. + // `server.closed` settles once the client has closed its side. Only then does the consumer attach. test.each([ ["res.bytes()", `[await res.bytes()]`], ["a streaming reader", `await Array.fromAsync(res.body)`],