diff --git a/packages/bun-types/bun.d.ts b/packages/bun-types/bun.d.ts index 44ca4f6af31f..52612e3b6bef 100644 --- a/packages/bun-types/bun.d.ts +++ b/packages/bun-types/bun.d.ts @@ -2134,7 +2134,7 @@ declare module "bun" { */ function write( destination: BunFile | S3File | PathLike, - input: Blob | NodeJS.TypedArray | ArrayBufferLike | string | BlobPart[] | Archive, + input: Blob | NodeJS.TypedArray | ArrayBufferLike | string | BlobPart[] | Archive | ReadableStream, options?: { /** * If writing to a PathLike, set the permissions of the file. @@ -2152,19 +2152,20 @@ declare module "bun" { ): Promise; /** - * Persist a {@link Response} body to disk. + * Persist a {@link Response} or {@link Request} body to disk. The body is + * streamed into the file as it arrives. * * @param destination The file to write to. If the file doesn't exist, it is * created; if it does, it is overwritten. If `input` is smaller than * `destination`, `destination` is truncated. - * @param input The `Response` whose body is written + * @param input The `Response` or `Request` whose body is written * @param options Options for the write * * @returns A promise that resolves with the number of bytes written. */ function write( destination: BunFile, - input: Response, + input: Response | Request, options?: { /** * If `true`, create the parent directory if it doesn't exist. @@ -2178,17 +2179,18 @@ declare module "bun" { ): Promise; /** - * Persist a {@link Response} body to disk. + * Persist a {@link Response} or {@link Request} body to disk. The body is + * streamed into the file as it arrives. * * @param destinationPath The file path to write to. If the file doesn't * exist, it is created; if it does, it is overwritten. If `input` is * smaller than the existing file, the file is truncated. - * @param input The `Response` whose body is written + * @param input The `Response` or `Request` whose body is written * @returns A promise that resolves with the number of bytes written. */ function write( destinationPath: PathLike, - input: Response, + input: Response | Request, options?: { /** * If `true`, create the parent directory if it doesn't exist. @@ -2740,7 +2742,7 @@ declare module "bun" { * @param options - The options to use for the write. */ write( - data: string | ArrayBufferView | ArrayBuffer | SharedArrayBuffer | Request | Response | BunFile, + data: string | ArrayBufferView | ArrayBuffer | SharedArrayBuffer | Request | Response | BunFile | ReadableStream, options?: { highWaterMark?: number }, ): Promise; diff --git a/src/http/Signals.rs b/src/http/Signals.rs index bf744ebd0cd8..1d343542a181 100644 --- a/src/http/Signals.rs +++ b/src/http/Signals.rs @@ -12,17 +12,22 @@ pub struct Signals { pub body_receive_mode: Option>, } +/// Receive backpressure high-water mark: bytes no consumer has taken, on either side of the +/// HTTP→JS hop. A body shorter than this completes unread, which frees its connection. +pub const BODY_HIGH_WATER_MARK: usize = 256 * 1024; + +/// Receive backpressure for a body handed to JS. Whichever side holds bytes no consumer has +/// taken moves `Flowing -> Paused` once they reach the high-water mark; whoever takes them +/// 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). #[repr(u8)] #[derive(Copy, Clone, PartialEq, Eq, Debug)] pub enum BodyReceiveMode { - /// Pause the transport after each delivered body chunk until JS pulls. - AutoPause = 0, - /// `callback` won the CAS; transport should be paused until JS pulls. + Flowing = 0, Paused = 1, - /// `.arrayBuffer()`/`.text()`/etc attached — never pause. BufferAll = 2, - /// Cancelled or abandoned — never pause, callback discards bytes. - Ignore = 3, + Abandoned = 3, } impl BodyReceiveMode { @@ -31,8 +36,8 @@ impl BodyReceiveMode { match v { 1 => Self::Paused, 2 => Self::BufferAll, - 3 => Self::Ignore, - _ => Self::AutoPause, + 3 => Self::Abandoned, + _ => Self::Flowing, } } } @@ -96,7 +101,7 @@ impl Default for Store { response_body_streaming: AtomicBool::new(false), aborted: AtomicBool::new(false), cert_errors: AtomicBool::new(false), - body_receive_mode: AtomicU8::new(BodyReceiveMode::AutoPause as u8), + body_receive_mode: AtomicU8::new(BodyReceiveMode::Flowing as u8), } } } @@ -125,21 +130,38 @@ impl Store { } #[inline] - pub fn try_transition_receive_mode(&self, from: BodyReceiveMode, to: BodyReceiveMode) -> bool { + fn try_transition_receive_mode(&self, from: BodyReceiveMode, to: BodyReceiveMode) -> bool { self.body_receive_mode .compare_exchange(from as u8, to as u8, Ordering::AcqRel, Ordering::Relaxed) .is_ok() } - /// Unconditionally move to a terminal mode (`BufferAll`/`Ignore`). - /// Returns whether the previous state was `Paused`. + /// `Flowing -> Paused`. No-op in the other states. + #[inline] + pub fn pause_receive(&self) { + let _ = self.try_transition_receive_mode(BodyReceiveMode::Flowing, BodyReceiveMode::Paused); + } + + /// `Paused -> Flowing`. Returns whether it was paused, i.e. whether the caller has to + /// schedule the transport's resume. + #[inline] + pub fn unpause_receive(&self) -> bool { + self.try_transition_receive_mode(BodyReceiveMode::Paused, BodyReceiveMode::Flowing) + } + + /// Terminal: never pause again. Returns whether it was paused. + #[inline] + pub fn receive_all(&self) -> bool { + self.body_receive_mode + .swap(BodyReceiveMode::BufferAll as u8, Ordering::AcqRel) + == BodyReceiveMode::Paused as u8 + } + + /// Terminal. #[inline] - pub fn set_receive_mode_terminal(&self, mode: BodyReceiveMode) -> bool { - debug_assert!(matches!( - mode, - BodyReceiveMode::BufferAll | BodyReceiveMode::Ignore - )); - self.body_receive_mode.swap(mode as u8, Ordering::AcqRel) == BodyReceiveMode::Paused as u8 + pub fn abandon(&self) { + self.body_receive_mode + .store(BodyReceiveMode::Abandoned as u8, Ordering::Release); } } diff --git a/src/http/h2_client/ClientSession.rs b/src/http/h2_client/ClientSession.rs index 3b6dc4a1d23d..78ded8082e10 100644 --- a/src/http/h2_client/ClientSession.rs +++ b/src/http/h2_client/ClientSession.rs @@ -683,8 +683,8 @@ impl ClientSession { self.by_http_id.get(&async_http_id).copied() } - /// JS just enabled `response_body_streaming` on the request, so flush any - /// body bytes that arrived between metadata delivery and `getReader()`. + /// A body consumer attached on the JS side: flush any body bytes that arrived between + /// metadata delivery and `getReader()`. fn drain_response_body(&mut self, async_http_id: u32) { let Some(stream) = self.stream_for_http_id(async_http_id) else { return; diff --git a/src/jsc/bindings/webcore/streams/BunStreamSource.cpp b/src/jsc/bindings/webcore/streams/BunStreamSource.cpp index 387dc7956ce7..c958585b118f 100644 --- a/src/jsc/bindings/webcore/streams/BunStreamSource.cpp +++ b/src/jsc/bindings/webcore/streams/BunStreamSource.cpp @@ -1017,6 +1017,10 @@ static std::optional rsisWriteChunk(JSC::VM& vm, JSGlobalObject* globalObj bool shouldSuspend = wrote.isNumber() && wrote.asNumber() < 0; if (auto* wrotePromise = dynamicDowncast(wrote)) { markPromiseAsHandled(vm, wrotePromise); + if (wrotePromise->status() == JSPromise::Status::Rejected) { + throwException(globalObject, scope, wrotePromise->result()); + return std::nullopt; + } shouldSuspend = wrotePromise->status() == JSPromise::Status::Pending; } if (shouldSuspend) { diff --git a/src/runtime/webcore/Blob.rs b/src/runtime/webcore/Blob.rs index 04debc966e8e..6cc2b5c1a206 100644 --- a/src/runtime/webcore/Blob.rs +++ b/src/runtime/webcore/Blob.rs @@ -234,7 +234,7 @@ pub trait BlobExt { &self, global_this: &JSGlobalObject, readable_stream: ReadableStream, - extra_options: Option, + options: &WriteFileOptions, ) -> JsResult; fn get_writer(&self, global_this: &JSGlobalObject, callframe: &CallFrame) -> JsResult; fn get_slice_from( @@ -1366,8 +1366,9 @@ impl BlobExt for Blob { &self, global_this: &JSGlobalObject, readable_stream: ReadableStream, - extra_options: Option, + options: &WriteFileOptions, ) -> JsResult { + let extra_options = options.extra_options; let Some(store) = self.store.get().clone() else { return Ok( JSPromise::dangerously_create_rejected_promise_value_without_notifying_vm( @@ -1447,11 +1448,23 @@ impl BlobExt for Blob { } else { let mut file_path = bun_paths::PathBuffer::uninit(); let path = pathlike.path().slice_z(&mut file_path); - match bun_sys::open( - path, - bun_sys::O::WRONLY | bun_sys::O::CREAT | bun_sys::O::NONBLOCK, - WRITE_PERMISSIONS, - ) { + let flags = bun_sys::O::WRONLY + | bun_sys::O::CREAT + | bun_sys::O::TRUNC + | bun_sys::O::NONBLOCK; + let mode = options.mode.unwrap_or(WRITE_PERMISSIONS); + let mut result = bun_sys::open(path, flags, mode); + if let bun_sys::Result::Err(err) = &result { + if err.get_errno() == bun_sys::E::ENOENT + && options.mkdirp_if_not_exists.unwrap_or(true) + { + result = match mkdirp_parent(path.as_bytes()) { + Ok(()) => bun_sys::open(path, flags, mode), + Err(err) => Err(err), + }; + } + } + match result { bun_sys::Result::Ok(result) => result, bun_sys::Result::Err(err) => { return Ok(JSPromise::dangerously_create_rejected_promise_value_without_notifying_vm( @@ -1555,6 +1568,9 @@ impl BlobExt for Blob { }; let stream_start = streams::Start::FileSink(streams::FileSinkOptions { + truncate: matches!(input_path, webcore::PathOrFileDescriptor::Path(_)), + mkdirp: options.mkdirp_if_not_exists.unwrap_or(true), + mode: options.mode.unwrap_or(WRITE_PERMISSIONS), input_path, ..Default::default() }); @@ -1574,8 +1590,13 @@ impl BlobExt for Blob { } }; - // Stay on the JS pump here: the native ByteStream path returns UNDEFINED before - // completion, which would resolve-0 early. + // SAFETY: file_sink is a live +1 *mut FileSink; `pipe_stream` takes its own refs. + if let Some(promise) = unsafe { (*file_sink).pipe_stream(&readable_stream, global_this) } { + // SAFETY: release our +1 ref on the sink. + unsafe { webcore::FileSink::deref(file_sink) }; + return Ok(promise); + } + let assignment_result: JSValue = webcore::file_sink::JSSink::assign_to_stream( global_this, readable_stream.value, @@ -1624,18 +1645,21 @@ impl BlobExt for Blob { return Ok(promise_value); } jsc::js_promise::Status::Fulfilled => { - // SAFETY: release our +1 ref on the sink. + // SAFETY: live until the deref below. + let written = unsafe { (*file_sink).stream_bytes.get().unwrap_or(0) }; + // SAFETY: release our +1 ref on the sink; not used after this. unsafe { webcore::FileSink::deref(file_sink) }; readable_stream.done(); return Ok(JSPromise::resolved_promise_value( global_this, - JSValue::js_number(0.0), + JSValue::js_number(written as f64), )); } jsc::js_promise::Status::Rejected => { // SAFETY: release our +1 ref on the sink. unsafe { webcore::FileSink::deref(file_sink) }; readable_stream.cancel(global_this)?; + promise.set_handled(global_this.vm()); return Ok(JSPromise::dangerously_create_rejected_promise_value_without_notifying_vm( global_this, promise.result(global_this.vm()), @@ -1654,12 +1678,14 @@ impl BlobExt for Blob { ); } } - // SAFETY: release our +1 ref on the sink. + // SAFETY: live until the deref below. + let written = unsafe { (*file_sink).stream_bytes.get().unwrap_or(0) }; + // SAFETY: release our +1 ref on the sink; not used after this. unsafe { webcore::FileSink::deref(file_sink) }; Ok(JSPromise::resolved_promise_value( global_this, - JSValue::js_number(0.0), + JSValue::js_number(written as f64), )) } @@ -4200,8 +4226,26 @@ pub enum Retry { No, } +/// Create the parent directories of `path`. +#[inline(never)] +pub(crate) fn mkdirp_parent(path: &[u8]) -> bun_sys::Result<()> { + let Some(dirname) = bun_core::dirname(path) else { + return Err(bun_sys::Error::from_code( + bun_sys::E::ENOENT, + bun_sys::Tag::mkdir, + )); + }; + node::fs::NodeFS::default() + .mkdir_recursive(&node::fs::args::Mkdir { + path: PathLike::borrowed(dirname), + recursive: true, + always_return_none: true, + ..Default::default() + }) + .map(|_| ()) +} + // TODO: move this to bun_sys? -// we choose not to inline this so that the path buffer is not on the stack unless necessary. #[inline(never)] pub(crate) fn mkdir_if_not_exists( this: &mut T, @@ -4210,24 +4254,16 @@ pub(crate) fn mkdir_if_not_exists( err_path: &[u8], ) -> Retry { if err.get_errno() == bun_sys::E::ENOENT && this.mkdirp_if_not_exists() { - if let Some(dirname) = bun_core::dirname(path_string.as_bytes()) { - let mut node_fs = node::fs::NodeFS::default(); - match node_fs.mkdir_recursive(&node::fs::args::Mkdir { - path: PathLike::borrowed(dirname), - recursive: true, - always_return_none: true, - ..Default::default() - }) { - bun_sys::Result::Ok(_) => { - this.set_mkdirp_if_not_exists(false); - return Retry::Continue; - } - bun_sys::Result::Err(err2) => { - this.set_errno_if_present(bun_errno::from_errno(err2.errno as i32).into()); - this.set_system_error(err.with_path(err_path).to_system_error()); - this.set_opened_fd_if_present(Fd::INVALID); - return Retry::Fail; - } + match mkdirp_parent(path_string.as_bytes()) { + bun_sys::Result::Ok(()) => { + this.set_mkdirp_if_not_exists(false); + return Retry::Continue; + } + bun_sys::Result::Err(err2) => { + this.set_errno_if_present(bun_errno::from_errno(err2.errno as i32).into()); + this.set_system_error(err.with_path(err_path).to_system_error()); + this.set_opened_fd_if_present(Fd::INVALID); + return Retry::Fail; } } } @@ -4263,6 +4299,18 @@ pub trait MkdirpTarget { // writeFileWithEmptySourceToDestination / writeFileWithSourceDestination // ────────────────────────────────────────────────────────────────────────── +fn body_used_rejection(global: &JSGlobalObject) -> JSValue { + JSPromise::dangerously_create_rejected_promise_value_without_notifying_vm( + global, + global + .err( + jsc::ErrorCode::BODY_ALREADY_USED, + format_args!("Body already used"), + ) + .to_js(), + ) +} + #[derive(Default, Clone, Copy)] pub struct WriteFileOptions { pub(crate) mkdirp_if_not_exists: Option, @@ -4580,11 +4628,7 @@ pub(crate) fn write_file_with_source_destination( )?, ctx, )? { - return destination_blob.pipe_readable_stream_to_blob( - ctx, - stream, - options.extra_options, - ); + return destination_blob.pipe_readable_stream_to_blob(ctx, stream, options); } else { return Ok( JSPromise::dangerously_create_rejected_promise_value_without_notifying_vm( @@ -4942,10 +4986,31 @@ pub(crate) fn write_file_internal( // `body_value` is `&mut Body::Value` from a live JS heap // Response/Request `m_ctx`, held raw so every borrow below is // scoped and none spans the JS-running calls in the arms. + // A stream someone holds a reader on, or has read from, is theirs. + let existing = get_stream().or_else(|| { + // SAFETY: scoped exclusive borrow; runs no JS. + match unsafe { &mut *body_value } { + BodyValue::Locked(locked) => locked.readable.get(), + _ => None, + } + }); + if let Some(readable) = existing { + if readable.is_locked(global_this) || readable.is_disturbed(global_this) { + destination_blob.detach(); + return Ok(ControlFlow::Break(body_used_rejection(global_this))); + } + } + // A body that is all here (also behind an untouched `.body` stream) is written as a blob. + // SAFETY: scoped exclusive borrow; runs no JS. + unsafe { (*body_value).to_blob_if_possible() }; // SAFETY: scoped shared read of the variant tag. let tag = match unsafe { &*body_value } { BodyValue::Error(_) => BodyTag::Error, BodyValue::Locked(_) => BodyTag::Locked, + BodyValue::Used => { + destination_blob.detach(); + return Ok(ControlFlow::Break(body_used_rejection(global_this))); + } _ => BodyTag::Use, }; match tag { @@ -5025,6 +5090,66 @@ pub(crate) fn write_file_internal( "ReadableStream has already been used" ))); } + // A body that is a stream, or that its producer can stream (fetch, the + // server, HTMLRewriter): pipe it into the file instead of collecting it in + // memory first. The stream also outlives the Response it came from. + let streamable = get_stream().is_some() || { + // SAFETY: scoped exclusive borrow; runs no JS. + let BodyValue::Locked(locked) = (unsafe { &mut *body_value }) else { + unreachable!() + }; + locked.readable.has() || locked.on_start_streaming.is_some() + }; + if streamable { + // SAFETY: exclusive borrow scoped to the call (may run JS). + let _ = unsafe { (*body_value).to_readable_stream(global_this) }?; + let readable = get_stream().or_else(|| { + // SAFETY: re-borrow after `to_readable_stream`. + let BodyValue::Locked(locked) = (unsafe { &mut *body_value }) else { + return None; + }; + locked.readable.get() + }); + // SAFETY: scoped; `to_readable_stream` may have replaced the value. + let body = unsafe { &*body_value }; + if let (Some(readable), BodyValue::Locked(_)) = (readable, body) { + let promise = destination_blob.pipe_readable_stream_to_blob( + global_this, + readable, + &options, + )?; + // The destination could not be opened: the stream was not touched and + // is still the body's. + let failed = promise.as_any_promise().is_some_and(|p| { + matches!(p.status(), jsc::js_promise::Status::Rejected) + }); + if !failed { + // SAFETY: scoped exclusive write; the stream now belongs to the sink. + unsafe { *body_value = BodyValue::Used }; + } + return Ok(ControlFlow::Break(promise)); + } + // The producer settled the body while the stream was being made. + // SAFETY: scoped borrows, as in the arms above. + match unsafe { &mut *body_value } { + BodyValue::Locked(_) => {} + BodyValue::Error(err) => { + let err_js = err.to_js(global_this); + destination_blob.detach(); + // SAFETY: `err` is not used after `to_js`, so this is the only + // live borrow of the value. + let _ = unsafe { (*body_value).use_() }; + return Ok(ControlFlow::Break( + JSPromise::dangerously_create_rejected_promise_value_without_notifying_vm( + global_this, + err_js, + ), + )); + } + // SAFETY: the match borrow ended with the pattern; no other borrow is live. + _ => return Ok(ControlFlow::Continue(unsafe { (*body_value).use_() })), + } + } let task = bun_core::heap::into_raw(Box::new(WriteFileWaitFromLockedValueTask { global_this: bun_ptr::BackRef::new(global_this), @@ -5079,6 +5204,14 @@ pub(crate) fn write_file_internal( break 'brk Blob::init_with_store(archive.store_ref().clone(), global_this); } + if let Some(readable) = ReadableStream::from_js_direct(data) { + if readable.is_locked(global_this) || readable.is_disturbed(global_this) { + destination_blob.detach(); + return Ok(body_used_rejection(global_this)); + } + return destination_blob.pipe_readable_stream_to_blob(global_this, readable, &options); + } + break 'brk Blob::get::(global_this, data)?; }; // Detach the source blob on scope exit. @@ -5789,7 +5922,10 @@ pub(crate) fn on_file_stream_resolve_request_stream( if let Some(stream) = strong.get() { stream.done(); } - this.promise.resolve(global_this, JSValue::js_number(0.0))?; + // SAFETY: the wrapper holds a ref on `sink` until it is dropped below. + let written = unsafe { (*this.sink).stream_bytes.get().unwrap_or(0) }; + this.promise + .resolve(global_this, JSValue::js_number(written as f64))?; Ok(JSValue::UNDEFINED) } diff --git a/src/runtime/webcore/Body.rs b/src/runtime/webcore/Body.rs index 92f7d6ad7051..63a242fe9e6b 100644 --- a/src/runtime/webcore/Body.rs +++ b/src/runtime/webcore/Body.rs @@ -774,10 +774,16 @@ impl Value { self.locked_to_native_stream(global_this, false) } Value::Error(err) => { - // Leave `self` as `Error` so the promise-returning readers - // (`handle_body_error`) still reject too. let reason = err.to_js(global_this); - ReadableStream::errored(global_this, reason) + let value = ReadableStream::errored(global_this, reason)?; + // As for a blob above: this stream is the body from here on, so `.body` hands it + // out again, `bodyUsed` follows it, and the promise readers reject through it. + let stream = ReadableStream::from_js_direct(value).unwrap(); + *self = Value::Locked(PendingValue { + readable: webcore::readable_stream::Strong::init(stream, global_this), + ..PendingValue::new(global_this) + }); + Ok(value) } } } @@ -825,7 +831,9 @@ impl Value { Value::Locked(_) => self.locked_to_native_stream(global_this, true), Value::Error(err) => { let reason = err.to_js(global_this); - ReadableStream::errored(global_this, reason) + let stream = ReadableStream::errored(global_this, reason)?; + *self = Value::Used; + Ok(stream) } } } diff --git a/src/runtime/webcore/ByteStream.rs b/src/runtime/webcore/ByteStream.rs index a475699818e7..f6c4a611b6db 100644 --- a/src/runtime/webcore/ByteStream.rs +++ b/src/runtime/webcore/ByteStream.rs @@ -65,6 +65,137 @@ impl Default for ByteStream { /// ReadableStream source backed by a ByteStream. pub type Source = readable_stream::NewSource; +/// A network body producer's (fetch, S3) hold on the stream it feeds: a counted ref on the stream's +/// `Source`, so delivery and unhooking go through memory the producer keeps alive rather than the +/// JS wrapper (which the VM's last sweep destroys in no particular order), plus the parked bit of +/// the receive backpressure. The ref roots the wrapper except while parked, so an unread stream +/// can be collected (`SourceHandle::consumer_collected`). +#[derive(Default)] +pub struct ProducerHold { + source: Cell>>, + parked: Cell, +} + +/// The JS-thread half of `BODY_HIGH_WATER_MARK`, decided from the stream's buffer after a +/// delivery. The HTTP thread does the other half on its hop buffer. +pub enum AfterDelivery { + /// Under the mark, or a whole-body consumer (`readableStreamTo*`) is collecting: keep going. + Resume, + /// At the mark with a back-pressured sink: it resumes the producer when it drains. + Pause, + /// At the mark and nothing reads: pause, release the loop, leave the stream collectable. + Park, +} + +impl ProducerHold { + /// Take the producer ref on the stream's source (JS thread). + /// + /// # Safety + /// `bytes` is the live ByteStream of a stream the caller holds. + pub unsafe fn hold(&self, bytes: *mut ByteStream) { + self.release(); + // SAFETY: fn contract; the ref keeps the Source alive past this call. + unsafe { + let source = Source::from_context_ptr(bytes); + (*source).increment_count(); + self.source.set(core::ptr::NonNull::new(source)); + } + } + + pub fn is_held(&self) -> bool { + self.source.get().is_some() + } + + /// The held stream, pinned for the guard's life: a consumer inside `on_data` can cancel the + /// producer (which drops the hold), and while parked the wrapper is not rooted. + pub fn bytes(&self) -> Option { + let source = self.source.get()?; + // SAFETY: live through our ref; no borrow of the source exists yet. + unsafe { (*source.as_ptr()).increment_count() }; + Some(PinnedBytes(source)) + } + + /// Stop being the producer. The source stays pinned by the returned guard, so the caller can + /// still deliver a terminal chunk. Touches no JS cell. + pub fn take(&self) -> Option { + let source = self.source.take()?; + self.parked.set(false); + // SAFETY: still pinned by our ref, which the guard now owns. + unsafe { + (*source.as_ptr()).producer.set(streams::SourceHandle::None); + (*source.as_ptr()).wrapper_unrooted.set(false); + } + Some(PinnedBytes(source)) + } + + /// `take` and drop. Touches no JS cell (safe inside a GC sweep). + pub fn release(&self) { + drop(self.take()); + } + + pub fn after_delivery(bytes: &ByteStream) -> AfterDelivery { + if bytes.buffered_len() < bun_http::signals::BODY_HIGH_WATER_MARK + || bytes.buffer_action.get().is_some() + { + AfterDelivery::Resume + } else if bytes.sink.get().is_some() { + AfterDelivery::Pause + } else { + AfterDelivery::Park + } + } + + /// Returns whether this call parked (the caller then releases its loop ref). + pub fn park(&self) -> bool { + if self.parked.replace(true) { + return false; + } + if let Some(source) = self.source.get() { + // SAFETY: live through our ref. The caller may hold the `&ByteStream` of this very + // source (the chunk it just delivered), which is why this is not a method call. + unsafe { Source::unroot_wrapper(source.as_ptr()) }; + } + true + } + + /// Returns whether this call unparked (the caller then re-takes its loop ref). Reached from a + /// consumer holding the stream. + pub fn unpark(&self) -> bool { + if !self.parked.replace(false) { + return false; + } + if let Some(source) = self.source.get() { + // SAFETY: as in `park`. + unsafe { Source::root_wrapper(source.as_ptr()) }; + } + true + } +} + +impl Drop for ProducerHold { + fn drop(&mut self) { + self.release(); + } +} + +/// A counted ref on a stream's `Source` for the guard's life; derefs to its ByteStream. +pub struct PinnedBytes(core::ptr::NonNull); + +impl core::ops::Deref for PinnedBytes { + type Target = ByteStream; + fn deref(&self) -> &ByteStream { + // SAFETY: pinned by this guard's ref; ByteStream is `&self`-only. + unsafe { &(*self.0.as_ptr()).context } + } +} + +impl Drop for PinnedBytes { + fn drop(&mut self) { + // SAFETY: balances the ref this guard owns. Can free the source. + unsafe { Source::decrement_count(self.0.as_ptr()) }; + } +} + impl readable_stream::SourceContext for ByteStream { const NAME: &'static str = "Bytes"; // setRefUnrefFn = null @@ -663,6 +794,15 @@ impl ByteStream { } self.done.set(true); self.pending_value.with_mut(|pv| pv.deinit()); + // A native sink wired to this stream must fail, not later see an EOF and commit what it + // has (an S3 upload would complete with a truncated object). + let sink = self.sink.replace(SinkHandle::None); + if sink.is_some() { + self.sink_paused.set(false); + sink.end(Some(streams::StreamError::AbortReason( + jsc::CommonAbortReason::UserAbort, + ))); + } if !view.is_empty() { self.pending_buffer.set(Self::empty_pending_buffer()); diff --git a/src/runtime/webcore/FileSink.rs b/src/runtime/webcore/FileSink.rs index 174ced256e97..31cbef5857b1 100644 --- a/src/runtime/webcore/FileSink.rs +++ b/src/runtime/webcore/FileSink.rs @@ -57,6 +57,13 @@ pub struct FileSink { /// Currently, only used when `stdin` in `Bun.spawn` is a ReadableStream. pub(crate) readable_stream: JsCell, + /// `pipe_stream`: settled from `on_close` with `stream_bytes`, or with the error that ended + /// the stream or the write. + stream_done: JsCell, + stream_error: JsCell>, + /// Bytes accepted since `pipe_stream` (`written` counts buffered bytes again when flushed). + pub(crate) stream_bytes: Cell>, + /// Strong reference to the JS wrapper object to prevent GC from collecting it /// while an async operation is pending. This is set when endFromJS returns a /// pending Promise and cleared when the operation completes. @@ -193,6 +200,10 @@ pub struct Options { pub(crate) input_path: PathOrFileDescriptor, pub close: bool, pub(crate) mode: bun_sys::Mode, + /// `Bun.write(path, stream)`: replace the file's contents. + pub(crate) truncate: bool, + /// `Bun.write(path, stream)`: create missing parent directories. + pub(crate) mkdirp: bool, } impl Default for Options { @@ -201,14 +212,21 @@ impl Default for Options { input_path: PathOrFileDescriptor::Fd(Fd::INVALID), close: false, mode: 0o664, + truncate: false, + mkdirp: false, } } } impl Options { pub(crate) fn flags(&self) -> i32 { - let _ = self; - bun_sys::O::NONBLOCK | bun_sys::O::CLOEXEC | bun_sys::O::CREAT | bun_sys::O::WRONLY + let flags = + bun_sys::O::NONBLOCK | bun_sys::O::CLOEXEC | bun_sys::O::CREAT | bun_sys::O::WRONLY; + if self.truncate { + flags | bun_sys::O::TRUNC + } else { + flags + } } } @@ -456,7 +474,7 @@ impl FileSink { if status == WriteStatus::EndOfFile { (*this).writer.with_mut(|w| w.close()); } else { - (*this).writer.with_mut(|w| w.end()); + (*this).end_writer(); } } @@ -479,6 +497,7 @@ impl FileSink { // drop the last reference and free `this` before that `close()` runs. // SAFETY: caller contract — `this` is live with write+dealloc provenance. unsafe { + (*this).record_stream_error(streams::StreamError::Error(err.clone())); if (*this).pending.get().state == streams::PendingState::Pending { (*this) .pending @@ -529,6 +548,8 @@ impl FileSink { let mut src = *(*this).source.get(); src.close(None); + (*this).settle_stream_done(); + // The writer is fully closed; no further callbacks will arrive. Release // the ref taken when a write returned `.pending`. This must be the last // thing we do as it may free `this`. @@ -536,6 +557,34 @@ impl FileSink { } } + /// `writer.end()`; `on_close` follows and may free `self`, except (Windows) for an fd the writer + /// does not own, which is never closed: settle a piped stream here then. + fn end_writer(&self) { + #[cfg(windows)] + if !self.writer.get().owns_fd { + self.settle_stream_done(); + } + self.writer.with_mut(|w| w.end()); + } + + fn settle_stream_done(&self) { + let mut promise = self.stream_done.replace(bun_jsc::JSPromiseStrong::empty()); + if !promise.has_value() { + return; + } + let Some(global) = self.js_global() else { + return; + }; + let result = match self.stream_error.replace(None) { + Some(err) => promise.reject(global, Ok(err.to_js(global))), + None => promise.resolve( + global, + JSValue::js_number(self.stream_bytes.get().unwrap_or(0) as f64), + ), + }; + crate::dispatch::fold(result); + } + /// Release the ref taken in `toResult`/`end`/`endFromJS` when a write /// returned `.pending` and we needed to stay alive until it completed. /// Idempotent via the flag check. May free `this`. @@ -623,24 +672,52 @@ impl FileSink { PathOrFileDescriptor::Fd(fd) => bun_io::PathOrFileDescriptor::Fd(*fd), PathOrFileDescriptor::Path(slice) => bun_io::PathOrFileDescriptor::Path(slice.slice()), }; - let result = bun_io::open_for_writing( - Fd::cwd(), - &io_path, - options.flags(), - options.mode, + let open = |pollable_out: &mut bool, + is_socket_out: &mut bool, + nonblocking_out: &mut bool, + force_sync_out: &mut bool| { + bun_io::open_for_writing( + Fd::cwd(), + &io_path, + options.flags(), + options.mode, + pollable_out, + is_socket_out, + self.force_sync.get(), + nonblocking_out, + force_sync_out, + |_fs: &mut bool| { + #[cfg(unix)] + { + *_fs = true; + } + }, + is_pollable, + ) + }; + let mut result = open( &mut pollable_out, &mut is_socket_out, - self.force_sync.get(), &mut nonblocking_out, &mut force_sync_out, - |_fs: &mut bool| { - #[cfg(unix)] - { - *_fs = true; - } - }, - is_pollable, ); + if options.mkdirp { + if let (sys::Result::Err(err), bun_io::PathOrFileDescriptor::Path(path)) = + (&result, &io_path) + { + if err.get_errno() == sys::E::ENOENT { + result = match webcore::blob::mkdirp_parent(path) { + Ok(()) => open( + &mut pollable_out, + &mut is_socket_out, + &mut nonblocking_out, + &mut force_sync_out, + ), + Err(err) => Err(err), + }; + } + } + } self.pollable.set(pollable_out); self.is_socket.set(is_socket_out); self.nonblocking.set(nonblocking_out); @@ -840,6 +917,7 @@ impl FileSink { // `run_pending_later()` alone would resolve it as if every // buffered byte had reached the reader. Latch the error and // move the sink to its terminal state (mirrors `end_from_js`). + (*this).record_stream_error(streams::StreamError::Error(err.clone())); (*this).done.set(true); if (*this).pending.get().state == streams::PendingState::Pending { (*this) @@ -1031,10 +1109,35 @@ impl FileSink { let buffered_before = self.writer.get().buffered_len(); // SAFETY(JsCell): `IOWriter::write` buffers/writes to fd; does not call JS. let rc = self.writer.with_mut(|w| w.write(data.slice())); + if self.counting_stream_bytes() { + self.count_stream_bytes(&rc, data.slice().len()); + } let accepted = self.bytes_accepted(buffered_before, &rc); self.to_result(rc, accepted) } + fn count_stream_bytes(&self, rc: &WriteResult, encoded_len: usize) { + let counted = self.stream_bytes.get().unwrap_or(0); + match rc { + WriteResult::Err(err) => { + self.record_stream_error(streams::StreamError::Error(err.clone())) + } + WriteResult::Done(n) => self.stream_bytes.set(Some(counted + *n as u64)), + _ => self.stream_bytes.set(Some(counted + encoded_len as u64)), + } + } + + /// Only `Bun.write(dest, stream)` reads the count; nothing else pays for the encoded length. + fn counting_stream_bytes(&self) -> bool { + self.stream_bytes.get().is_some() + } + + fn record_stream_error(&self, err: streams::StreamError) { + if self.stream_error.get().is_none() { + self.stream_error.set(Some(err)); + } + } + pub(crate) fn write_latin1(&self, data: &streams::Result) -> streams::Writable { if self.done.get() { return streams::Writable::Done; @@ -1042,6 +1145,12 @@ impl FileSink { let buffered_before = self.writer.get().buffered_len(); // SAFETY(JsCell): `IOWriter::write_latin1` buffers/writes; no JS. let rc = self.writer.with_mut(|w| w.write_latin1(data.slice())); + if self.counting_stream_bytes() { + self.count_stream_bytes( + &rc, + bun_core::strings::element_length_latin1_into_utf8(data.slice()), + ); + } let accepted = self.bytes_accepted(buffered_before, &rc); self.to_result(rc, accepted) } @@ -1053,6 +1162,12 @@ impl FileSink { let buffered_before = self.writer.get().buffered_len(); // SAFETY(JsCell): `IOWriter::write_utf16` buffers/writes; no JS. let rc = self.writer.with_mut(|w| w.write_utf16(data.slice16())); + if self.counting_stream_bytes() { + self.count_stream_bytes( + &rc, + bun_core::strings::element_length_utf16_into_utf8(data.slice16()), + ); + } let accepted = self.bytes_accepted(buffered_before, &rc); self.to_result(rc, accepted) } @@ -1067,11 +1182,17 @@ impl FileSink { // detach so the writer's `on_close` → `source.close()` is a no-op. self.source.with_mut(|s| s.clear()); } - if err.is_none() || !is_byte_stream { - let sys_err = match err { - Some(streams::StreamError::Error(e)) => Some(e), - _ => None, - }; + let errored = err.is_some(); + // A failed `write()` recorded its error before the source called back here. + let write_failed = self.stream_error.get().is_some(); + let sys_err = match &err { + Some(streams::StreamError::Error(e)) => Some(e.clone()), + _ => None, + }; + if let Some(err) = err { + self.record_stream_error(err); + } + if !errored || !is_byte_stream { let _ = self.end(sys_err); return; } @@ -1079,8 +1200,16 @@ impl FileSink { return; } self.done.set(true); - self.readable_stream - .with_mut(|rs| *rs = readable_stream::Strong::default()); + // The source stopped because a write failed (it cleared its sink first): nothing reads + // the rest, so cancel it. A source that failed on its own is already done. + let readable_stream = self + .readable_stream + .replace(readable_stream::Strong::default()); + if write_failed { + if let (Some(stream), Some(global)) = (readable_stream.get(), self.js_global()) { + crate::dispatch::fold(stream.cancel(global)); + } + } self.writer.with_mut(|w| w.close()); } @@ -1101,23 +1230,24 @@ impl FileSink { match self.writer.with_mut(|w| w.flush()) { WriteResult::Done(written) | WriteResult::Wrote(written) => { self.written.set(self.written.get() + written as usize); // @truncate - self.writer.with_mut(|w| w.end()); if has_pending { // `to_result` already seeded `Owned(consumed)`; just deliver it. self.run_pending_later(); } + self.end_writer(); sys::Result::Ok(()) } WriteResult::Err(e) => { + self.record_stream_error(streams::StreamError::Error(e.clone())); self.done.set(true); if has_pending { self.pending .with_mut(|p| p.result = streams::Writable::Err(e)); - self.writer.with_mut(|w| w.end()); self.run_pending_later(); + self.end_writer(); return sys::Result::Ok(()); } - self.writer.with_mut(|w| w.end()); + self.end_writer(); sys::Result::Err(e) } WriteResult::Pending(written) => { @@ -1436,6 +1566,9 @@ impl FileSink { auto_flusher: JsCell::new(AutoFlusher::default()), run_pending_later: FlushPendingTask::default(), readable_stream: JsCell::new(readable_stream::Strong::default()), + stream_done: JsCell::new(bun_jsc::JSPromiseStrong::empty()), + stream_error: JsCell::new(None), + stream_bytes: Cell::new(None), js_sink_ref: JsCell::new(bun_jsc::strong::Optional::empty()), } } @@ -1537,6 +1670,56 @@ fn on_reject_stream(global_this: &JSGlobalObject, callframe: &CallFrame) -> JsRe } impl FileSink { + /// `Bun.write(file, stream)`: wire `stream`'s native source straight to this sink and return a + /// promise for the byte count once the file is closed. `None` if the stream is not a native + /// source; the caller falls back to the JS pump. + pub fn pipe_stream( + &mut self, + stream: &ReadableStream, + global_this: &JSGlobalObject, + ) -> Option { + // SAFETY: `&mut self` carries write+dealloc provenance over the allocation. + let _guard = unsafe { FileSinkRef::new_ref(std::ptr::from_mut::(self)) }; + + self.stream_bytes.set(Some(0)); + self.stream_done + .set(bun_jsc::JSPromiseStrong::init(global_this)); + let promise = self.stream_done.get().value(); + self.readable_stream + .set(readable_stream::Strong::init(*stream, global_this)); + + match stream.wire_native_sink( + global_this, + webcore::SinkHandle::FileSink(bun_ptr::BackRef::new(&*self)), + JSValue::UNDEFINED, + |src| self.source.set(src), + ) { + readable_stream::NativeWireResult::Wired => { + // A synchronous source (FileReader over a regular file) may have run to the + // end inside `wire_native_sink`; the promise is settled then. + if self.stream_done.get().has_value() && !self.done.get() { + self.writer + .with_mut(|w| w.enable_keeping_process_alive(self.io_evtloop())); + if !self.must_be_kept_alive_until_eof.get() { + self.must_be_kept_alive_until_eof.set(true); + self.ref_(); + } + } + Some(promise) + } + readable_stream::NativeWireResult::EndedInline(err) => { + self.source.set(streams::SourceHandle::None); + self.end_from_stream(err); + Some(promise) + } + readable_stream::NativeWireResult::NotNative => { + self.stream_done.set(bun_jsc::JSPromiseStrong::empty()); + self.readable_stream.set(readable_stream::Strong::default()); + None + } + } + } + pub fn assign_to_stream( &mut self, stream: &mut ReadableStream, diff --git a/src/runtime/webcore/ReadableStream.rs b/src/runtime/webcore/ReadableStream.rs index 5f43cc0469d1..4577965b709b 100644 --- a/src/runtime/webcore/ReadableStream.rs +++ b/src/runtime/webcore/ReadableStream.rs @@ -1140,14 +1140,20 @@ impl NewSource { // `on_js_close`, reached from `on_reader_done` off the event loop with // no JS frame on the stack, never reads a dead-but-unswept cell. if !self.wrapper_unrooted.get() { - self.upgrade_wrapper(); + // SAFETY: `self` is live for the call. + unsafe { Self::upgrade_wrapper(self) }; } } - fn upgrade_wrapper(&mut self) { - if let Some(global) = self.global_this.as_deref() { - if self.this_jsvalue.is_not_empty() { - self.this_jsvalue.upgrade(global); + /// # Safety + /// `this` points at a live `NewSource`. + unsafe fn upgrade_wrapper(this: *mut Self) { + // SAFETY: fn contract; field places only, see `unroot_wrapper`. + unsafe { + if let Some(global) = (*this).global_this.as_deref() { + if (*this).this_jsvalue.is_not_empty() { + (*this).this_jsvalue.upgrade(global); + } } } } @@ -1155,18 +1161,32 @@ impl NewSource { /// The producer keeps its native ref but stops rooting the wrapper: nothing /// is reading, so the stream should be collectable. [`SourceContext::wrapper_finalized`] /// tells the producer if that happens. - /// Same access pattern as [`Self::increment_count`]: reached through the - /// producer's raw pointer while the context may be borrowed. - pub fn unroot_wrapper(&mut self) { - self.wrapper_unrooted.set(true); - self.this_jsvalue.downgrade(); + /// + /// Takes a raw pointer: the producer reaches this while it holds a `&C` into + /// `this` (the chunk it is delivering to), so only the fields written here + /// are touched, never a `&mut Self` that would cover the context too. + /// + /// # Safety + /// `this` points at a live `NewSource`. + pub unsafe fn unroot_wrapper(this: *mut Self) { + // SAFETY: fn contract. + unsafe { + (*this).wrapper_unrooted.set(true); + (*this).this_jsvalue.downgrade(); + } } /// Undo [`Self::unroot_wrapper`]: a consumer is reading again. - pub fn root_wrapper(&mut self) { - self.wrapper_unrooted.set(false); - if self.ref_count > 1 { - self.upgrade_wrapper(); + /// + /// # Safety + /// As [`Self::unroot_wrapper`]. + pub unsafe fn root_wrapper(this: *mut Self) { + // SAFETY: fn contract. + unsafe { + (*this).wrapper_unrooted.set(false); + if (*this).ref_count > 1 { + Self::upgrade_wrapper(this); + } } } diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index d5ab7d5d7632..1da0c765a883 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -126,7 +126,7 @@ impl BodyAbortListener { // `attach_abort_signal`; `clean_native_bindings` removes it before the // box is dropped, so it is live here. Copy out up front: erroring a // still-streaming body can re-enter `Response::unref` via - // `FetchTasklet::ignore_remaining_response_body` and destroy this box. + // `FetchTasklet::abandon_response_body` and destroy this box. let (response, global) = unsafe { ((*ctx.cast::()).response, (*ctx.cast::()).global) }; Response::ref_(response.as_mut_ptr()); diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 3a854c808c48..3f4f15a1cdf8 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -38,10 +38,7 @@ use crate::webcore::{AbortSignal, DrainResult, FetchHeaders, InternalBlob, Respo // ConcurrentTask callbacks at the tier-3 layer. type ElJsResult = bun_event_loop::JsResult; -/// How much of a body nothing is reading is taken off the socket before the transport is left -/// paused (`FetchTasklet::after_body_chunk_delivered`). Below this an unread body still -/// completes, which frees the connection for reuse; above it memory stays bounded. -const UNOBSERVED_BODY_HIGH_WATER_MARK: usize = 256 * 1024; +use http::signals::BODY_HIGH_WATER_MARK; use boringssl::c::{X509_free, d2i_X509}; @@ -119,17 +116,8 @@ pub struct FetchTasklet { // Response is intrusively refcounted; modeled as a raw ptr. `Cell`: released from // `on_body_stream_collected`, which only has a shared ref. pub(crate) native_response: Cell>, - /// A counted ref on the response body stream's ByteStream source for as long as this tasklet is its - /// `producer`, so delivery and unhooking go through native memory we keep alive rather than - /// through the JS wrappers, which the VM's last sweep destroys in no particular order. The - /// ref also roots the stream's JS wrapper, except while the body is parked - /// (`body_stream_parked`). - pub(crate) response_stream_source: Cell>>, - /// Nothing reads the body stream and it holds `UNOBSERVED_BODY_HIGH_WATER_MARK` bytes: the - /// transport is left paused, the loop is released, and the stream is left collectable - /// (`on_body_stream_collected`). A consumer that drains the stream's buffer or a native - /// sink that attaches unparks it. - pub(crate) body_stream_parked: Cell, + /// The response body stream while this tasklet is its producer. + pub(crate) response_stream: crate::webcore::byte_stream::ProducerHold, pub(crate) request_headers: Headers, pub(crate) promise: jsc::JSPromiseStrong, pub(crate) concurrent_task: ConcurrentTask, @@ -163,10 +151,6 @@ pub struct FetchTasklet { // Custom Hostname pub(crate) hostname: Option>, pub(crate) is_waiting_body: bool, - /// Set by `on_start_buffering_callback` (JS thread) and read by - /// `callback()` (HTTP thread, under `mutex`): the body is being - /// accumulated in `scheduled_response_buffer` for a buffered consumer. - pub(crate) is_buffering_body: AtomicBool, pub(crate) is_waiting_abort: bool, pub(crate) is_waiting_request_stream_start: bool, pub(crate) mutex: Mutex, @@ -784,13 +768,12 @@ impl FetchTasklet { let mut err = scopeguard::guard(self.on_reject(), |mut e| e.reset()); let mut js_err = JSValue::ZERO; // if we are streaming update with error - if let Some(source) = self.take_response_stream_source() { + if let Some(bytes) = self.response_stream.take() { js_err = err.to_js(&global_this); js_err.ensure_still_alive(); - Self::response_bytes(source).on_data(StreamResult::Err(StreamError::JSValue( + bytes.on_data(StreamResult::Err(StreamError::JSValue( bun_jsc::strong::Optional::create(js_err, &global_this), ))); - Self::release_response_stream_source(source); } // A failure result is terminal (`to_result` forces `has_more = // false` once `fail` is set), so everything pending must settle @@ -823,30 +806,22 @@ impl FetchTasklet { if !self.result.has_more { // Unhook before the final delivery so it cannot signal a producer that is done; // release after it so the bytes land in memory we still pin. - if let Some(source) = self.take_response_stream_source() { - bun_output::scoped_log!(FetchTasklet, "onBodyReceived response_stream_source done"); - let bytes = Self::response_bytes(source); + if let Some(bytes) = self.response_stream.take() { + bun_output::scoped_log!(FetchTasklet, "onBodyReceived response_stream done"); bytes.size_hint.set(self.get_size_hint()); buffer_reset.set(false); let chunk = self.scheduled_response_buffer.list.as_slice(); bytes.on_data(Self::temporary_chunk(chunk, true)); - Self::release_response_stream_source(source); return Ok(()); } - } else if let Some(source) = self.response_stream_source.get() { - bun_output::scoped_log!(FetchTasklet, "onBodyReceived response_stream_source"); - // Pin across the delivery: while parked the wrapper is not rooted, and a consumer - // inside `on_data` can cancel us (`on_stream_cancelled` drops the producer ref). - // SAFETY: live through the producer ref; no borrow of the source exists yet. - unsafe { (*source.as_ptr()).increment_count() }; - let bytes = Self::response_bytes(source); + } else if let Some(bytes) = self.response_stream.bytes() { + bun_output::scoped_log!(FetchTasklet, "onBodyReceived response_stream"); bytes.size_hint.set(self.get_size_hint()); let chunk = self.scheduled_response_buffer.list.as_slice(); bytes.on_data(Self::temporary_chunk(chunk, false)); - if self.response_stream_source.get().is_some() { + if self.response_stream.is_held() { self.after_body_chunk_delivered(&bytes); } - Self::release_response_stream_source(source); return Ok(()); } @@ -1617,20 +1592,12 @@ impl FetchTasklet { readable: ReadableStream, ) { let this = Self::from_ctx(ctx); - this.clear_stream_handlers(); if let crate::webcore::readable_stream::Source::Bytes(bytes) = readable.ptr { - // SAFETY: the caller holds the stream, which owns the live ByteStream embedded in - // its Source; the ref taken here keeps the Source (and, until parked, the stream) - // alive past this call. JS thread. - unsafe { - let source = crate::webcore::readable_stream::NewSource::from_context_ptr(bytes); - (*source).increment_count(); - this.response_stream_source.set(NonNull::new(source)); - } + // SAFETY: the caller holds the stream, which owns the live ByteStream. JS thread. + unsafe { this.response_stream.hold(bytes) }; + } else { + this.response_stream.release(); } - // A ByteStream now drains scheduled_response_buffer per chunk; undo any - // buffered-consumer reservation request so callback() stops growing it. - this.is_buffering_body.store(false, Ordering::Release); } fn on_start_streaming_http_response_body_callback(ctx: NonNull) -> DrainResult { @@ -1645,28 +1612,16 @@ impl FetchTasklet { this.poll_ref .with_mut(|poll_ref| poll_ref.ref_(bun_io::js_vm_ctx())); - if let Some(http_) = this.http.as_mut() { - http_.enable_response_body_streaming(); - } - this.mutex.lock(); - // A ByteStream is attaching; clear the buffered-consumer reserve gate - // under the mutex so the HTTP thread cannot observe the stale `true` - // in `callback()` between this unlock and `on_readable_stream_available`. - this.is_buffering_body.store(false, Ordering::Release); let size_hint = this.get_size_hint() as usize; - // This means we have received part of the body but not the whole thing. let drained = core::mem::take(&mut this.scheduled_response_buffer.list); this.mutex.unlock(); - // The bytes taken above are the drain `Paused` was waiting for: flip back to - // `AutoPause` and resume. After the drain, not before: a chunk the HTTP thread - // appends (and pauses for) in between would otherwise be handed to the stream here - // with its task finding the buffer empty, and nothing left to undo that pause. - // Also covers headers and body arriving in separate writes with no follow-up - // data: the HTTP thread must flush what it has. - this.signal_store - .try_transition_receive_mode(BodyReceiveMode::Paused, BodyReceiveMode::AutoPause); + // After the take, not before: a chunk the HTTP thread appends (and pauses for) in + // 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.schedule_receive_resume(); if drained.is_empty() { @@ -1687,63 +1642,23 @@ impl FetchTasklet { } } - /// The ByteStream behind a source one of our refs still pins. - fn response_bytes( - source: NonNull, - ) -> bun_ptr::BackRef { - // SAFETY: the source is live (caller's ref); ByteStream is `&self`-only (R-2). - bun_ptr::BackRef::new(unsafe { &(*source.as_ptr()).context }) - } - - /// Stop being the producer. The source stays pinned by our ref, so the caller can still - /// deliver a terminal chunk before `release_response_stream_source`. Touches no JS cell. - fn take_response_stream_source(&self) -> Option> { - let source = self.response_stream_source.take()?; - self.body_stream_parked.set(false); - // SAFETY: still pinned by the ref from `on_readable_stream_available`. - unsafe { - (*source.as_ptr()).producer.set(SourceHandle::None); - (*source.as_ptr()).wrapper_unrooted.set(false); - } - Some(source) - } - - /// Drop one of our refs (the producer's, or a delivery's pin). Can free the source: the - /// JS wrapper may be gone already. - fn release_response_stream_source(source: NonNull) { - // SAFETY: balances a ref this tasklet took; `source` is not used after this. - unsafe { crate::webcore::byte_stream::Source::decrement_count(source.as_ptr()) }; - } - - /// Unhook from the response ByteStream (the stream can outlive us in JS) and release - /// it. Touches no JS cell. + /// Unhook from the response ByteStream (the stream can outlive us in JS). Touches no JS cell. fn clear_stream_handlers(&self) { - if let Some(source) = self.take_response_stream_source() { - Self::release_response_stream_source(source); - } + self.response_stream.release(); } - pub(crate) fn on_stream_cancelled(&mut self) { - if self.signal_store.body_receive_mode() == BodyReceiveMode::Ignore { - return; - } - // reader.cancel() / body.cancel() aborts the fetch so the server sees the - // close (Node/Deno/browsers abort unconditionally). abort_task() is idempotent. + /// reader.cancel() / body.cancel(): the server has to see the close (Node, Deno and browsers + /// abort too). `&self` because a failed sink write reaches here from inside `on_body_received`. + pub(crate) fn on_stream_cancelled(&self) { self.abort_task(); - self.ignore_remaining_response_body(); + self.abandon_response_body(); } /// `SourceHandle::consumer_collected`: the parked stream's wrapper was swept, so nothing - /// can read the rest of the body. Inside a GC sweep: native state only, like - /// `on_response_finalize`. The abort comes back through the HTTP callback and the usual - /// teardown runs there. + /// can read the rest of the body. Inside a GC sweep, like `on_response_finalize`. pub(crate) fn on_body_stream_collected(&self) { bun_output::scoped_log!(FetchTasklet, "onBodyStreamCollected"); - if self.signal_store.body_receive_mode() == BodyReceiveMode::Ignore { - return; - } - self.abort_transport(); - self.ignore_remaining_response_body(); + self.abandon_response_body(); } pub(crate) fn on_stream_drained(&self) { @@ -1752,10 +1667,7 @@ impl FetchTasklet { } fn resume_receive(&self) { - if self - .signal_store - .try_transition_receive_mode(BodyReceiveMode::Paused, BodyReceiveMode::AutoPause) - { + if self.signal_store.unpause_receive() { self.schedule_receive_resume(); } } @@ -1765,55 +1677,37 @@ impl FetchTasklet { self.unpark_body_stream(); } - /// After a non-terminal delivery. A consumer waiting on the stream resumes the transport - /// itself as it drains. Bytes nothing took stay in the stream's buffer: keep receiving until - /// it holds `UNOBSERVED_BODY_HIGH_WATER_MARK` (a small body still completes, which frees the - /// connection), then park. + /// The other half of this rule is in `callback` (HTTP thread). fn after_body_chunk_delivered(&self, bytes: &crate::webcore::ByteStream) { + use crate::webcore::byte_stream::{AfterDelivery, ProducerHold}; bun_output::scoped_log!( FetchTasklet, - "afterBodyChunkDelivered sink={} action={} pending={} buffered={}", - bytes.sink.get().is_some(), - bytes.buffer_action.get().is_some(), - bytes.pending.get().state == crate::webcore::streams::PendingState::Pending, + "afterBodyChunkDelivered buffered={}", bytes.buffered_len() ); - if bytes.sink.get().is_some() - || bytes.buffer_action.get().is_some() - || bytes.pending.get().state == crate::webcore::streams::PendingState::Pending - { - return; - } - if bytes.buffered_len() < UNOBSERVED_BODY_HIGH_WATER_MARK { - self.resume_receive(); - } else { - self.park_body_stream(); + match ProducerHold::after_delivery(bytes) { + AfterDelivery::Resume => self.resume_receive(), + AfterDelivery::Pause => self.signal_store.pause_receive(), + AfterDelivery::Park => { + self.signal_store.pause_receive(); + self.park_body_stream(); + } } } fn park_body_stream(&self) { - if self.body_stream_parked.replace(true) { - return; - } - bun_output::scoped_log!(FetchTasklet, "parkBodyStream"); - self.poll_ref - .with_mut(|poll_ref| poll_ref.unref(bun_io::js_vm_ctx())); - if let Some(source) = self.response_stream_source.get() { - // SAFETY: live through the producer ref (same pattern as `increment_count`). - unsafe { (*source.as_ptr()).unroot_wrapper() }; + if self.response_stream.park() { + bun_output::scoped_log!(FetchTasklet, "parkBodyStream"); + self.poll_ref + .with_mut(|poll_ref| poll_ref.unref(bun_io::js_vm_ctx())); } } fn unpark_body_stream(&self) { - if !self.body_stream_parked.replace(false) { - return; - } - bun_output::scoped_log!(FetchTasklet, "unparkBodyStream"); - self.poll_ref - .with_mut(|poll_ref| poll_ref.ref_(bun_io::js_vm_ctx())); - if let Some(source) = self.response_stream_source.get() { - // SAFETY: as in `park_body_stream`; reached from a consumer holding the stream. - unsafe { (*source.as_ptr()).root_wrapper() }; + if self.response_stream.unpark() { + bun_output::scoped_log!(FetchTasklet, "unparkBodyStream"); + self.poll_ref + .with_mut(|poll_ref| poll_ref.ref_(bun_io::js_vm_ctx())); } } @@ -1821,11 +1715,7 @@ impl FetchTasklet { let this = Self::from_ctx(ctx); this.poll_ref .with_mut(|poll_ref| poll_ref.ref_(bun_io::js_vm_ctx())); - this.is_buffering_body.store(true, Ordering::Release); - if this - .signal_store - .set_receive_mode_terminal(BodyReceiveMode::BufferAll) - { + if this.signal_store.receive_all() { this.schedule_receive_resume(); } } @@ -1878,14 +1768,12 @@ impl FetchTasklet { .is_some_and(|http_| http_.method() == Method::HEAD) } - /// Content the server frames anyway (a 205 with content) is dropped, and the - /// rest of it drained as for a Response collected unread, which keeps the - /// connection poolable. No Response is attached yet, so that is all the - /// ignore does. `is_waiting_body` stays false: nothing may reach this body. + /// Content the server frames anyway (a 205 with content) is dropped and the connection + /// closed. `is_waiting_body` stays false: nothing may reach this body. fn null_body_value(&mut self) -> BodyValue { self.scheduled_response_buffer = MutableString::default(); if self.result.has_more { - self.ignore_remaining_response_body(); + self.abandon_response_body(); } BodyValue::Null } @@ -1936,35 +1824,18 @@ impl FetchTasklet { ) } - fn ignore_remaining_response_body(&self) { - bun_output::scoped_log!(FetchTasklet, "ignoreRemainingResponseBody"); - // enabling streaming will make the http thread to drain into the main thread (aka stop buffering) - // without a stream ref, response body or response instance alive it will just ignore the result - // An aborted fetch is already shutting down; don't re-arm receive/resume - // draining, which would read the rest of an unbounded body and hold the - // socket open (drain_events resumes before shutdowns). - let aborted = self.signal_store.aborted.load(Ordering::Relaxed); - if self - .signal_store - .set_receive_mode_terminal(BodyReceiveMode::Ignore) - && !aborted - { - self.schedule_receive_resume(); - } - if let Some(http_) = self.http.as_deref() { - if !aborted { - http_.enable_response_body_streaming(); - } - } - // we should not keep the process alive if we are ignoring the body + /// Nothing will read the rest of the body: abort the transport, let go of the loop and of the + /// response; `callback` drops whatever still arrives. Safe inside a GC sweep + /// (`on_response_finalize`, `on_body_stream_collected`): no JS cell is touched; the + /// request-body sink is left for `clear_sink()` in `deinit()`. + fn abandon_response_body(&self) { + bun_output::scoped_log!(FetchTasklet, "abandonResponseBody"); + self.signal_store.abandon(); + self.abort_transport(); self.poll_ref .with_mut(|poll_ref| poll_ref.unref(bun_io::js_vm_ctx())); - // Also fine from GC sweeps (`on_response_finalize`, `on_body_stream_collected`): - // unhooking touches no JS cell. The request-body sink is left for `clear_sink()` in - // `deinit()` (an event-loop task, outside sweep) to detach. self.clear_stream_handlers(); self.response.clear(); - if let Some(response) = self.native_response.take() { // SAFETY: `response` is the +1 ref held in `native_response`. Response::unref(response); @@ -2024,8 +1895,7 @@ impl FetchTasklet { scheduled_response_buffer: MutableString::default(), response: jsc::Weak::default(), native_response: Cell::new(None), - response_stream_source: Cell::new(None), - body_stream_parked: Cell::new(false), + response_stream: Default::default(), request_headers: fetch_options.headers, promise, concurrent_task: ConcurrentTask::default(), @@ -2042,7 +1912,6 @@ impl FetchTasklet { upgraded_connection: fetch_options.upgraded_connection, hostname: fetch_options.hostname, is_waiting_body: false, - is_buffering_body: AtomicBool::new(false), is_waiting_abort: false, is_waiting_request_stream_start: false, mutex: Mutex::new(), @@ -2411,15 +2280,14 @@ impl FetchTasklet { FetchTasklet::deref(this_ptr); } - fn abort_task(&mut self) { + fn abort_task(&self) { if self.abort_transport() { self.tracker.did_cancel(&self.global_this); } } - /// Idempotent: reader.cancel(), an AbortSignal and a collected body stream can all - /// reach here for the same fetch. Only the first abort enqueues a shutdown; a second - /// would append a redundant ShutdownMessage for an already-closing socket. No JS. + /// Idempotent: an AbortSignal, VM teardown and `abandon_response_body` can all reach here for + /// the same fetch. Only the first enqueues a shutdown. No JS. fn abort_transport(&self) -> bool { if self.signal_store.aborted.swap(true, Ordering::Relaxed) { return false; @@ -2577,12 +2445,11 @@ impl FetchTasklet { let success = task_ref.result.is_success(); - if task_ref.signal_store.body_receive_mode() == BodyReceiveMode::Ignore { + if task_ref.signal_store.body_receive_mode() == BodyReceiveMode::Abandoned { if task_ref.scheduled_response_buffer.list.capacity() > 0 { task_ref.scheduled_response_buffer = MutableString::default(); } if success && task_ref.result.has_more { - // we are ignoring the body so we should not receive more data, so will only signal when result.has_more = true task_ref.mutex.unlock(); return; } @@ -2593,9 +2460,8 @@ impl FetchTasklet { } else { // Grow to Content-Length once so the per-packet append below // doesn't leave the ~2x doubling over-capacity that the - // ArrayBuffer would adopt. Gated on `is_buffering_body` - // (set by `on_start_buffering_callback`). - if task_ref.is_buffering_body.load(Ordering::Acquire) { + // ArrayBuffer would adopt. Only for a consumer that wants the whole body. + if task_ref.signal_store.body_receive_mode() == BodyReceiveMode::BufferAll { if let http::BodySize::ContentLength(n) = task_ref.body_size { if n > scheduled.list.capacity() { let additional = n @@ -2614,11 +2480,10 @@ impl FetchTasklet { bun_core::handle_oom(scheduled.write(chunk)); } } - if task_ref.result.has_more && !task_ref.scheduled_response_buffer.list.is_empty() { - let _ = task_ref.signal_store.try_transition_receive_mode( - BodyReceiveMode::AutoPause, - BodyReceiveMode::Paused, - ); + if task_ref.result.has_more + && task_ref.scheduled_response_buffer.list.len() >= BODY_HIGH_WATER_MARK + { + task_ref.signal_store.pause_receive(); } } @@ -2731,38 +2596,24 @@ impl FetchTasklet { #[bun_uws::uws_callback(export = "Bun__FetchResponse_finalize", no_catch)] pub(crate) fn on_response_finalize(&mut self) { bun_output::scoped_log!(FetchTasklet, "onResponseFinalize"); - let this = self; - if let Some(response) = this.native_response.get() { - // SAFETY: native_response is intrusively-ref'd by FetchTasklet; alive until unref. - let body = unsafe { (*response).get_body_value() }; - // Three scenarios: - // - // 1. We are streaming, in which case we should not ignore the body. - // 2. We were buffering, in which case - // 2a. if we have no promise, we should ignore the body. - // 2b. if we have a promise, we should keep loading the body. - // 3. We never started buffering, in which case we should ignore the body. - // - // Inside a finalizer: decide from native state only. - if !matches!(body, BodyValue::Locked(_)) || this.response_stream_source.get().is_some() - { - // Scenario 1 or 3. In scenario 1 the stream can outlive its Response; if - // nothing reads it, it parks and its own collection ends the body - // (`on_body_stream_collected`). - return; - } - - if let BodyValue::Locked(locked) = body { - if let Some(promise) = locked.promise { - if promise.is_empty_or_undefined_or_null() { - // Scenario 2b. - this.ignore_remaining_response_body(); - } - } else { - // Scenario 3. - this.ignore_remaining_response_body(); - } - } + let Some(response) = self.native_response.get() else { + return; + }; + // SAFETY: native_response is intrusively-ref'd by FetchTasklet; alive until unref. + let BodyValue::Locked(locked) = (unsafe { (*response).get_body_value() }) else { + // The body arrived or failed; nothing is underway. + return; + }; + // What can outlive the Response and still take the body: its stream (whose own collection + // is `on_body_stream_collected`), or a whole-body consumer (`.text()` and friends hold a + // promise, `Bun.write` an `on_receive_value`). + let outlived = self.response_stream.is_held() + || locked.on_receive_value.is_some() + || locked + .promise + .is_some_and(|promise| !promise.is_empty_or_undefined_or_null()); + if !outlived { + self.abandon_response_body(); } } } diff --git a/src/runtime/webcore/s3/client.rs b/src/runtime/webcore/s3/client.rs index 05c55a2b7ac9..0c1480673460 100644 --- a/src/runtime/webcore/s3/client.rs +++ b/src/runtime/webcore/s3/client.rs @@ -424,8 +424,12 @@ pub(crate) fn writable_stream( storage_class: Option, request_payer: bool, ) -> JsResult { - // Local callback wrapper - fn wrapper_callback(result: S3UploadResult, sink: &mut NetworkSink) -> JsResult<()> { + // Local callback wrapper. `uploaded` is read off the upload (see `MultiPartUpload::callback`). + fn wrapper_callback( + result: S3UploadResult, + uploaded: u64, + sink: &mut NetworkSink, + ) -> JsResult<()> { // `global_this` is a `BackRef` set at construction; copy it so the // re-borrow does not hold `&sink` across the `&mut sink` calls below. let global = sink @@ -445,7 +449,8 @@ pub(crate) fn writable_stream( .resolve(global, JSValue::js_number(0.0))?; } if sink.end_promise.has_value() { - sink.end_promise.resolve(global, JSValue::js_number(0.0))?; + sink.end_promise + .resolve(global, JSValue::js_number(uploaded as f64))?; } } S3UploadResult::Failure(err) => { @@ -468,9 +473,18 @@ pub(crate) fn writable_stream( // Thunks adapting typed callbacks to the erased `*mut c_void` signatures stored on // MultiPartUpload. - fn wrapper_callback_thunk(result: S3UploadResult, ctx: *mut c_void) -> JsResult<()> { - // SAFETY: ctx was set to `response_stream: *mut NetworkSink` below. - wrapper_callback(result, unsafe { bun_ptr::callback_ctx::(ctx) }) + fn wrapper_callback_thunk( + task: &MultiPartUpload, + result: S3UploadResult, + ctx: *mut c_void, + ) -> JsResult<()> { + let sink = ctx.cast::(); + // SAFETY: ctx was set to `response_stream: *mut NetworkSink` below; the box is live + // while the upload holds it. + let r = wrapper_callback(result, task.uploaded_bytes.get(), unsafe { &mut *sink }); + // SAFETY: the upload's hold on the box ends here; `sink` is not used afterwards. + unsafe { NetworkSink::release_writer_holder(sink) }; + r } fn on_writable_thunk(task: &MultiPartUpload, ctx: *mut c_void, flushed: u64) { NetworkSink::on_writable(task, ctx.cast::(), flushed); @@ -502,6 +516,7 @@ pub(crate) fn writable_stream( vm: VirtualMachine::get(), global_this: global_static, buffered: JsCell::new(StreamBuffer::default()), + uploaded_bytes: Cell::new(0), path: Box::<[u8]>::from(path), proxy: if !proxy_url.is_empty() { Box::<[u8]>::from(proxy_url) @@ -536,6 +551,7 @@ pub(crate) fn writable_stream( task: Some(unsafe { bun_ptr::BackRef::from_raw_mut(task_ptr) }), global_this: Some(bun_ptr::BackRef::new(global_this)), high_water_mark: part_size as BlobSizeType, + writer_holders: Cell::new(2), ..Default::default() })); @@ -674,18 +690,19 @@ impl S3UploadStreamWrapper { let mut settled: JsResult<()> = Ok(()); match &result { S3UploadResult::Success => { + let uploaded = JSValue::js_number(self_.task_ref().uploaded_bytes.get() as f64); if let Some(sink) = self_.sink_mut() { sink.pending.run(); if settled.is_ok() && sink.flush_promise.has_value() { settled = sink.flush_promise.resolve(&global, JSValue::js_number(0.0)); } if settled.is_ok() && sink.end_promise.has_value() { - settled = sink.end_promise.resolve(&global, JSValue::js_number(0.0)); + settled = sink.end_promise.resolve(&global, uploaded); } } if self_.end_promise.has_value() { if settled.is_ok() { - settled = self_.end_promise.resolve(&global, JSValue::js_number(0.0)); + settled = self_.end_promise.resolve(&global, uploaded); } self_.end_promise = bun_jsc::JSPromiseStrong::empty(); } @@ -913,7 +930,11 @@ pub(crate) fn upload_stream( // Thunks adapting typed callbacks to the erased `*mut c_void` signatures stored on // MultiPartUpload. - fn resolve_thunk(result: S3UploadResult, ctx: *mut c_void) -> JsResult<()> { + fn resolve_thunk( + _: &MultiPartUpload, + result: S3UploadResult, + ctx: *mut c_void, + ) -> JsResult<()> { // SAFETY: ctx was set to `*mut S3UploadStreamWrapper` below. S3UploadStreamWrapper::resolve(result, unsafe { bun_ptr::callback_ctx::(ctx) @@ -951,6 +972,7 @@ pub(crate) fn upload_stream( vm: VirtualMachine::get(), global_this: global_static, buffered: JsCell::new(StreamBuffer::default()), + uploaded_bytes: Cell::new(0), path: Box::<[u8]>::from(path), proxy: if !proxy_url.is_empty() { Box::<[u8]>::from(proxy_url) @@ -1273,7 +1295,7 @@ fn download_stream( None }; - task.signals = task.signal_store.to(); + task.signals = task.signal_store.to_with_backpressure(); let vm = VirtualMachine::get(); let verbose = vm.get_verbose_fetch(); @@ -1320,13 +1342,13 @@ fn download_stream( } pub struct S3DownloadStreamWrapper { - pub readable_stream_ref: ReadableStreamStrong, + stream: crate::webcore::byte_stream::ProducerHold, pub path: Box<[u8]>, pub global: GlobalRef, // JSC_BORROW /// Non-owning. The task frees itself on the main thread once `has_more == false`, /// which first drops this wrapper (clearing the stream's producer handle), so this /// pointer is never observed dangling from `on_stream_cancelled`. - pub task: *mut S3HttpDownloadStreamingTask, + pub task: Cell<*mut S3HttpDownloadStreamingTask>, } impl S3DownloadStreamWrapper { @@ -1334,65 +1356,119 @@ impl S3DownloadStreamWrapper { bun_core::heap::into_raw(Box::new(init)) } + /// `this` is the heap pointer from `new` (write and dealloc provenance): the terminal + /// callback frees the wrapper through it, so no reference derived from it may outlive + /// this call. fn callback( chunk: &MutableString, has_more: bool, request_err: Option, - self_: &mut Self, + this: *mut Self, ) { - // scope-exit cleanup via guard (keeps borrowck happy) - let _guard = scopeguard::guard(std::ptr::from_mut::(self_), move |s| { + let _guard = scopeguard::guard(this, move |s| { if !has_more { - // SAFETY: s is a live Box-allocated pointer (heap::alloc in S3DownloadStreamWrapper::new); - // reconstituting and dropping the Box runs Drop::drop and frees the allocation + // SAFETY: `s` is the live allocation from `new`; the HTTP thread does not call + // back after the terminal chunk, so this is the only owner left. drop(unsafe { bun_core::heap::take(s) }); } }); + // SAFETY: live until the guard runs, which is after the last use of this borrow. + let self_ = unsafe { &*this }; - if let Some(readable) = self_.readable_stream_ref.get() { - // BACKREF: see `Source::bytes()` — payload live while the - // readable stream is rooted. R-2: `&` — `on_data` re-enters JS. - if let Some(bytes) = readable.ptr.bytes() { - if let Some(err) = request_err { - bytes.on_data(crate::webcore::streams::StreamResult::Err( - crate::webcore::streams::StreamError::JSValue( - bun_jsc::strong::Optional::create( - s3_error_to_js(&err, &self_.global, Some(&self_.path)), - &self_.global, - ), - ), - )); - return; - } - if has_more { - bytes.on_data(crate::webcore::streams::StreamResult::Temporary( - // chunk.list is borrowed for the duration of on_data. - bun_ptr::RawSlice::new(chunk.list.as_slice()), - )); - return; - } - - bytes.on_data(crate::webcore::streams::StreamResult::TemporaryAndDone( - // chunk.list is borrowed for the duration of on_data. - bun_ptr::RawSlice::new(chunk.list.as_slice()), - )); + if let Some(err) = request_err { + let Some(bytes) = self_.stream.take() else { + return; + }; + bytes.on_data(crate::webcore::streams::StreamResult::Err( + crate::webcore::streams::StreamError::JSValue(bun_jsc::strong::Optional::create( + s3_error_to_js(&err, &self_.global, Some(&self_.path)), + &self_.global, + )), + )); + return; + } + if has_more { + let Some(bytes) = self_.stream.bytes() else { return; + }; + bytes.on_data(crate::webcore::streams::StreamResult::Temporary( + // chunk.list is borrowed for the duration of on_data. + bun_ptr::RawSlice::new(chunk.list.as_slice()), + )); + // `on_data` can cancel us, which releases the hold. + if self_.stream.is_held() { + self_.after_chunk_delivered(&bytes); + } + return; + } + let Some(bytes) = self_.stream.take() else { + return; + }; + bytes.on_data(crate::webcore::streams::StreamResult::TemporaryAndDone( + // chunk.list is borrowed for the duration of on_data. + bun_ptr::RawSlice::new(chunk.list.as_slice()), + )); + } + + /// The other half of this rule is in `S3HttpDownloadStreamingTask::process_http_callback` + /// (HTTP thread). + fn after_chunk_delivered(&self, bytes: &ByteStream) { + use crate::webcore::byte_stream::{AfterDelivery, ProducerHold}; + let task = self.task.get(); + if task.is_null() { + return; + } + match ProducerHold::after_delivery(bytes) { + // SAFETY: see `task`. + AfterDelivery::Resume => unsafe { (*task).resume_receive() }, + // SAFETY: see `task`. + AfterDelivery::Pause => unsafe { (*task).signal_store.pause_receive() }, + AfterDelivery::Park => { + // SAFETY: see `task`. + unsafe { (*task).signal_store.pause_receive() }; + if self.stream.park() { + // SAFETY: see `task`; `park` touched the stream's source, not the task. + unsafe { (*task).poll_ref.unref(bun_io::js_vm_ctx()) }; + } } } } - pub(crate) fn on_stream_cancelled(&mut self) { - let self_ = self; - // Release the Strong ref so the ReadableStream can be GC'd. - // The download may still be in progress, but the callback will - // see readable_stream_ref.get() return null and skip data delivery. - // When the download finishes (has_more == false), deinit() will - // clean up the remaining resources. - self_.readable_stream_ref.deinit(); + fn unpark(&self) { + let task = self.task.get(); + if self.stream.unpark() && !task.is_null() { + // SAFETY: see `task`; reached from a consumer, so the task is still live. + unsafe { (*task).poll_ref.ref_(bun_io::js_vm_ctx()) }; + } + } + + pub(crate) fn on_stream_drained(&self) { + self.unpark(); + let task = self.task.get(); + if !task.is_null() { + // SAFETY: see `task`. + unsafe { (*task).resume_receive() }; + } + } + + pub(crate) fn on_consumer_attached(&self) { + self.unpark(); + } + + /// The parked stream's wrapper was collected: nothing can read the rest. Inside a GC sweep; + /// touches no JS cell. + pub(crate) fn on_stream_collected(&self) { + self.on_stream_cancelled(); + } + + pub(crate) fn on_stream_cancelled(&self) { + // The download may still be in progress, but the callback will see no stream and skip + // delivery. When the download finishes (has_more == false) the task frees this wrapper. + self.stream.release(); // Abort the in-flight HTTP request so the HTTP thread delivers a final // callback with `has_more == false`, which frees the task and this wrapper. // Without this, a server that never sends the terminal chunk would leak both. - let task = core::mem::replace(&mut self_.task, core::ptr::null_mut()); + let task = self.task.replace(core::ptr::null_mut()); if !task.is_null() { // SAFETY: task is live until its own `on_response` frees it on this thread, // which has not happened yet (it would have dropped this wrapper first). @@ -1401,6 +1477,7 @@ impl S3DownloadStreamWrapper { .signal_store .aborted .store(true, core::sync::atomic::Ordering::Relaxed); + (*task).poll_ref.unref(bun_io::js_vm_ctx()); // Wake the HTTP thread so it observes the abort even when the // socket is idle; otherwise the final `has_more == false` // callback never fires and both the task and wrapper leak. @@ -1415,25 +1492,9 @@ impl S3DownloadStreamWrapper { err: Option, opaque_self: *mut c_void, ) { - // SAFETY: opaque_self points to a S3DownloadStreamWrapper allocated in readable_stream - let self_: &mut Self = unsafe { bun_ptr::callback_ctx::(opaque_self) }; - Self::callback(chunk, has_more, err, self_); - } -} - -impl Drop for S3DownloadStreamWrapper { - /// readable_stream_ref / path are freed by their own field Drop. - fn drop(&mut self) { - // Clear the ByteStream's producer handle before `readable_stream_ref` - // drops so the stream never calls back into a freed wrapper. - if let Some(readable) = self.readable_stream_ref.get() { - if let Some(bytes) = readable.ptr.bytes() { - bytes - .parent_const() - .producer - .set(crate::webcore::streams::SourceHandle::None); - } - } + // `opaque_self` is the wrapper allocated in `readable_stream`; handed on as the raw + // pointer so that the terminal callback can free it. + Self::callback(chunk, has_more, err, opaque_self.cast::()); } } @@ -1467,25 +1528,20 @@ pub(crate) fn readable_stream( let readable_value = reader_mut.to_readable_stream(global_this)?; let wrapper = S3DownloadStreamWrapper::new(S3DownloadStreamWrapper { - readable_stream_ref: ReadableStreamStrong::init( - ReadableStream { - ptr: ReadableStreamPtr::Bytes(&raw mut reader_mut.context), - value: readable_value, - }, - global_this, - ), + stream: Default::default(), path: Box::<[u8]>::from(path), global: global_static, - task: core::ptr::null_mut(), + task: Cell::new(core::ptr::null_mut()), }); + // SAFETY: `reader` is the live source made above; `wrapper` the live heap allocation. + unsafe { (*wrapper).stream.hold(&raw mut reader_mut.context) }; reader_mut .producer .set(crate::webcore::streams::SourceHandle::S3DownloadBody( - // SAFETY: `wrapper` is the live heap allocation (write provenance). - unsafe { - bun_ptr::BackRef::from_raw_mut(NonNull::new(wrapper).expect("heap::alloc").as_ptr()) - }, + // SAFETY: `wrapper` is the live heap allocation; cleared from the producer slot before + // it is freed (`ProducerHold::take`). + unsafe { bun_ptr::BackRef::from_raw(wrapper) }, )); let task = download_stream( @@ -1502,7 +1558,7 @@ pub(crate) fn readable_stream( // SAFETY: on the success path `download_stream` only schedules work onto the HTTP // thread; the wrapper is freed via `opaque_callback` on this (main) thread, which // cannot run until we return to the event loop, so `wrapper` is still live here. - unsafe { (*wrapper).task = task }; + unsafe { (*wrapper).task.set(task) }; } Ok(readable_value) } diff --git a/src/runtime/webcore/s3/download_stream.rs b/src/runtime/webcore/s3/download_stream.rs index 9e4cc0c09469..1025d73b98c6 100644 --- a/src/runtime/webcore/s3/download_stream.rs +++ b/src/runtime/webcore/s3/download_stream.rs @@ -70,18 +70,30 @@ impl S3HttpDownloadStreamingTask { self.state.store(state.0, Ordering::Relaxed); } - fn report_progress(&mut self, state: State) { + /// The chunk callback runs JS, which reaches back into this task through the stream + /// wrapper's pointer (`pause_receive`, `poll_ref`). No borrow of `this` may span that call, + /// so this takes the raw pointer and scopes every access to its statement. + /// + /// # Safety + /// `this` is live and exclusively accessed by this thread for the duration of the call + /// (`on_response`). + unsafe fn report_progress(this: *mut Self, state: State) { let has_more = state.has_more(); let failed = match state.status_code() { 200 | 204 | 206 => state.request_error() != 0, _ => true, }; + // SAFETY: fn contract; `callback` and `callback_context` are set once, before the task + // is queued. + let (callback, callback_context) = + unsafe { ((*this).callback, (*this).callback_context.as_ptr().cast()) }; bun_core::scoped_log!( S3, "reportProgres failed: {} has_more: {} len: {}", failed, has_more, - self.reported_response_buffer.list.len() + // SAFETY: fn contract. + unsafe { (*this).reported_response_buffer.list.len() } ); if failed { @@ -92,10 +104,13 @@ impl S3HttpDownloadStreamingTask { let mut code: &[u8] = b"UnknownError"; let mut message: &[u8] = b"an unexpected error has occurred"; let parsed; - if let Some(req_err) = self.request_error { + // SAFETY: fn contract. + if let Some(req_err) = unsafe { (*this).request_error } { code = req_err.name().as_bytes(); } else { - let bytes = self.reported_response_buffer.list.as_slice(); + // SAFETY: fn contract; the buffer is not touched again before the callback + // returns, and `message` is not used after it. + let bytes = unsafe { (*this).reported_response_buffer.list.as_slice() }; if !bytes.is_empty() { message = bytes; } @@ -105,27 +120,26 @@ impl S3HttpDownloadStreamingTask { message = error.message.as_deref().unwrap_or(message); } } - (self.callback)( + callback( &empty, false, Some(S3Error { code, message }), - self.callback_context.as_ptr().cast(), + callback_context, ); return; } // dont report empty chunks if we have more data to read - if !has_more || self.reported_response_buffer.list.len() > 0 { + // SAFETY: fn contract. + if !has_more || unsafe { (*this).reported_response_buffer.list.len() } > 0 { // `core::mem::take` transfers ownership of the buffer, leaving an // empty MutableString behind. - let chunk = core::mem::take(&mut self.reported_response_buffer); - (self.callback)( - &chunk, - has_more, - None, - self.callback_context.as_ptr().cast(), - ); - self.reported_response_buffer.reset(); + // SAFETY: fn contract; the borrow ends with the statement. + let chunk = unsafe { core::mem::take(&mut (*this).reported_response_buffer) }; + callback(&chunk, has_more, None, callback_context); + // SAFETY: fn contract; the callback does not free the task, `on_response` does, + // after this returns. + unsafe { (*this).reported_response_buffer.reset() }; } } @@ -169,8 +183,8 @@ impl S3HttpDownloadStreamingTask { .store(false, Ordering::Relaxed) }; } - // SAFETY: as above; exclusive borrow scoped to the call. - unsafe { (*this).report_progress(state) }; + // SAFETY: as above. + unsafe { Self::report_progress(this, state) }; } /// this function is only called from the http callback in the HTTPThread and returns true if we @@ -247,6 +261,15 @@ impl S3HttpDownloadStreamingTask { ); result.body_into(&mut self.reported_response_buffer.list); + // Only a body that is being delivered can be resumed: the wrapper resumes as JS takes + // chunks. An error body is collected whole before it is reported, with no chunk + // delivered in between, so pausing it would never be undone. + if !is_done + && !wait_until_done + && self.reported_response_buffer.list.len() >= bun_http::signals::BODY_HIGH_WATER_MARK + { + self.signal_store.pause_receive(); + } if should_enqueue { if self.reported_response_buffer.list.is_empty() && !is_done { return false; @@ -354,6 +377,13 @@ impl S3HttpDownloadStreamingTask { } } + /// A consumer took bytes: undo a pause from either side of the hop. + pub(crate) fn resume_receive(&self) { + if self.signal_store.unpause_receive() { + bun_http::http_thread().schedule_receive_resume(self.async_http_id); + } + } + fn release_portable(&mut self) { // SAFETY: `http` is always initialised before the task is scheduled / dropped. let http = unsafe { self.http.assume_init_mut() }; diff --git a/src/runtime/webcore/s3/multipart.rs b/src/runtime/webcore/s3/multipart.rs index 2780c3a36ea2..fb5b8986013c 100644 --- a/src/runtime/webcore/s3/multipart.rs +++ b/src/runtime/webcore/s3/multipart.rs @@ -141,6 +141,9 @@ pub struct MultiPartUpload { pub global_this: GlobalRef, pub(crate) buffered: JsCell, + /// Bytes accepted by `write*` (after encoding): what a streamed `Bun.write`/`writer.end()` + /// resolves with. + pub(crate) uploaded_bytes: Cell, pub path: Box<[u8]>, pub(crate) proxy: Box<[u8]>, @@ -154,7 +157,9 @@ pub struct MultiPartUpload { pub(crate) state: Cell, - pub callback: fn(S3UploadResult, *mut c_void) -> bun_jsc::JsResult<()>, + /// Completion. The upload is passed so `uploaded_bytes` can be read by a callee that no + /// longer holds a ref to it (a `writer()` sink whose JS wrapper was collected). + pub callback: fn(&MultiPartUpload, S3UploadResult, *mut c_void) -> bun_jsc::JsResult<()>, pub(crate) on_writable: Option, pub(crate) callback_context: Cell<*mut c_void>, } @@ -574,7 +579,11 @@ impl MultiPartUpload { } if self.state.get() != State::Finished { let old_state = self.state.replace(State::Finished); - (self.callback)(S3UploadResult::Failure(err), self.callback_context.get())?; + (self.callback)( + self, + S3UploadResult::Failure(err), + self.callback_context.get(), + )?; if old_state == State::MultipartCompleted { // we are a multipart upload so we need to rollback @@ -621,7 +630,7 @@ impl MultiPartUpload { self.state.set(State::Finished); // single file upload no need to commit // The deref must run after the callback: - let r = (self.callback)(S3UploadResult::Success, self.callback_context.get()); + let r = (self.callback)(self, S3UploadResult::Success, self.callback_context.get()); MultiPartUpload::deref_(self.root_ptr()); r } else { @@ -735,15 +744,19 @@ impl MultiPartUpload { } self_.state.set(State::Finished); // The deref must run after the callback: - let r = - (self_.callback)(S3UploadResult::Failure(err), self_.callback_context.get()); + let r = (self_.callback)( + self_, + S3UploadResult::Failure(err), + self_.callback_context.get(), + ); MultiPartUpload::deref_(this); r } S3CommitResult::Success => { self_.state.set(State::Finished); // The deref must run after the callback: - let r = (self_.callback)(S3UploadResult::Success, self_.callback_context.get()); + let r = + (self_.callback)(self_, S3UploadResult::Success, self_.callback_context.get()); MultiPartUpload::deref_(this); r } @@ -1043,6 +1056,7 @@ impl MultiPartUpload { } fn append_chunk(&self, encoding: WriteEncoding, chunk: &[u8]) -> Result<(), AllocError> { + let before = self.buffered.get().size(); self.buffered.with_mut(|buffered| match encoding { WriteEncoding::Bytes => buffered.write(chunk), WriteEncoding::Latin1 => buffered.write_latin1::(chunk), @@ -1052,6 +1066,8 @@ impl MultiPartUpload { buffered.write_utf16(utf16) } })?; + self.uploaded_bytes + .set(self.uploaded_bytes.get() + (self.buffered.get().size() - before) as u64); Ok(()) } diff --git a/src/runtime/webcore/streams.rs b/src/runtime/webcore/streams.rs index 417f9a78e828..e499ab86ce3f 100644 --- a/src/runtime/webcore/streams.rs +++ b/src/runtime/webcore/streams.rs @@ -873,7 +873,7 @@ pub enum SourceHandle { ShellWritable(BackRef), FetchResponseBody(BackRef), ServerRequestBody(crate::server::AnyRequestContext), - S3DownloadBody(BackRef), + S3DownloadBody(BackRef), HTMLRewriter(BackRef), /// `bun:internal-for-testing` only: `ready()` re-enters the stream's /// `on_cancel`, making consumed-during-`signal_drained` re-entrancy @@ -912,10 +912,8 @@ impl SourceHandle { SourceHandle::Subprocess(p) => p.on_close(err), // SAFETY: live backref; cleared before the pointee is freed. SourceHandle::ShellWritable(mut p) => unsafe { p.get_mut() }.on_close(err), - // SAFETY: live backref; cleared before the pointee is freed. - SourceHandle::FetchResponseBody(mut p) => unsafe { p.get_mut() }.on_stream_cancelled(), - // SAFETY: live backref; cleared before the pointee is freed. - SourceHandle::S3DownloadBody(mut p) => unsafe { p.get_mut() }.on_stream_cancelled(), + SourceHandle::FetchResponseBody(p) => p.on_stream_cancelled(), + SourceHandle::S3DownloadBody(p) => p.on_stream_cancelled(), SourceHandle::ServerRequestBody(_) => {} SourceHandle::HTMLRewriter(p) => p.on_close(err), SourceHandle::TestingCancelOnDrain(_) => {} @@ -939,13 +937,12 @@ impl SourceHandle { SourceHandle::FetchResponseBody(p) => p.on_ready(), SourceHandle::ServerRequestBody(any) => any.on_request_body_stream_drained(), SourceHandle::HTMLRewriter(p) => p.on_ready(), + SourceHandle::S3DownloadBody(p) => p.on_stream_drained(), SourceHandle::TestingCancelOnDrain(p) => { p.on_cancel(); } // Remaining variants leave `on_ready` at the trait default (no-op). - SourceHandle::Subprocess(_) - | SourceHandle::ShellWritable(_) - | SourceHandle::S3DownloadBody(_) => {} + SourceHandle::Subprocess(_) | SourceHandle::ShellWritable(_) => {} } } @@ -955,6 +952,7 @@ impl SourceHandle { pub fn consumer_collected(self) { match self { SourceHandle::FetchResponseBody(p) => p.on_body_stream_collected(), + SourceHandle::S3DownloadBody(p) => p.on_stream_collected(), SourceHandle::None | SourceHandle::JSController(_) | SourceHandle::ServerRequestBody(_) @@ -962,7 +960,6 @@ impl SourceHandle { | SourceHandle::FileReader(_) | SourceHandle::Subprocess(_) | SourceHandle::ShellWritable(_) - | SourceHandle::S3DownloadBody(_) | SourceHandle::HTMLRewriter(_) | SourceHandle::TestingCancelOnDrain(_) => {} } @@ -971,6 +968,7 @@ impl SourceHandle { pub fn start(&mut self) { match *self { SourceHandle::FetchResponseBody(p) => p.on_start(), + SourceHandle::S3DownloadBody(p) => p.on_consumer_attached(), // Remaining variants leave `on_start` at the trait default (no-op). SourceHandle::None | SourceHandle::JSController(_) @@ -979,7 +977,6 @@ impl SourceHandle { | SourceHandle::FileReader(_) | SourceHandle::Subprocess(_) | SourceHandle::ShellWritable(_) - | SourceHandle::S3DownloadBody(_) | SourceHandle::HTMLRewriter(_) | SourceHandle::TestingCancelOnDrain(_) => {} } @@ -2143,6 +2140,10 @@ pub struct NetworkSink { pub(crate) upstream_error: jsc::strong::Optional, pub(crate) ended: bool, pub(crate) done: bool, + /// `s3file.writer()`: the box is referenced by the JS wrapper (`finalize`) and by the upload's + /// completion callback; whichever lets go last frees it. 0 = owned elsewhere + /// (`S3UploadStreamWrapper`). + pub(crate) writer_holders: core::cell::Cell, } impl Default for NetworkSink { @@ -2158,6 +2159,7 @@ impl Default for NetworkSink { upstream_error: jsc::strong::Optional::empty(), ended: false, done: false, + writer_holders: core::cell::Cell::new(0), } } } @@ -2215,6 +2217,24 @@ impl NetworkSink { self.detach_writable(); } + /// One of the `writer_holders` is done with the box. + /// + /// # Safety + /// `this` is the live heap box from `writable()`; not used by the caller afterwards. + pub(crate) unsafe fn release_writer_holder(this: *mut NetworkSink) { + // SAFETY: fn contract. + unsafe { + let holders = (*this).writer_holders.get(); + if holders == 0 { + return; + } + (*this).writer_holders.set(holders - 1); + if holders == 1 { + drop(bun_core::heap::take(this)); + } + } + } + fn detach_writable(&mut self) { if let Some(task) = self.task.take() { // task is ref-counted; deref releases our ref @@ -2507,10 +2527,11 @@ impl crate::webcore::sink::JsSinkType for NetworkSink { crate::impl_js_sink_forwarders!(); unsafe fn finalize(this: *mut Self) { - // SAFETY: trait contract — `this` is live, and the inherent `finalize` - // only releases the ref on the separate `MultiPartUpload`, never this - // sink, so the `&mut` scoped to this call stays valid throughout. - unsafe { (*this).finalize() } + // SAFETY: trait contract — `this` is live and not used after this call. + unsafe { + (*this).finalize(); + Self::release_writer_holder(this); + } } fn end_from_js(&mut self, global: &JSGlobalObject) -> bun_sys::Result { Self::end_from_js(self, global) diff --git a/test/js/bun/io/bun-write.test.js b/test/js/bun/io/bun-write.test.js index f02bdf4f8e17..1806305f3200 100644 --- a/test/js/bun/io/bun-write.test.js +++ b/test/js/bun/io/bun-write.test.js @@ -11,6 +11,8 @@ import { tempDir, withoutAggressiveGC, } from "harness"; +import { once } from "node:events"; +import http from "node:http"; import path, { join } from "path"; let i = 0; @@ -900,9 +902,8 @@ int posix_fadvise(int fd, off_t offset, off_t len, int advice) { expect(await Bun.file(dest).text()).toBe("

hi

"); }); - // Bun.write owns the Locked body (on_receive_value + retargeted task); - // clone()'s tee must see that and yield a used body instead of dispatching - // the producer callbacks with Bun.write's task as ctx. + // Bun.write is reading the body, so it is used: clone() throws (as it does after .text()), and + // the write is unaffected. it("Bun.write(path, HTMLRewriter.transform(resp)) survives clone() while a handler is suspended", async () => { using dir = tempDir("bun-write-htmlrewriter-clone", {}); const dest = join(String(dir), "out.html"); @@ -919,13 +920,229 @@ int posix_fadvise(int fd, off_t offset, off_t len, int advice) { .transform(new Response("

y

")); const write = Bun.write(dest, out); await suspended; - const clone = out.clone(); - expect(clone).toBeInstanceOf(Response); + expect(out.bodyUsed).toBe(true); + expect(() => out.clone()).toThrow(expect.objectContaining({ code: "ERR_BODY_ALREADY_USED" })); openGate(); - await write; + expect(await write).toBe(8); expect(await Bun.file(dest).text()).toBe("

x

"); }); + describe("Bun.write(path, response) streams the body to the file", () => { + const CHUNK = 64 * 1024; + const COUNT = 64; // 4 MiB + // Serves COUNT chunks; with `gate`, the second half only once it opens. + // node:http rather than Bun.serve: no server-side Response objects to muddy a Response count. + async function origin(gate) { + const payload = Buffer.alloc(CHUNK, "a"); + const server = http.createServer(async (req, res) => { + // The client may go away mid-body (a failed write cancels the source, or the test ends). + res.on("error", () => {}); + if (req.url.endsWith("/small")) return res.end(payload.subarray(0, 1000)); + res.writeHead(200, { "content-length": String(CHUNK * COUNT) }); + for (let i = 0; i < COUNT && !res.destroyed; i++) { + if (gate && i === COUNT / 2) await gate; + if (!res.write(payload)) await once(res, "drain").catch(() => {}); + } + if (!res.destroyed) res.end(); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + return { + url: new URL(`http://127.0.0.1:${server.address().port}/`), + [Symbol.asyncDispose]: () => new Promise(resolve => server.closeAllConnections() || server.close(resolve)), + }; + } + + it("resolves with the byte count and replaces a longer existing file", async () => { + using dir = tempDir("bun-write-response-stream", { "out.bin": Buffer.alloc(CHUNK * COUNT + 12345, "z") }); + await using server = await origin(); + const dest = join(String(dir), "out.bin"); + expect(await Bun.write(dest, await fetch(server.url))).toBe(CHUNK * COUNT); + expect(fs.statSync(dest).size).toBe(CHUNK * COUNT); + }); + + it("writes as the body arrives, and a collected Response does not stop it", async () => { + using dir = tempDir("bun-write-response-collected", {}); + const dest = join(String(dir), "deep", "er", "out.bin"); + const { promise: gate, resolve: openGate } = Promise.withResolvers(); + await using server = await origin(gate); + fs.mkdirSync(path.dirname(dest), { recursive: true }); + const onDisk = new Promise(resolve => { + const watcher = fs.watch(path.dirname(dest), () => { + if (fs.existsSync(dest) && fs.statSync(dest).size > 0) { + watcher.close(); + resolve(); + } + }); + }); + // Its own frame: after it returns only Bun.write refers to the body. + let response; + async function start() { + const res = await fetch(server.url); + response = new WeakRef(res); + return Bun.write(dest, res); + } + const written = start(); + // The first half is on disk before the second half is sent. Before, nothing was written + // until the whole body had been collected in memory. + await onDisk; + // One full collection per event-loop turn (an fs round trip) until the Response is gone. + do { + await fs.promises.stat(dest); + Bun.gc(true); + } while (response.deref()); + openGate(); + // Before #40278 was fixed this never settled once the Response had been collected. + expect(await written).toBe(CHUNK * COUNT); + expect(fs.statSync(dest).size).toBe(CHUNK * COUNT); + }); + + it("a body whose stream was already touched", async () => { + using dir = tempDir("bun-write-response-touched", {}); + await using server = await origin(); + const res = await fetch(server.url); + expect(res.body).toBeInstanceOf(ReadableStream); + expect(await Bun.write(join(String(dir), "out.bin"), res)).toBe(CHUNK * COUNT); + expect(res.bodyUsed).toBe(true); + }); + + // https://github.com/oven-sh/bun/issues/13237: this never settled. + it("a Response around a JS ReadableStream, counting string chunks by their UTF-8 length", async () => { + using dir = tempDir("bun-write-response-js-stream", {}); + const dest = join(String(dir), "out.txt"); + const stream = new ReadableStream({ + start(ctrl) { + ctrl.enqueue("héllo "); + ctrl.enqueue(new TextEncoder().encode("stream ")); + ctrl.enqueue("\u{1f600}"); + ctrl.close(); + }, + }); + const expected = Buffer.from("héllo stream \u{1f600}"); + expect(await Bun.write(dest, new Response(stream))).toBe(expected.length); + expect(Buffer.from(await Bun.file(dest).arrayBuffer())).toEqual(expected); + }); + + // /dev/full: every write fails with ENOSPC. + it.skipIf(process.platform !== "linux")("rejects with the write error, for each kind of body", async () => { + await using server = await origin(); + const streamed = await fetch(server.url); + // A body that is all here behind an untouched `.body` stream is written as a blob. + const arrived = new Response(await (await fetch(server.url + "small")).blob()); + expect(arrived.body).toBeInstanceOf(ReadableStream); + const js = new Response( + new ReadableStream({ + start(ctrl) { + ctrl.enqueue(new Uint8Array(1000)); + ctrl.close(); + }, + }), + ); + for (const res of [streamed, arrived, js]) { + await expect(Bun.write("/dev/full", res)).rejects.toThrow(expect.objectContaining({ code: "ENOSPC" })); + } + }); + + // https://github.com/oven-sh/bun/issues/31681: these wrote "[object ReadableStream]". + it("a bare ReadableStream: res.body, a JS stream, file.write(stream)", async () => { + using dir = tempDir("bun-write-readable-stream", {}); + await using server = await origin(); + const viaBody = join(String(dir), "body.bin"); + expect(await Bun.write(viaBody, (await fetch(server.url)).body)).toBe(CHUNK * COUNT); + expect(fs.statSync(viaBody).size).toBe(CHUNK * COUNT); + + const viaJs = join(String(dir), "js.txt"); + const js = () => + new ReadableStream({ + start(ctrl) { + ctrl.enqueue("one "); + ctrl.enqueue(new TextEncoder().encode("two")); + ctrl.close(); + }, + }); + expect(await Bun.write(viaJs, js())).toBe(7); + expect(await Bun.file(viaJs).text()).toBe("one two"); + expect(await Bun.file(join(String(dir), "file-write.txt")).write(js())).toBe(7); + expect(await Bun.file(join(String(dir), "file-write.txt")).text()).toBe("one two"); + + const locked = js(); + locked.getReader(); + await expect(Bun.write(viaJs, locked)).rejects.toThrow( + expect.objectContaining({ code: "ERR_BODY_ALREADY_USED" }), + ); + expect(await Bun.file(viaJs).text()).toBe("one two"); + }); + + it("a Request body inside Bun.serve", async () => { + using dir = tempDir("bun-write-request-stream", {}); + const dest = join(String(dir), "upload.bin"); + await using server = Bun.serve({ + port: 0, + async fetch(req) { + return new Response(String(await Bun.write(dest, req))); + }, + }); + const body = Buffer.alloc(3 * CHUNK * COUNT, "b"); + const res = await fetch(server.url, { method: "POST", body }); + expect({ written: Number(await res.text()), size: fs.statSync(dest).size }).toEqual({ + written: body.length, + size: body.length, + }); + }); + + it("rejects with the network error when the body is cut short", async () => { + using dir = tempDir("bun-write-response-truncated", {}); + using listener = Bun.listen({ + port: 0, + hostname: "127.0.0.1", + socket: { + data(socket) { + socket.write("HTTP/1.1 200 OK\r\nContent-Length: 1000000\r\n\r\n" + Buffer.alloc(1000, "a").toString()); + socket.flush(); + socket.end(); + }, + }, + }); + const res = await fetch(`http://127.0.0.1:${listener.port}/`); + await expect(Bun.write(join(String(dir), "out.bin"), res)).rejects.toThrow( + expect.objectContaining({ code: "ECONNRESET" }), + ); + }); + + it("rejects a body that was already used, and createPath: false into a missing directory", async () => { + using dir = tempDir("bun-write-response-rejects", {}); + await using server = await origin(); + const used = await fetch(server.url); + await used.arrayBuffer(); + await expect(Bun.write(join(String(dir), "a"), used)).rejects.toThrow( + expect.objectContaining({ code: "ERR_BODY_ALREADY_USED" }), + ); + const reading = await fetch(server.url); + const reader = reading.body.getReader(); + await expect(Bun.write(join(String(dir), "b"), reading)).rejects.toThrow( + expect.objectContaining({ code: "ERR_BODY_ALREADY_USED" }), + ); + reader.releaseLock(); + // Also when the whole body is already here. + const small = new Response("hello"); + const smallReader = small.body.getReader(); + await expect(Bun.write(join(String(dir), "s"), small)).rejects.toThrow( + expect.objectContaining({ code: "ERR_BODY_ALREADY_USED" }), + ); + expect(await smallReader.read().then(r => r.value.byteLength)).toBe(5); + // A destination that cannot be opened (the directory itself) leaves the body usable. + const retry = await fetch(server.url); + await expect(Bun.write(String(dir), retry)).rejects.toThrow( + expect.objectContaining({ code: "EISDIR", syscall: "open" }), + ); + expect(retry.bodyUsed).toBe(false); + expect((await retry.arrayBuffer()).byteLength).toBe(CHUNK * COUNT); + await expect( + Bun.write(join(String(dir), "missing", "c"), await fetch(server.url), { createPath: false }), + ).rejects.toThrow(expect.objectContaining({ code: "ENOENT" })); + }); + }); + it("BunFile.name survives concurrent write() calls + GC", async () => { using dir = tempDir("bun-file-name-concurrent-write-gc", {}); const filePath = join(String(dir), "out.txt"); diff --git a/test/js/web/fetch/body-mixin-errors.test.ts b/test/js/web/fetch/body-mixin-errors.test.ts index ac5d747effc9..9dbf3953c93e 100644 --- a/test/js/web/fetch/body-mixin-errors.test.ts +++ b/test/js/web/fetch/body-mixin-errors.test.ts @@ -112,6 +112,68 @@ describe("body-mixin-errors", () => { }); }); + // The body can also have failed before anything reads it: here the whole response arrives at + // once and its body does not decode, so the Response is created with the failure already in + // hand. `.body` then has to be the body's stream all the same: the same stream each time, read + // once it counts as used, and the readers see "already used" afterwards instead of the + // network error again. + async function withUndecodableBodyServer(fn: (url: string) => Promise): Promise { + const server = net.createServer(socket => { + socket.resume(); + socket.end( + "HTTP/1.1 200 OK\r\nContent-Encoding: gzip\r\nContent-Length: 16\r\nConnection: close\r\n\r\nthis is not gzip", + ); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const { port } = server.address() as net.AddressInfo; + try { + return await fn(`http://127.0.0.1:${port}/`); + } finally { + await new Promise(r => server.close(() => r())); + } + } + + it.concurrent("fetch: .body of a body that failed before it was read is still the body", async () => { + await withUndecodableBodyServer(async url => { + const res = await fetch(url); + const body = res.body!; + expect(res.body).toBe(body); + expect(res.bodyUsed).toBe(false); + + let firstErr: unknown; + await body + .getReader() + .read() + .catch(e => (firstErr = e)); + expect(firstErr).toBeInstanceOf(TypeError); + expect(res.bodyUsed).toBe(true); + + let secondErr: unknown; + await res.text().catch(e => (secondErr = e)); + expectBodyAlreadyUsed(secondErr); + }); + }); + + // textStream() is one shot: handing it out uses the body up, failed or not. + it.concurrent("fetch: textStream() of a body that failed before it was read uses the body up", async () => { + await withUndecodableBodyServer(async url => { + const res = await fetch(url); + expect(res.bodyUsed).toBe(false); + + let firstErr: unknown; + await res + .textStream() + .getReader() + .read() + .catch(e => (firstErr = e)); + expect(firstErr).toBeInstanceOf(TypeError); + expect(res.bodyUsed).toBe(true); + + expect(() => res.textStream()).toThrow(TypeError); + }); + }); + it.concurrent.each(["arrayBuffer", "bytes", "blob", "json"] as const)( "fetch: truncated body %s() marks body used", async method => { diff --git a/test/js/web/fetch/body.test.ts b/test/js/web/fetch/body.test.ts index fbdf0ce7553b..65ba890247a9 100644 --- a/test/js/web/fetch/body.test.ts +++ b/test/js/web/fetch/body.test.ts @@ -1415,11 +1415,11 @@ describe.concurrent("a fetch() Response that cannot have a body", () => { expect(await response.text()).toBe(""); }); - test("content that arrives after a 205 resolved is drained, so the process can exit", async () => { + test("content still arriving after a 205 resolved does not hold the process", async () => { // The server sends 3 of the 5 declared bytes with the head and the other 2 // only once fetch() has resolved. Nothing but the fetch refs the event loop, - // so the process only exits if the fetch takes those 2 bytes off the socket - // instead of keeping them for a body reader that cannot exist. + // so the process only exits if the fetch stops waiting for a body no reader + // can exist for: it closes the connection instead (the server sees the reset). await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -1430,6 +1430,7 @@ describe.concurrent("a fetch() Response that cannot have a body", () => { const server = net.createServer(socket => { upstream = socket; socket.unref(); + socket.on("error", () => {}); socket.once("data", () => { socket.write("HTTP/1.1 205 Reset Content\\r\\nContent-Length: 5\\r\\n\\r\\nhel"); }); diff --git a/test/js/web/fetch/fetch-backpressure.test.ts b/test/js/web/fetch/fetch-backpressure.test.ts index 564b15c206b5..babca0edde63 100644 --- a/test/js/web/fetch/fetch-backpressure.test.ts +++ b/test/js/web/fetch/fetch-backpressure.test.ts @@ -1,12 +1,17 @@ // Receive-side backpressure: a stalled `res.body.getReader()` must stop the // HTTP thread from buffering the entire response in memory. +import { S3Client } from "bun"; import { describe, expect, test } from "bun:test"; -import { bunEnv, bunExe, isASAN, isDebug, isWindows, tls } from "harness"; +import { bunEnv, bunExe, isASAN, isDebug, isWindows, tempDir, tls } from "harness"; import { randomBytes } from "node:crypto"; import { once } from "node:events"; +import { statSync } from "node:fs"; +import { stat } from "node:fs/promises"; import { createServer } from "node:http"; import { createSecureServer } from "node:http2"; import { createServer as createHttpsServer } from "node:https"; +import { createServer as createTcpServer } from "node:net"; +import { join } from "node:path"; import { Readable, Writable } from "node:stream"; import { pipeline } from "node:stream/promises"; import { gzipSync } from "node:zlib"; @@ -443,26 +448,33 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" ["end (FIN)", (s: import("bun").Socket) => s.end()], ] as const) { test(`peer ${name} while receive is paused rejects the body`, async () => { - const { promise, resolve } = Promise.withResolvers(); + // Declares far more than it will send and writes until the kernel stops taking it, which, + // with an untouched body, is once the client holds the high-water mark and has paused. + const declared = 1 << 30; + const payload = Buffer.alloc(CHUNK, 65); + const blocked = Promise.withResolvers(); + let sent = 0; + const push = (s: import("bun").Socket) => { + while (sent < declared) { + const n = s.write(payload); + sent += Math.max(n, 0); + if (n < payload.length) return void (sent > 4 * CHUNK && blocked.resolve(s)); + } + }; using listener = Bun.listen({ port: 0, hostname: "127.0.0.1", socket: { open(s) { - // Declared length far exceeds what is sent, so the client parks - // in the body stage (and pauses) right after this first chunk. - s.write(`HTTP/1.1 200 OK\r\nContent-Length: ${TOTAL}\r\n\r\n` + Buffer.alloc(CHUNK, 65).toString()); - s.flush(); - resolve(s); + s.write(`HTTP/1.1 200 OK\r\nContent-Length: ${declared}\r\n\r\n`); + push(s); }, + drain: push, data() {}, }, }); - // By the time fetch() resolves, the first body chunk was delivered with - // more expected, so the transport is paused; nothing re-arms it until - // the body is pulled. const res = await fetch(`http://127.0.0.1:${listener.port}/`); - kill(await promise); + kill(await blocked.promise); const reader = res.body!.getReader(); let total = 0; const err = await (async () => { @@ -471,7 +483,7 @@ describe.concurrent("fetch() receive backpressure — streaming consumer shapes" () => null, e => e, ); - expect({ code: err?.code, partial: total < TOTAL }).toEqual({ code: "ECONNRESET", partial: true }); + expect({ code: err?.code, partial: total < declared }).toEqual({ code: "ECONNRESET", partial: true }); }); } @@ -504,16 +516,26 @@ async function serveUntilBlocked() { let sent = 0; let closed = 0; const payload = Buffer.alloc(CHUNK, 65); + // The kernel stopped taking writes with more than the client's high-water mark outstanding: + // from here only a reader can make room. + const blocked = Promise.withResolvers(); + const firstClosed = Promise.withResolvers(); const srv = createServer((_req, res) => { res.on("error", () => {}); - res.on("close", () => closed++); + res.on("close", () => { + closed++; + firstClosed.resolve(); + }); res.flushHeaders(); let i = 0; const push = () => { while (i < BIG && !res.destroyed) { i++; sent += CHUNK; - if (!res.write(payload)) return void res.once("drain", push); + if (!res.write(payload)) { + if (i > 8) blocked.resolve(); + return void res.once("drain", push); + } } if (!res.destroyed) res.end(); }; @@ -538,6 +560,8 @@ async function serveUntilBlocked() { async untilClosed() { while (closed === 0) await Bun.sleep(5); }, + blocked: blocked.promise, + closed: firstClosed.promise, [Symbol.asyncDispose]: () => { srv.closeAllConnections(); return new Promise(r => srv.close(() => r(undefined))); @@ -671,6 +695,363 @@ describe("fetch() receive backpressure — body stream nothing is reading", () = } }); +// A Response whose body nothing ever touches (looked at for its status, then forgotten). The +// rule above applies to it as well: its body is received up to the mark, so a short one +// completes and its connection goes back to the pool while the Response is still around, and a +// long one leaves the transport paused until the Response is collected, at which point its fetch +// is aborted like an abandoned stream's. + +type Framing = "content-length" | "chunked" | "close-delimited"; + +// A raw HTTP/1.1 origin, so that each test says exactly how its bodies are framed. Every +// response is `length` bytes of body, written as fast as the socket takes them. With `holdTail` +// the last CHUNK of every body is held back until `finishHeld()`, which also stops holding. +// The origin never ends a body by closing, so every close it sees is the client's. +async function rawOrigin(framing: Framing, length: number, holdTail = false) { + const payload = Buffer.alloc(CHUNK, 65); + const frame = + framing === "chunked" + ? Buffer.concat([Buffer.from(`${CHUNK.toString(16)}\r\n`), payload, Buffer.from("\r\n")]) + : payload; + const tail = framing === "chunked" ? Buffer.concat([frame, Buffer.from("0\r\n\r\n")]) : frame; + const head = + framing === "content-length" + ? `HTTP/1.1 200 OK\r\nContent-Length: ${length}\r\n\r\n` + : framing === "chunked" + ? "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n" + : "HTTP/1.1 200 OK\r\nConnection: close\r\n\r\n"; + + let connections = 0; + let closed = 0; + let holding = holdTail; + const held: (() => void)[] = []; + const sockets = new Set(); + const closeWaiters: [number, () => void][] = []; + let requested = 0; + const requestWaiters: [number, () => void][] = []; + + function respond(socket: import("node:net").Socket) { + socket.write(head); + let left = length / CHUNK; + const push = () => { + while (left > 1 && !socket.destroyed) { + left--; + if (!socket.write(frame)) return void socket.once("drain", push); + } + const end = () => void (socket.destroyed || socket.write(tail)); + if (holding) held.push(end); + else end(); + }; + push(); + } + + const srv = createTcpServer(socket => { + connections++; + sockets.add(socket); + socket.on("error", () => {}); + socket.on("close", () => { + closed++; + sockets.delete(socket); + for (const [n, resolve] of closeWaiters) if (closed >= n) resolve(); + }); + // A pooled connection carries one request after another. + let pending = ""; + socket.on("data", data => { + pending += data.toString("latin1"); + for (let end; (end = pending.indexOf("\r\n\r\n")) !== -1; ) { + pending = pending.slice(end + 4); + requested++; + for (const [n, resolve] of requestWaiters) if (requested >= n) resolve(); + respond(socket); + } + }); + }); + srv.listen(0, "127.0.0.1"); + await once(srv, "listening"); + const { port } = srv.address() as import("node:net").AddressInfo; + return { + url: `http://127.0.0.1:${port}/`, + connections: () => connections, + closed: () => closed, + closedAtLeast: (n: number) => + closed >= n ? Promise.resolve() : new Promise(resolve => closeWaiters.push([n, resolve])), + requests: (n: number) => + requested >= n ? Promise.resolve() : new Promise(resolve => requestWaiters.push([n, resolve])), + finishHeld() { + holding = false; + for (const end of held.splice(0)) end(); + }, + [Symbol.asyncDispose]: () => { + for (const socket of sockets) socket.destroy(); + return new Promise(resolve => srv.close(() => resolve())); + }, + }; +} + +// One full collection per event-loop turn (an fs round trip, not a timer) until `event` settles. +// What `event` waits for is several hops away from the collection itself: the finalizer runs in +// the sweep, the abort is a message to the HTTP thread, and the origin sees the close on its own +// socket. Loop on the event, not on a count of collections. +async function collectUntil(event: Promise): Promise { + let settled = false; + const result = event.finally(() => (settled = true)); + while (!settled) { + await stat(import.meta.path); + Bun.gc(true); + } + return result; +} + +// Sequential on purpose, as above: these watch connections and closes on their own origin. +describe("fetch() receive backpressure — a Response whose body nothing touches", () => { + const N = 4; + + test("a long body, Response held untouched: the process is not held and the body stays where it is", async () => { + await using server = await serveUntilBlocked(); + 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 }); + expect(server.sent()).toBeLessThan(BODY); + }); + + for (const framing of ["content-length", "chunked", "close-delimited"] as Framing[]) { + test(`a long ${framing} body, Response collected: its fetch is aborted`, async () => { + await using origin = await rawOrigin(framing, TOTAL); + // Its own frame, so that nothing on this one still refers to a response afterwards. + async function abandonOne() { + expect((await fetch(origin.url)).status).toBe(200); + } + for (let i = 0; i < N; i++) await abandonOne(); + await collectUntil(origin.closedAtLeast(N)); + }); + } + + // 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 () => { + // Each body's tail is held back until its Response exists, so every body is still underway + // when fetch() resolves, as it is over a real network. + await using origin = await rawOrigin(framing, 2 * CHUNK, true); + const responses: Response[] = []; + for (let i = 0; i < N; i++) responses.push(await fetch(origin.url)); + // Every body is underway, so no connection was free for the next request. + expect(origin.connections()).toBe(N); + + origin.finishHeld(); + // The held bodies complete on their own and give their connections back. One request after + // another from here needs at most one more connection (the first can leave before the tails + // were taken); before, every one of them did, since each held body pinned its connection. + for (let i = 0; i < N; i++) expect((await (await fetch(origin.url)).arrayBuffer()).byteLength).toBe(2 * CHUNK); + expect(origin.connections() - N).toBeLessThanOrEqual(1); + expect({ closed: origin.closed(), held: responses.length }).toEqual({ closed: 0, held: N }); + }); + } + + // The boundary of the abort above: a consumer that waits for the whole body (`.text()` through + // a promise, `Bun.write()` through a native callback) may be all that is left of a Response. + // Its body still has to arrive. + const wholeBodyConsumers: [string, (res: Response, dir: string) => Promise][] = [ + ["res.text()", res => res.text().then(text => text.length)], + ["Bun.write(file, res)", (res, dir) => Bun.write(join(dir, "body"), res)], + ]; + for (const [name, consume] of wholeBodyConsumers) { + test(`a Response collected while ${name} waits for its body: the body still arrives`, async () => { + using dir = tempDir("fetch-collected-while-consumed", {}); + await using origin = await rawOrigin("content-length", 2 * CHUNK, true); + let response!: WeakRef; + // Its own frame: once it returns, the consumer's promise is all that is held. + async function start() { + const res = await fetch(origin.url); + response = new WeakRef(res); + return consume(res, String(dir)); + } + const received = start(); + await origin.requests(1); + // One full collection per event-loop turn (an fs round trip) until the Response is gone. + do { + await stat(import.meta.path); + Bun.gc(true); + } while (!response || response.deref()); + origin.finishHeld(); + // Before, the collection let go of the body instead, and Bun.write() never settled. + expect({ received: await received, closed: origin.closed() }).toEqual({ received: 2 * CHUNK, closed: 0 }); + }); + } +}); + +// S3 downloads go through the same HTTP client with their own body producer +// (S3DownloadStreamWrapper). The same rule applies: a reader that stalls pauses the transport, +// an unread stream does not hold the process, and a collected one aborts the download. +describe("S3 receive backpressure", () => { + // A GET-only fake bucket: every object is BODY bytes written as fast as the socket takes them. + async function fakeBucket() { + const server = await serveUntilBlocked(); + const s3 = new S3Client({ accessKeyId: "test", secretAccessKey: "test", endpoint: server.url, bucket: "b" }); + return Object.assign(server, { s3 }); + } + + test("a reader that stalls and comes back drains the body; cancel() closes the connection", async () => { + await using bucket = await fakeBucket(); + const reader = bucket.s3.file("big").stream().getReader(); + let got = (await reader.read()).value!.byteLength; + await bucket.blocked; + while (got < 64 * CHUNK) got += (await reader.read()).value!.byteLength; + await reader.cancel(); + await bucket.closed; + }); + + test("Bun.write(file, s3file) streams to disk with the byte count", async () => { + using dir = tempDir("s3-to-file", {}); + await using server = await serve("h1"); + const s3 = new S3Client({ accessKeyId: "test", secretAccessKey: "test", endpoint: server.url, bucket: "b" }); + const dest = join(String(dir), "out.bin"); + expect(await Bun.write(dest, s3.file("k"))).toBe(TOTAL); + expect(statSync(dest).size).toBe(TOTAL); + }); + + // A paused, unread stream releases the loop: the process exits with most of the body unsent. + // Without the pause it would either read all of BODY first or never exit. + test("an unread S3 stream does not hold the process", async () => { + await using bucket = await fakeBucket(); + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + `const s3 = new Bun.S3Client({ accessKeyId: "t", secretAccessKey: "t", endpoint: ${JSON.stringify(bucket.url)}, bucket: "b" }); + globalThis.keep = s3.file("k").stream(); + const r = globalThis.keep.getReader(); await r.read(); r.releaseLock();`, + ], + 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 }); + expect(bucket.sent()).toBeLessThan(BODY); + }); + + test("a collected S3 stream aborts the download", async () => { + await using bucket = await fakeBucket(); + // Its own frame: after it returns nothing refers to the stream. + await (async () => { + const r = bucket.s3.file("k").stream().getReader(); + await r.read(); + r.releaseLock(); + })(); + await bucket.blocked; + await collectUntil(bucket.closed); + expect(bucket.sent()).toBeLessThan(BODY); + }); + + // An error body is collected whole for the error message; the mark must not pause it. + test("a non-2xx response larger than the mark rejects instead of stalling", async () => { + const message = Buffer.alloc(400 * 1024, "e").toString(); + await using bucket = Bun.serve({ + port: 0, + fetch: () => + new Response(`AccessDenied${message}`, { status: 403 }), + }); + const s3 = new S3Client({ accessKeyId: "t", secretAccessKey: "t", endpoint: bucket.url.href, bucket: "b" }); + await expect(s3.file("denied").stream().getReader().read()).rejects.toThrow( + expect.objectContaining({ code: "AccessDenied" }), + ); + }); + + // The other direction: a fetch body uploaded to S3. The multipart sink's queue back-pressures + // the fetch, and both `Bun.write(s3file, res)` and `s3file.writer()` resolve with the bytes sent. + async function fakeUploadBucket(holdParts?: Promise) { + let uploaded = 0; + let parts = 0; + let completed = 0; + const firstPart = Promise.withResolvers(); + const server = Bun.serve({ + port: 0, + async fetch(req) { + const url = new URL(req.url); + if (req.method === "POST" && url.searchParams.has("uploads")) + return new Response("u"); + if (req.method === "PUT") { + const body = await req.arrayBuffer(); + firstPart.resolve(); + await holdParts; + uploaded += body.byteLength; + return new Response("", { headers: { etag: `"e${++parts}"` } }); + } + if (req.method === "POST" && url.searchParams.has("uploadId")) { + await req.text(); + completed++; + return new Response( + "bke", + ); + } + if (req.method === "DELETE") return new Response(null, { status: 204 }); + return new Response("", { status: 400 }); + }, + }); + const s3 = new S3Client({ accessKeyId: "t", secretAccessKey: "t", endpoint: server.url.href, bucket: "b" }); + return Object.assign(server, { + s3, + firstPart: firstPart.promise, + uploaded: () => uploaded, + completed: () => completed, + }); + } + + test("fetch → Bun.write(s3file, res) is paced by the part uploads, and an aborted source commits nothing", async () => { + const hold = Promise.withResolvers(); + await using origin = await serveUntilBlocked(); + await using bucket = await fakeUploadBucket(hold.promise); + const abort = new AbortController(); + // partSize 5 MiB × queueSize 1: with the first part held, the sink fills and the origin has + // to stop long before its 1 GiB is out. + const written = Bun.write( + bucket.s3.file("up", { partSize: 5 * 1024 * 1024, queueSize: 1 }), + await fetch(origin.url, { signal: abort.signal }), + ); + await bucket.firstPart; + await origin.blocked; + expect(origin.sent()).toBeLessThan(BODY); + // Aborting the source fails the upload: nothing is committed. + abort.abort(); + hold.resolve(); + await expect(written).rejects.toThrow(expect.objectContaining({ name: "AbortError" })); + expect(bucket.completed()).toBe(0); + }); + + test("Bun.write(s3file, res) resolves with the byte count", async () => { + await using origin = await serve("h1"); + await using bucket = await fakeUploadBucket(); + expect(await Bun.write(bucket.s3.file("up"), await fetch(origin.url))).toBe(TOTAL); + expect(bucket.uploaded()).toBe(TOTAL); + }); + + test("s3file.writer().end() resolves with the byte count, also once the writer is collected", async () => { + const hold = Promise.withResolvers(); + await using bucket = await fakeUploadBucket(hold.promise); + const collected = Promise.withResolvers(); + const registry = new FinalizationRegistry(() => collected.resolve()); + // Its own frame: once it returns, only the pending end() refers to the upload. + function start() { + const writer = bucket.s3.file("up2").writer(); + registry.register(writer, null); + writer.write(Buffer.alloc(1000, 1)); + writer.write("héllo"); + return writer.end(); + } + const ended = start(); + await bucket.firstPart; + // The count must not depend on the writer object: it is collectable while the PUT is out. + await collectUntil(collected.promise); + hold.resolve(); + expect(await ended).toBe(1006); + }); +}); + describe.concurrent("fetch() receive backpressure — a body nothing waits for does not hold the process", () => { // The client keeps the unread body reachable and has nothing else to do: it has to exit // on its own. Before, it stayed alive draining the body into memory, or, when the pause diff --git a/test/js/web/fetch/fetch-response-finalizer-sweep.test.ts b/test/js/web/fetch/fetch-response-finalizer-sweep.test.ts index 9df2a9b32fd9..6ba1c1eb3ce9 100644 --- a/test/js/web/fetch/fetch-response-finalizer-sweep.test.ts +++ b/test/js/web/fetch/fetch-response-finalizer-sweep.test.ts @@ -1,5 +1,5 @@ // fetch(): the JSResponse Weak finalizer (WeakBlock::sweep) reaches -// FetchTasklet::ignore_remaining_response_body. That path used to call +// FetchTasklet::abandon_response_body. That path used to call // ResumableSink::detach_js(), which writes the wrapper's cached ondrain/ // oncancel/stream slots via generated *SetCachedValue helpers. Those helpers // uncheckedDowncast(...) -> JSCell::classInfo(), and