Skip to content
26 changes: 18 additions & 8 deletions src/brotli/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,17 +101,23 @@ impl StreamingDecoder {
unsafe { self.brotli.as_mut() }
}

/// Consume all of `input`, appending decompressed bytes to `out`
/// (growing in 4096-byte steps). Returns `ShortRead` when more input is
/// required and `is_done` is false.
#[inline]
pub fn is_inflating(&self) -> bool {
matches!(self.state, ReaderState::Inflating)
}

/// Append decompressed bytes to `out` (growing in 4096-byte steps) until `input` is
/// consumed or `out.len()` reaches `max_output`. Returns the input bytes consumed.
/// Returns `ShortRead` when more input is required and `is_done` is false.
Comment thread
robobun marked this conversation as resolved.
pub fn decompress(
&mut self,
input: &[u8],
out: &mut Vec<u8>,
max_output: usize,
is_done: bool,
) -> crate::Result<()> {
) -> crate::Result<usize> {
if matches!(self.state, ReaderState::End | ReaderState::Error) {
return Ok(());
return Ok(input.len());
}
debug_assert!(out.as_ptr() != input.as_ptr());

Expand All @@ -120,12 +126,16 @@ impl StreamingDecoder {
self.state,
ReaderState::Uninitialized | ReaderState::Inflating
) {
if out.len() >= max_output {
return Ok(total_in);
}
if out.try_reserve(4096).is_err() {
self.state = ReaderState::Error;
return Err(crate::Error::OutOfMemory);
}
let budget = max_output - out.len();
let spare = out.spare_capacity_mut();
let out_len = spare.len();
let out_len = spare.len().min(budget);
let mut next_out: *mut u8 = spare.as_mut_ptr().cast::<u8>();

let next_in = &input[total_in..];
Expand Down Expand Up @@ -159,7 +169,7 @@ impl StreamingDecoder {
match result {
c::BrotliDecoderResult::success => {
self.state = ReaderState::End;
return Ok(());
return Ok(input.len());
}
c::BrotliDecoderResult::err => {
self.state = ReaderState::Error;
Expand Down Expand Up @@ -192,7 +202,7 @@ impl StreamingDecoder {
}
}
}
Ok(())
Ok(total_in)
}
}

Expand Down
29 changes: 21 additions & 8 deletions src/http/Decompressor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,16 @@ impl Decompressor {
// explicit `Drop` is unnecessary. Callers that want a mid-lifecycle reset
// assign `*self = Decompressor::None`.

/// Inside a stream: another `decompress_chunk` may produce output with no new input.
pub(crate) fn is_mid_stream(&self) -> bool {
match self {
Decompressor::Zlib(r) => r.is_inflating(),
Decompressor::Brotli(r) => r.is_inflating(),
Decompressor::Zstd(r) => r.is_inflating(),
Decompressor::None => false,
}
}

fn init(&mut self, encoding: Encoding, first_chunk: &[u8]) -> crate::Result<()> {
match encoding {
Encoding::Gzip | Encoding::Deflate => {
Expand Down Expand Up @@ -49,27 +59,30 @@ impl Decompressor {
}

/// Feed one body chunk `buffer` through the decoder, appending the
/// decompressed output to `body_out_str`. Creates the decoder on first
/// call. Returns `ShortRead` when more input is needed and the stream is
/// not yet done.
/// decompressed output to `body_out_str` 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.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn decompress_chunk(
&mut self,
encoding: Encoding,
buffer: &[u8],
body_out_str: &mut MutableString,
max_output: usize,
is_done: bool,
) -> crate::Result<()> {
) -> crate::Result<usize> {
if !encoding.is_compressed() {
return Ok(());
return Ok(buffer.len());
}
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, is_done)?),
Decompressor::Brotli(reader) => Ok(reader.decompress(buffer, out, is_done)?),
Decompressor::Zstd(reader) => Ok(reader.decompress(buffer, out, is_done)?),
Decompressor::Zlib(reader) => Ok(reader.decompress(buffer, out, max_output, is_done)?),
Decompressor::Brotli(reader) => {
Ok(reader.decompress(buffer, out, max_output, is_done)?)
}
Decompressor::Zstd(reader) => Ok(reader.decompress(buffer, out, max_output, is_done)?),
Decompressor::None => {
unreachable!("Invalid encoding. This code should not be reachable")
}
Expand Down
80 changes: 62 additions & 18 deletions src/http/InternalState.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,8 @@ pub struct InternalState<'a> {
/// (cap-bounded) after the callback returns.
pub(crate) decoded_body: MutableString,
pub(crate) compressed_body: MutableString,
/// Prefix of `compressed_body` the decoder has already taken.
compressed_body_consumed: usize,
pub(crate) content_length: Option<usize>,
pub(crate) total_body_received: usize,
// Self-borrow into `original_request_body.bytes`; `RawSlice` carries the
Expand Down Expand Up @@ -80,6 +82,8 @@ pub struct InternalStateFlags {
/// `reset()`/`init()` so each redirect/retry hop re-compresses from the
/// original uncompressed `original_request_body`.
pub(crate) body_compressed: bool,
/// Held input or buffered decoder output remains for `HTTPClient::drain_response_body`.
pub(crate) decompress_output_pending: bool,
}

impl InternalStateFlags {
Expand All @@ -95,6 +99,7 @@ impl InternalStateFlags {
is_waiting_for_cert_check: false,
receive_paused: false,
body_compressed: false,
decompress_output_pending: false,
}
}
}
Expand All @@ -113,6 +118,7 @@ impl Default for InternalState<'_> {
stage: Stage::Pending,
decoded_body: MutableString::init_empty(),
compressed_body: MutableString::init_empty(),
compressed_body_consumed: 0,
content_length: None,
total_body_received: 0,
request_body: bun_ptr::RawSlice::EMPTY,
Expand Down Expand Up @@ -208,10 +214,19 @@ impl<'a> InternalState<'a> {
self.flags.received_last_chunk
}

#[inline]
pub(crate) fn has_pending_compressed(&self) -> bool {
self.flags.decompress_output_pending
}

/// 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).
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;
}
Expand All @@ -227,27 +242,30 @@ impl<'a> InternalState<'a> {
pub(crate) fn finalize_body_on_eof(&mut self) -> Result<(), Error> {
self.flags.received_last_chunk = true;
let buffer_snap = core::mem::take(&mut self.get_body_buffer().list);
self.process_body_buffer(buffer_snap, true).map(drop)
self.process_body_buffer(buffer_snap, true, usize::MAX)
Comment thread
robobun marked this conversation as resolved.
.map(drop)
}

pub(crate) fn decompress_bytes(
&mut self,
buffer: &[u8],
is_final_chunk: bool,
) -> Result<(), Error> {
max_output: usize,
) -> Result<usize, Error> {
// A response that declared a Content-Encoding but sent zero body bytes
// (e.g. an empty chunked gzip response) has nothing to decompress.
// Running the decompressor anyway makes it report a truncated stream
// (ZlibError); Node treats this as an empty body.
if buffer.is_empty() && self.total_body_received == 0 {
self.compressed_body.reset();
return Ok(());
return Ok(0);
}

// `self.compressed_body.reset()` must run on every exit. scopeguard would
// hold &mut self.compressed_body across the body and conflict with &mut self.decompressor,
// so each early-return below calls it explicitly.
let mut still_needs_to_decompress = true;
let mut consumed = buffer.len();

if bun_core::feature_flags::is_libdeflate_enabled() {
// Fast-path: use libdeflate
Expand All @@ -256,6 +274,7 @@ impl<'a> InternalState<'a> {
use bun_libdeflate_sys::libdeflate as bun_libdeflate;
if !(is_final_chunk
&& !self.flags.is_libdeflate_fast_path_disabled
&& matches!(self.decompressor, Decompressor::None)
&& self.encoding.can_use_lib_deflate()
&& self.is_done())
{
Expand All @@ -280,6 +299,12 @@ impl<'a> InternalState<'a> {
.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
{
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
Expand Down Expand Up @@ -356,33 +381,40 @@ impl<'a> InternalState<'a> {
let min = ((buffer.len() as f64) * 1.5)
.ceil()
.min(1024.0 * 1024.0 * 2.0);
if let Err(err) = self.decoded_body.grow_by((min as usize).max(32)) {
if let Err(err) = self
.decoded_body
.grow_by((min as usize).max(32).min(max_output))
{
self.compressed_body.reset();
return Err(err.into());
}
}

let is_done = self.is_done();
if let Err(err) = self.decompressor.decompress_chunk(
match self.decompressor.decompress_chunk(
self.encoding,
buffer,
&mut self.decoded_body,
max_output,
is_done,
) {
if is_done || err != crate::Error::ShortRead {
bun_core::pretty_errorln!(
"<r><red>Decompression error: {}<r>",
bstr::BStr::new(err.name()),
);
Output::flush();
self.compressed_body.reset();
return Err(err);
Ok(n) => consumed = n,
Err(err) => {
if is_done || err != crate::Error::ShortRead {
bun_core::pretty_errorln!(
"<r><red>Decompression error: {}<r>",
bstr::BStr::new(err.name()),
);
Output::flush();
self.compressed_body.reset();
return Err(err);
}
}
}
}

self.compressed_body.reset();
Ok(())
Ok(consumed)
}

// `buffer` is always the current body buffer's bytes. To avoid aliased &mut/& under
Expand All @@ -393,6 +425,7 @@ impl<'a> InternalState<'a> {
&mut self,
mut buffer: Vec<u8>,
is_final_chunk: bool,
max_output: usize,
) -> Result<bool, Error> {
if self.flags.is_redirect_pending {
// Caller moved the bytes out of the body buffer; put them back so the
Expand All @@ -403,10 +436,21 @@ impl<'a> InternalState<'a> {

match self.encoding {
Encoding::Brotli | Encoding::Gzip | Encoding::Deflate | Encoding::Zstd => {
self.decompress_bytes(&buffer, is_final_chunk)?;
// Retain capacity by
// returning the (cleared) allocation to compressed_body instead of dropping it.
buffer.clear();
let start = self.compressed_body_consumed;
let consumed =
start + self.decompress_bytes(&buffer[start..], is_final_chunk, max_output)?;
let held = buffer.len() - consumed;
// A decoder can hold output with no input left (brotli copy command, zstd flush).
self.flags.decompress_output_pending = self.decoded_body.list.len() >= max_output
&& (held != 0 || self.decompressor.is_mid_stream());
// Shifting only once the taken prefix is the larger part moves each byte once.
if consumed >= held {
buffer.drain(..consumed);
Comment thread
robobun marked this conversation as resolved.
self.compressed_body_consumed = 0;
} else {
self.compressed_body_consumed = consumed;
}
// Retain capacity by returning the allocation to compressed_body.
self.compressed_body.list = buffer;
}
_ => {
Expand Down
13 changes: 13 additions & 0 deletions src/http/Signals.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,19 @@ impl Signals {
.map(bun_ptr::BackRef::from)
.is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Paused as u8)
}

/// `Flowing` or `Paused`: a consumer takes the body piece by piece.
#[inline]
pub(crate) fn is_demand_driven(self) -> bool {
self.body_receive_mode
.map(bun_ptr::BackRef::from)
.is_some_and(|a| {
matches!(
BodyReceiveMode::from_u8(a.load(Ordering::Acquire)),
BodyReceiveMode::Flowing | BodyReceiveMode::Paused
)
})
}
}

pub struct Store {
Expand Down
Loading
Loading