Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions src/http/AsyncHTTP.rs
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,8 @@ pub struct Options<'a> {
pub verbose: Option<HTTPVerboseLevel>,
pub disable_keepalive: Option<bool>,
pub disable_decompression: Option<bool>,
/// The consumer takes `HTTPClientResult::held_body`. Others get it decoded with the last chunk.
pub takes_held_body: bool,
pub max_redirects: Option<u8>,
pub reject_unauthorized: Option<bool>,
pub tls_props: Option<SSLConfigSharedPtr>,
Expand Down Expand Up @@ -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;
}
Expand Down
6 changes: 2 additions & 4 deletions src/http/Decompressor.rs
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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<u8>,
max_output: usize,
is_done: bool,
) -> crate::Result<usize> {
Expand All @@ -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) => {
Expand Down
128 changes: 100 additions & 28 deletions src/http/InternalState.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,28 @@ fn gzip_trailer_size(buffer: &[u8]) -> Option<usize> {
.map(|size| u32::from_le_bytes(*size) as usize)
}

/// One libdeflate call over a whole gzip body into `size` bytes. False: zlib decodes it.
fn inflate_exact_size(
decompressor: &mut bun_libdeflate_sys::libdeflate::Decompressor,
buffer: &[u8],
size: usize,
out: &mut Vec<u8>,
) -> bool {
use bun_libdeflate_sys::libdeflate::{Encoding, Status};
out.clear();
// A trailer can lie; the streaming path allocates only what is really there.
if out.try_reserve_exact(size).is_err() {
return false;
}
let result = decompressor.decompress_to_vec(buffer, out, Encoding::Gzip);
// Unconsumed input means a multi-member stream (RFC 1952 §2.2), which libdeflate does not decode.
if result.status == Status::Success && result.read == buffer.len() {
return true;
}
out.clear();
false
}

// TODO: reduce the size of this struct
// Many of these fields can be moved to a packed struct and use less space

Expand Down Expand Up @@ -250,16 +272,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
Expand Down Expand Up @@ -324,32 +360,12 @@ impl<'a> InternalState<'a> {
break 'libdeflate;
}
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)
.is_err()
{
break 'libdeflate;
}
let result = deflater.decompressor_mut().decompress_to_vec(
still_needs_to_decompress = !inflate_exact_size(
deflater.decompressor_mut(),
buffer,
estimated_size,
&mut self.decoded_body.list,
bun_libdeflate::Encoding::Gzip,
);
// libdeflate decodes a single gzip member; unconsumed
// input means this is a multi-member stream (RFC 1952
// §2.2). Let the zlib path handle it.
if result.status == bun_libdeflate::Status::Success
&& result.read == buffer.len()
{
still_needs_to_decompress = false;
} else {
self.decoded_body.list.clear();
}

break 'libdeflate;
}
}
Expand Down Expand Up @@ -409,7 +425,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,
) {
Expand Down Expand Up @@ -492,6 +508,62 @@ 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<u8>,
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<u8>, max_output: usize) -> Result<bool, Error> {
if max_output == usize::MAX && self.inflate_exact_size(out) {
return Ok(true);
}
let input = &self.input[self.consumed..];
log!("Decompressing {} bytes of a held body\n", input.len());
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)
}
}

impl HeldBody {
/// A gzip body that no decoder has started, for a consumer that takes all of it.
fn inflate_exact_size(&mut self, out: &mut Vec<u8>) -> bool {
if !bun_core::feature_flags::is_libdeflate_enabled()
|| self.encoding != Encoding::Gzip
|| !matches!(self.decompressor, Decompressor::None)
|| self.consumed != 0
|| !out.is_empty()
{
return false;
}
let Some(size) = gzip_trailer_size(&self.input).filter(|&n| n < EXACT_SIZE_INFLATE_MAX)
else {
return false;
};
let Some(mut decompressor) = bun_libdeflate_sys::libdeflate::OwnedDecompressor::new()
else {
return false;
};
log!("Decompressing {} bytes with libdeflate\n", self.input.len());
if !inflate_exact_size(&mut decompressor, &self.input, size, out) {
return false;
}
self.consumed = self.input.len();
true
}
}

#[derive(Clone, Copy, PartialEq, Eq)]
pub enum HTTPStage {
Pending,
Expand Down
2 changes: 1 addition & 1 deletion src/http/ProxyTunnel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -506,7 +506,7 @@ fn on_close(ctx: *mut HTTPClient) {
&& !this.state.flags.is_redirect_pending;
let mut fail_err: Option<crate::Error> = 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);
Expand Down
34 changes: 17 additions & 17 deletions src/http/Signals.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +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`: `Flowing` before a consumer attaches. See `Signals::hold_for_consumer`.
/// `Unclaimed`: `Flowing` before a consumer attaches. A body that completes in it stays undecoded.
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Debug)]
pub enum BodyReceiveMode {
Expand Down Expand Up @@ -88,32 +88,32 @@ impl Signals {
.is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Paused as u8)
}

/// `Flowing`, `Paused` or `Unclaimed`: a consumer takes the body piece by piece.
/// Nothing will read the body, and its consumer is shutting the transport down.
#[inline]
pub(crate) fn is_demand_driven(self) -> bool {
pub(crate) fn is_body_abandoned(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 | BodyReceiveMode::Unclaimed
)
})
.is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Abandoned as u8)
}

/// No consumer has attached to the body yet.
#[inline]
pub(crate) fn is_body_unclaimed(self) -> bool {
self.body_receive_mode
.map(bun_ptr::BackRef::from)
.is_some_and(|a| a.load(Ordering::Acquire) == BodyReceiveMode::Unclaimed as u8)
}

/// `Unclaimed -> Paused`: a whole body waits for the consumer, which resumes the transport.
/// `Flowing`, `Paused` or `Unclaimed`: a consumer takes the body piece by piece.
#[inline]
pub(crate) fn hold_for_consumer(self) -> bool {
pub(crate) fn is_demand_driven(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,
matches!(
BodyReceiveMode::from_u8(a.load(Ordering::Acquire)),
BodyReceiveMode::Flowing | BodyReceiveMode::Paused | BodyReceiveMode::Unclaimed
)
.is_ok()
})
}
}
Expand Down
Loading
Loading