Skip to content
Merged
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
4 changes: 3 additions & 1 deletion src/http/HTTPThread.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<bun_libdeflate_sys::libdeflate::OwnedDecompressor>,
pub(crate) compressor: Option<bun_libdeflate_sys::libdeflate::OwnedCompressor>,
pub(crate) shared_buffer: [u8; 512 * 1024],
pub(crate) shared_buffer: [u8; LIBDEFLATE_SHARED_BUFFER_LEN],
}

// SAFETY: `Option<Owned{De,}Compressor>` is `#[repr(transparent)]` over
Expand Down
49 changes: 32 additions & 17 deletions src/http/InternalState.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,19 @@ use crate::{CertificateInfo, Decompressor, Encoding, HTTPRequestBody, HTTPRespon

bun_core::define_scoped_log!(log, HTTPInternalState, hidden);

/// 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;

/// 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<usize> {
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

Expand Down Expand Up @@ -219,6 +232,20 @@ impl<'a> InternalState<'a> {
self.flags.decompress_output_pending
}

/// 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
&& !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).
Expand Down Expand Up @@ -285,36 +312,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;
Expand Down
57 changes: 53 additions & 4 deletions src/http/Signals.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,15 @@ 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`.
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Debug)]
pub enum BodyReceiveMode {
Flowing = 0,
Paused = 1,
BufferAll = 2,
Abandoned = 3,
Unclaimed = 4,
}

impl BodyReceiveMode {
Expand All @@ -37,6 +39,7 @@ impl BodyReceiveMode {
1 => Self::Paused,
2 => Self::BufferAll,
3 => Self::Abandoned,
4 => Self::Unclaimed,
_ => Self::Flowing,
}
}
Expand Down Expand Up @@ -85,18 +88,34 @@ 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
.map(bun_ptr::BackRef::from)
.is_some_and(|a| {
matches!(
BodyReceiveMode::from_u8(a.load(Ordering::Acquire)),
BodyReceiveMode::Flowing | BodyReceiveMode::Paused
BodyReceiveMode::Flowing | BodyReceiveMode::Paused | BodyReceiveMode::Unclaimed
)
})
}

/// `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
.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 {
Expand All @@ -120,6 +139,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)),
Expand Down Expand Up @@ -149,10 +176,32 @@ 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 resumes.
#[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
Expand Down
32 changes: 20 additions & 12 deletions src/http/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4117,14 +4117,22 @@ 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<bool> {
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.
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);
Expand Down Expand Up @@ -4700,12 +4708,12 @@ 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;
// 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;
}
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
Expand Down
4 changes: 2 additions & 2 deletions src/runtime/webcore/fetch/FetchTasklet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down Expand Up @@ -1985,7 +1985,7 @@ impl FetchTasklet {
abort_handle: jsc::AbortHandle::for_owner::<FetchTasklet>(),
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,
Expand Down
Loading
Loading