From ca373fb1952dfbf91eb228001f37eeb047e00ad7 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 22:30:16 +0000 Subject: [PATCH 1/4] fetch: restore the one-shot gzip inflate for text() and arrayBuffer() A gzip response body that arrives whole usually does so before the caller has the Response. The receive mode is still Flowing then, so the 256 KB output budget from #43123 sends the body through zlib passes. Before that budget, one exact-size libdeflate call inflated it. FetchTasklet now starts in a new receive mode, Unclaimed: no consumer has attached yet. When a complete gzip body whose trailer size is above the 512 KB shared buffer and below 32 MB has not been touched by a decoder, the HTTP thread moves Unclaimed to Paused and decodes nothing. The consumer that attaches resumes the transport. A buffered consumer (text, arrayBuffer) then gets the one libdeflate call. A reader gets budgeted zlib passes as before. Nothing is decoded for a Response that nobody reads. The tests cover a buffered consumer and a reader, with the body arriving with the head, after the head, and after the consumer, over TLS and through a CONNECT proxy, for a Response nobody reads, and for a collected Response. The CONNECT proxy cases clear NO_PROXY and the proxy variables for the child and assert that the proxy saw one CONNECT. --- src/http/HTTPThread.rs | 4 +- src/http/InternalState.rs | 52 ++++-- src/http/Signals.rs | 63 ++++++- src/http/lib.rs | 34 ++-- src/runtime/webcore/fetch/FetchTasklet.rs | 4 +- test/js/web/fetch/fetch-backpressure.test.ts | 179 ++++++++++++++++++- 6 files changed, 298 insertions(+), 38 deletions(-) diff --git a/src/http/HTTPThread.rs b/src/http/HTTPThread.rs index 7c61110e0b83..b4f2c2f208df 100644 --- a/src/http/HTTPThread.rs +++ b/src/http/HTTPThread.rs @@ -225,10 +225,12 @@ pub struct CertCheckResumeMessage { pub(crate) async_http_id: u32, } +pub(crate) const LIBDEFLATE_SHARED_BUFFER_LEN: usize = 512 * 1024; + pub struct LibdeflateState { pub(crate) decompressor: Option, pub(crate) compressor: Option, - pub(crate) shared_buffer: [u8; 512 * 1024], + pub(crate) shared_buffer: [u8; LIBDEFLATE_SHARED_BUFFER_LEN], } // SAFETY: `Option` is `#[repr(transparent)]` over diff --git a/src/http/InternalState.rs b/src/http/InternalState.rs index 1a083f5e38d2..907f3b036116 100644 --- a/src/http/InternalState.rs +++ b/src/http/InternalState.rs @@ -6,6 +6,21 @@ use crate::{CertificateInfo, Decompressor, Encoding, HTTPRequestBody, HTTPRespon bun_core::define_scoped_log!(log, HTTPInternalState, hidden); +/// A gzip trailer is arbitrary input from the internet, so this bounds the allocation it can +/// ask libdeflate's exact-size call for. +const EXACT_SIZE_INFLATE_MAX: usize = 32 * 1024 * 1024; + +/// gzip stores the size of the uncompressed data in the last 4 bytes of the stream. It is only +/// valid if the stream is less than 4 GB, and it counts the last member alone. +fn gzip_trailer_size(buffer: &[u8]) -> Option { + if buffer.len() <= 16 || buffer.len() >= 1024 * 1024 * 1024 { + return None; + } + buffer + .last_chunk::<4>() + .map(|size| u32::from_le_bytes(*size) as usize) +} + // TODO: reduce the size of this struct // Many of these fields can be moved to a packed struct and use less space @@ -219,6 +234,21 @@ impl<'a> InternalState<'a> { self.flags.decompress_output_pending } + /// A complete gzip body that no decoder has touched, of a size that libdeflate inflates in + /// one exact-size call (`decompress_bytes`). An output budget rules that call out. + pub(crate) fn wants_exact_size_inflate(&self) -> bool { + bun_core::feature_flags::is_libdeflate_enabled() + && self.encoding == Encoding::Gzip + && !self.flags.is_libdeflate_fast_path_disabled + && !self.flags.is_redirect_pending + && matches!(self.decompressor, Decompressor::None) + && self.is_done() + && gzip_trailer_size(&self.compressed_body.list).is_some_and(|size| { + size > crate::http_thread::LIBDEFLATE_SHARED_BUFFER_LEN + && size < EXACT_SIZE_INFLATE_MAX + }) + } + /// True when a socket close during `in_progress` completes the body rather /// than failing it: chunked decoder already in the trailers state, or a /// close-delimited response (no Content-Length, no Transfer-Encoding). @@ -285,36 +315,24 @@ impl<'a> InternalState<'a> { log!("Decompressing {} bytes with libdeflate\n", buffer.len()); let deflater = crate::http_thread().deflater(); - // gzip stores the size of the uncompressed data in the last 4 bytes of the stream - // But it's only valid if the stream is less than 4.7 GB, since it's 4 bytes. // If we know that the stream is going to be larger than our // pre-allocated buffer, then let's dynamically allocate the exact // size. if self.encoding == Encoding::Gzip - && buffer.len() > 16 - && buffer.len() < 1024 * 1024 * 1024 + && let Some(estimated_size) = gzip_trailer_size(buffer) + && estimated_size > deflater.shared_buffer.len() { - let estimated_size: u32 = u32::from_le_bytes( - buffer[buffer.len() - 4..][..4] - .try_into() - .expect("infallible: size matches"), - ); // Under an output budget only `shared_buffer`'s worth may come out in one shot. - if (estimated_size as usize) > deflater.shared_buffer.len() - && max_output != usize::MAX - { + if max_output != usize::MAX { break 'libdeflate; } - // Since this is arbtirary input from the internet, let's set an upper bound of 32 MB for the allocation size. - if (estimated_size as usize) > deflater.shared_buffer.len() - && estimated_size < 32 * 1024 * 1024 - { + if estimated_size < EXACT_SIZE_INFLATE_MAX { self.decoded_body.list.clear(); // A trailer can lie; the streaming path below allocates only what is really there. if self .decoded_body .list - .try_reserve_exact(estimated_size as usize) + .try_reserve_exact(estimated_size) .is_err() { break 'libdeflate; diff --git a/src/http/Signals.rs b/src/http/Signals.rs index 868a028bb1f0..5bd51496d94c 100644 --- a/src/http/Signals.rs +++ b/src/http/Signals.rs @@ -21,6 +21,9 @@ pub const BODY_HIGH_WATER_MARK: usize = 256 * 1024; /// moves `Paused -> Flowing` and schedules a resume. The transport applies `Paused` after the /// next read. Two terminal states: `BufferAll` (a consumer wants the whole body) and /// `Abandoned` (nothing will read it; the transport is being shut down, drop what arrives). +/// `Unclaimed` is `Flowing` while no consumer has attached yet, for a body whose consumer +/// attaches later (a `Response`). The transport may move it to `Paused` to keep a body whole +/// until the consumer says how it reads (`Signals::hold_for_consumer`). #[repr(u8)] #[derive(Copy, Clone, PartialEq, Eq, Debug)] pub enum BodyReceiveMode { @@ -28,6 +31,7 @@ pub enum BodyReceiveMode { Paused = 1, BufferAll = 2, Abandoned = 3, + Unclaimed = 4, } impl BodyReceiveMode { @@ -37,6 +41,7 @@ impl BodyReceiveMode { 1 => Self::Paused, 2 => Self::BufferAll, 3 => Self::Abandoned, + 4 => Self::Unclaimed, _ => Self::Flowing, } } @@ -85,7 +90,7 @@ impl Signals { .is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Paused as u8) } - /// `Flowing` or `Paused`: a consumer takes the body piece by piece. + /// `Flowing`, `Paused` or `Unclaimed`: a consumer takes the body piece by piece. #[inline] pub(crate) fn is_demand_driven(self) -> bool { self.body_receive_mode @@ -93,10 +98,29 @@ impl Signals { .is_some_and(|a| { matches!( BodyReceiveMode::from_u8(a.load(Ordering::Acquire)), - BodyReceiveMode::Flowing | BodyReceiveMode::Paused + BodyReceiveMode::Flowing + | BodyReceiveMode::Paused + | BodyReceiveMode::Unclaimed ) }) } + + /// `Unclaimed -> Paused`. Returns whether it held: no consumer has attached, and the one + /// that does resumes the transport (`Store::receive_all`, `Store::receive_on_demand`). + #[inline] + pub(crate) fn hold_for_consumer(self) -> bool { + self.body_receive_mode + .map(bun_ptr::BackRef::from) + .is_some_and(|a| { + a.compare_exchange( + BodyReceiveMode::Unclaimed as u8, + BodyReceiveMode::Paused as u8, + Ordering::AcqRel, + Ordering::Relaxed, + ) + .is_ok() + }) + } } pub struct Store { @@ -120,6 +144,14 @@ impl Default for Store { } impl Store { + /// For a body whose consumer attaches after the response head arrives: `Unclaimed`. + pub fn unclaimed() -> Self { + Self { + body_receive_mode: AtomicU8::new(BodyReceiveMode::Unclaimed as u8), + ..Self::default() + } + } + pub fn to(&mut self) -> Signals { Signals { header_progress: Some(NonNull::from(&self.header_progress)), @@ -149,10 +181,33 @@ impl Store { .is_ok() } - /// `Flowing -> Paused`. No-op in the other states. + /// `Flowing` or `Unclaimed -> Paused`. No-op in the other states. #[inline] pub fn pause_receive(&self) { - let _ = self.try_transition_receive_mode(BodyReceiveMode::Flowing, BodyReceiveMode::Paused); + let _ = self + .body_receive_mode + .try_update(Ordering::AcqRel, Ordering::Acquire, |mode| { + matches!( + BodyReceiveMode::from_u8(mode), + BodyReceiveMode::Flowing | BodyReceiveMode::Unclaimed + ) + .then_some(BodyReceiveMode::Paused as u8) + }); + } + + /// A streaming consumer attached: `Unclaimed` or `Paused -> Flowing`. The caller schedules + /// the transport's resume either way. + #[inline] + pub fn receive_on_demand(&self) { + let _ = self + .body_receive_mode + .try_update(Ordering::AcqRel, Ordering::Acquire, |mode| { + matches!( + BodyReceiveMode::from_u8(mode), + BodyReceiveMode::Unclaimed | BodyReceiveMode::Paused + ) + .then_some(BodyReceiveMode::Flowing as u8) + }); } /// `Paused -> Flowing`. Returns whether it was paused, i.e. whether the caller has to diff --git a/src/http/lib.rs b/src/http/lib.rs index ce86dead9135..bb1f5a8be669 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -4117,14 +4117,23 @@ 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); + let mut max_output = self.decompress_output_cap(); + if max_output != usize::MAX && self.state.encoding.is_compressed() { + // 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 waits whole for its consumer: + // `BufferAll` gets that call, a reader gets budgeted passes. + if is_final_chunk && self.state.wants_exact_size_inflate() { + if self.signals.hold_for_consumer() { + self.state.flags.decompress_output_pending = true; + return Ok(false); + } + // A consumer attached after the cap was read. + max_output = self.decompress_output_cap(); + } } // `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); @@ -4700,12 +4709,13 @@ impl<'a> HTTPClient<'a> { || self.signals.body_receive_mode.is_some(); if is_done || is_streaming || content_length.is_none() { let is_final_chunk = is_done; + // We can only use the libdeflate fast path when we are not streaming. A body that + // arrived whole keeps it: `process_received_body` may hold it for its consumer. + if !is_final_chunk { + self.state.flags.is_libdeflate_fast_path_disabled = true; + } let processed = self.process_received_body(is_final_chunk)?; - // We can only use the libdeflate fast path when we are not streaming - // If we ever call processBodyBuffer again, it cannot go through the fast path. - self.state.flags.is_libdeflate_fast_path_disabled = true; - let total_received = self.state.total_body_received; self.report_progress(total_received); // Close-delimited bodies still need per-packet decompression, but diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index ae77f31fa4bd..d62482222fbf 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -1707,7 +1707,7 @@ impl FetchTasklet { // between would otherwise reach the stream with its task finding the buffer empty, and // nothing left to undo that pause. Unconditional: also flushes body bytes the client // holds that arrived with no follow-up read (`drain_response_body`). - this.signal_store.unpause_receive(); + this.signal_store.receive_on_demand(); this.schedule_receive_resume(); if drained.is_empty() { @@ -1985,7 +1985,7 @@ impl FetchTasklet { abort_handle: jsc::AbortHandle::for_owner::(), context: cx.context().id(), signals: Signals::default(), - signal_store: http::signals::Store::default(), + signal_store: http::signals::Store::unclaimed(), has_schedule_callback: AtomicBool::new(false), abort_reason: StrongOptional::empty(), check_server_identity: fetch_options.check_server_identity, diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index f697b83759b3..66c768a46e2b 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -576,13 +576,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 +596,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 @@ -851,6 +853,179 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not }); }); }); + + // 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 + // read. So the client holds it undecoded until a consumer attaches: `.bytes()` gets the one + // call, a reader gets its passes. + describe("a gzip body that arrived whole waits for its consumer", () => { + const SIZE = 1024 * 1024 + 77; + const raw = Buffer.alloc(SIZE, "alpha beta gamma delta lorem ipsum "); + const digest = md5(raw); + // A few KB, so one packet carries it. + const wire = gzipSync(raw); + + type Framing = "content-length" | "chunked"; + // "with the head": one write, so the body is complete before the caller has a Response. + // "after the head": the body goes out when `sendBody()` says so. + type When = "with the head" | "after the head"; + async function serveWhole(framing: Framing, when: When, secure = false) { + const head = + framing === "chunked" + ? `HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nTransfer-Encoding: chunked\r\n\r\n` + : `HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: ${wire.length}\r\n\r\n`; + const body = + framing === "chunked" + ? Buffer.concat([Buffer.from(`${wire.length.toString(16)}\r\n`), wire, Buffer.from("\r\n0\r\n\r\n")]) + : wire; + const requested = Promise.withResolvers(); + const closed = Promise.withResolvers(); + const handler = (s: import("node:net").Socket) => { + s.on("error", () => {}); + s.on("close", () => closed.resolve()); + s.once("data", () => { + s.write(when === "with the head" ? Buffer.concat([Buffer.from(head), body]) : head); + requested.resolve(s); + }); + }; + const server = await listening(secure ? createTlsServer(tls, handler) : createTcpServer(handler)); + return { + ...server, + url: `${secure ? "https" : "http"}://127.0.0.1:${server.port}/`, + closed: closed.promise, + sendBody: async () => { + const socket = await requested.promise; + await new Promise(written => socket.write(body, () => written())); + }, + }; + } + + const consumers: [string, (res: Response) => Promise][] = [ + ["res.bytes()", async res => md5(await res.bytes())], + [ + "a streaming reader", + async res => { + const hasher = new Bun.CryptoHasher("md5"); + for await (const chunk of res.body!) hasher.update(chunk); + return hasher.digest("hex"); + }, + ], + ]; + + // The libdeflate call logs "Decompressing N bytes with libdeflate", and every zlib pass + // "Decompressing N bytes". Only a debug build has the log. + test.skipIf(!isDebug)("res.bytes() inflates it in one libdeflate call", async () => { + await using server = await serveWhole("content-length", "with the head"); + const script = /* js */ ` + const bytes = await (await fetch(${JSON.stringify(server.url)})).bytes(); + console.log("RESULT", bytes.byteLength, new Bun.CryptoHasher("md5").update(bytes).digest("hex")); + `; + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", script], + env: { ...bunEnv, 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), + passes: output.match(/Decompressing \d+ bytes.*/g), + }).toEqual({ + result: [`RESULT ${SIZE} ${digest}`], + passes: [`Decompressing ${wire.length} bytes with libdeflate`], + }); + expect(exitCode).toBe(0); + }); + + describe.each(["content-length", "chunked"] as Framing[])("%s", framing => { + test.each(consumers)("%s takes a body that came with the head", async (_, consume) => { + await using server = await serveWhole(framing, "with the head"); + expect(await consume(await fetch(server.url))).toBe(digest); + }); + + test.each(consumers)("%s takes a body that came after the head", async (_, consume) => { + await using server = await serveWhole(framing, "after the head"); + const res = await fetch(server.url); + await server.sendBody(); + expect(await consume(res)).toBe(digest); + }); + + // Nothing is held for a consumer that is already there. + test.each(consumers)("%s already waits when the body comes", async (_, consume) => { + await using server = await serveWhole(framing, "after the head"); + const consumed = consume(await fetch(server.url)); + await server.sendBody(); + expect(await consumed).toBe(digest); + }); + }); + + test.each(consumers)("%s takes a body that came with the head over TLS", async (_, consume) => { + await using server = await serveWhole("content-length", "with the head", true); + const res = await fetch(server.url, { tls: { rejectUnauthorized: false } }); + expect(await consume(res)).toBe(digest); + }); + + // A tunnelled socket is never paused. The consumer's resume reaches the client all the same. + test.each([ + ["res.bytes()", `[await (await fetch(url, opts)).bytes()]`], + ["a streaming reader", `await Array.fromAsync((await fetch(url, opts)).body)`], + ])("%s takes a body that came with the head through a CONNECT proxy", async (_, chunks) => { + await using server = await serveWhole("content-length", "with the head", true); + await using proxy = await serveConnectProxy(); + const opts = { proxy: proxy.url, tls: { rejectUnauthorized: false } }; + const script = /* js */ ` + const url = ${JSON.stringify(server.url)}, opts = ${JSON.stringify(opts)}; + const body = Buffer.concat(${chunks}); + const digest = new Bun.CryptoHasher("md5").update(body).digest("hex"); + process.stdout.write(JSON.stringify({ total: body.length, digest })); + `; + // An ambient NO_PROXY that lists 127.0.0.1 makes fetch() ignore its `proxy` option. + const env = { + ...bunEnv, + NO_PROXY: undefined, + no_proxy: undefined, + HTTP_PROXY: undefined, + http_proxy: undefined, + HTTPS_PROXY: undefined, + https_proxy: undefined, + ALL_PROXY: undefined, + all_proxy: undefined, + }; + await using proc = Bun.spawn({ cmd: [bunExe(), "-e", script], env, stdout: "pipe", stderr: "pipe" }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout, stderr, connects: proxy.connects() }).toEqual({ + stdout: JSON.stringify({ total: SIZE, digest }), + stderr: "", + connects: 1, + }); + expect(exitCode).toBe(0); + }); + + test("held and never read: the process is not held", async () => { + await using server = await serveWhole("content-length", "with the head"); + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", `globalThis.keep = await fetch(${JSON.stringify(server.url)});`], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout, stderr, exitCode }).toEqual({ stdout: "", stderr: "", exitCode: 0 }); + }); + + test("held, and its Response is collected: the fetch is aborted", async () => { + await using server = await serveWhole("content-length", "with the head"); + // Its own frame, so that nothing on this one still refers to the response afterwards. + async function abandon() { + expect((await fetch(server.url)).status).toBe(200); + } + await abandon(); + await collectUntil(server.closed); + }); + }); }); describe.concurrent("fetch() receive backpressure — buffered consumers are not throttled", () => { From 3ca187d4134aa6b2cf36fac70b642389a4bd39c0 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Fri, 18 Sep 2026 22:45:21 +0000 Subject: [PATCH 2/4] [autofix.ci] apply automated fixes --- src/http/Signals.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/http/Signals.rs b/src/http/Signals.rs index 5bd51496d94c..fe77b9ed0da3 100644 --- a/src/http/Signals.rs +++ b/src/http/Signals.rs @@ -98,9 +98,7 @@ impl Signals { .is_some_and(|a| { matches!( BodyReceiveMode::from_u8(a.load(Ordering::Acquire)), - BodyReceiveMode::Flowing - | BodyReceiveMode::Paused - | BodyReceiveMode::Unclaimed + BodyReceiveMode::Flowing | BodyReceiveMode::Paused | BodyReceiveMode::Unclaimed ) }) } From 8828f6a2153ff0f2f7a189a9fa0dcb944c1215b8 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 23:02:23 +0000 Subject: [PATCH 3/4] fetch: keep the new comments in the gzip hold to one line each No code changes. Each comment that spanned two or three lines now says only what the code cannot, in one line. --- src/http/InternalState.rs | 9 +++------ src/http/Signals.rs | 10 +++------- src/http/lib.rs | 6 ++---- 3 files changed, 8 insertions(+), 17 deletions(-) diff --git a/src/http/InternalState.rs b/src/http/InternalState.rs index 907f3b036116..041c611e109e 100644 --- a/src/http/InternalState.rs +++ b/src/http/InternalState.rs @@ -6,12 +6,10 @@ use crate::{CertificateInfo, Decompressor, Encoding, HTTPRequestBody, HTTPRespon bun_core::define_scoped_log!(log, HTTPInternalState, hidden); -/// A gzip trailer is arbitrary input from the internet, so this bounds the allocation it can -/// ask libdeflate's exact-size call for. +/// Bounds the allocation that an untrusted gzip trailer can ask libdeflate's exact-size call for. const EXACT_SIZE_INFLATE_MAX: usize = 32 * 1024 * 1024; -/// gzip stores the size of the uncompressed data in the last 4 bytes of the stream. It is only -/// valid if the stream is less than 4 GB, and it counts the last member alone. +/// ISIZE, the last 4 bytes of a gzip stream: the decoded size of its last member, modulo 4 GB. fn gzip_trailer_size(buffer: &[u8]) -> Option { if buffer.len() <= 16 || buffer.len() >= 1024 * 1024 * 1024 { return None; @@ -234,8 +232,7 @@ impl<'a> InternalState<'a> { self.flags.decompress_output_pending } - /// A complete gzip body that no decoder has touched, of a size that libdeflate inflates in - /// one exact-size call (`decompress_bytes`). An output budget rules that call out. + /// A complete gzip body that only an unbudgeted pass can inflate in one libdeflate call. pub(crate) fn wants_exact_size_inflate(&self) -> bool { bun_core::feature_flags::is_libdeflate_enabled() && self.encoding == Encoding::Gzip diff --git a/src/http/Signals.rs b/src/http/Signals.rs index fe77b9ed0da3..57a0696ef9e1 100644 --- a/src/http/Signals.rs +++ b/src/http/Signals.rs @@ -21,9 +21,7 @@ pub const BODY_HIGH_WATER_MARK: usize = 256 * 1024; /// moves `Paused -> Flowing` and schedules a resume. The transport applies `Paused` after the /// next read. Two terminal states: `BufferAll` (a consumer wants the whole body) and /// `Abandoned` (nothing will read it; the transport is being shut down, drop what arrives). -/// `Unclaimed` is `Flowing` while no consumer has attached yet, for a body whose consumer -/// attaches later (a `Response`). The transport may move it to `Paused` to keep a body whole -/// until the consumer says how it reads (`Signals::hold_for_consumer`). +/// `Unclaimed`: `Flowing` before a consumer attaches. See `Signals::hold_for_consumer`. #[repr(u8)] #[derive(Copy, Clone, PartialEq, Eq, Debug)] pub enum BodyReceiveMode { @@ -103,8 +101,7 @@ impl Signals { }) } - /// `Unclaimed -> Paused`. Returns whether it held: no consumer has attached, and the one - /// that does resumes the transport (`Store::receive_all`, `Store::receive_on_demand`). + /// `Unclaimed -> Paused`: a whole body waits for the consumer, which resumes the transport. #[inline] pub(crate) fn hold_for_consumer(self) -> bool { self.body_receive_mode @@ -193,8 +190,7 @@ impl Store { }); } - /// A streaming consumer attached: `Unclaimed` or `Paused -> Flowing`. The caller schedules - /// the transport's resume either way. + /// A streaming consumer attached: `Unclaimed` or `Paused -> Flowing`. The caller resumes. #[inline] pub fn receive_on_demand(&self) { let _ = self diff --git a/src/http/lib.rs b/src/http/lib.rs index bb1f5a8be669..08bfcb7c32f5 100644 --- a/src/http/lib.rs +++ b/src/http/lib.rs @@ -4124,8 +4124,7 @@ impl<'a> HTTPClient<'a> { self.state.flags.decompress_output_pending = true; return Ok(false); } - // A body that one libdeflate call can inflate waits whole for its consumer: - // `BufferAll` gets that call, a reader gets budgeted passes. + // A body that one libdeflate call can inflate waits whole for its consumer. if is_final_chunk && self.state.wants_exact_size_inflate() { if self.signals.hold_for_consumer() { self.state.flags.decompress_output_pending = true; @@ -4709,8 +4708,7 @@ impl<'a> HTTPClient<'a> { || self.signals.body_receive_mode.is_some(); if is_done || is_streaming || content_length.is_none() { let is_final_chunk = is_done; - // We can only use the libdeflate fast path when we are not streaming. A body that - // arrived whole keeps it: `process_received_body` may hold it for its consumer. + // A body that arrived whole keeps the libdeflate fast path: it may be held. if !is_final_chunk { self.state.flags.is_libdeflate_fast_path_disabled = true; } From e98f60c7ff61a925817cef41791d5f85b430c54e Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 18 Sep 2026 23:16:42 +0000 Subject: [PATCH 4/4] test: a held gzip body is delivered whole after the origin closes Four cases for a body that the client holds undecoded when the origin ends the connection before a consumer attaches. On a direct socket the client sees the close once the consumer resumes it. Through a CONNECT tunnel the close arrives while the body is held, and the test waits for the client to close its side before the consumer attaches. A buffered consumer and a reader get the exact bytes in both. --- test/js/web/fetch/fetch-backpressure.test.ts | 67 ++++++++++++++++++++ 1 file changed, 67 insertions(+) diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 66c768a46e2b..9ac7e31cabb3 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -898,6 +898,8 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not const socket = await requested.promise; await new Promise(written => socket.write(body, () => written())); }, + // What an origin's keep-alive timeout does: the connection ends after a complete response. + end: async () => void (await requested.promise).end(), }; } @@ -1004,6 +1006,71 @@ describe.concurrent("fetch() receive backpressure — the decompressor does not expect(exitCode).toBe(0); }); + // The socket of a held body is paused, so the client sees this close only once a consumer + // resumes it. + test.each(consumers)("%s takes a held body after the origin closed the connection", async (_, consume) => { + await using server = await serveWhole("content-length", "with the head"); + const res = await fetch(server.url); + await server.end(); + 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. + test.each([ + ["res.bytes()", `[await res.bytes()]`], + ["a streaming reader", `await Array.fromAsync(res.body)`], + ])("%s takes a held body after the origin closed the tunnel", async (_, chunks) => { + await using server = await serveWhole("content-length", "with the head", true); + await using proxy = await serveConnectProxy(); + const opts = { proxy: proxy.url, tls: { rejectUnauthorized: false } }; + const script = /* js */ ` + const res = await fetch(${JSON.stringify(server.url)}, ${JSON.stringify(opts)}); + process.stdout.write("held\\n"); + for await (const line of console) if (line === "go") break; + const body = Buffer.concat(${chunks}); + const digest = new Bun.CryptoHasher("md5").update(body).digest("hex"); + process.stdout.write(JSON.stringify({ total: body.length, digest }) + "\\n"); + `; + // An ambient NO_PROXY that lists 127.0.0.1 makes fetch() ignore its `proxy` option. + const env = { + ...bunEnv, + NO_PROXY: undefined, + no_proxy: undefined, + HTTP_PROXY: undefined, + http_proxy: undefined, + HTTPS_PROXY: undefined, + https_proxy: undefined, + ALL_PROXY: undefined, + all_proxy: undefined, + }; + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", script], + env, + stdin: "pipe", + stdout: "pipe", + stderr: "pipe", + }); + const stderr = proc.stderr.text(); + let result = ""; + for await (const line of forEachLine(proc.stdout)) { + if (line === "held") { + await server.end(); + await server.closed; + proc.stdin.write("go\n"); + proc.stdin.end(); + } else result = line; + } + const exitCode = await proc.exited; + expect({ result, stderr: await stderr, connects: proxy.connects() }).toEqual({ + result: JSON.stringify({ total: SIZE, digest }), + stderr: "", + connects: 1, + }); + expect(exitCode).toBe(0); + }); + test("held and never read: the process is not held", async () => { await using server = await serveWhole("content-length", "with the head"); await using proc = Bun.spawn({