Skip to content
12 changes: 12 additions & 0 deletions src/runtime/webcore/FileReader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -449,6 +449,18 @@ impl FileReader {
unsafe { (*self.parent()).increment_count() };
}
}
#[cfg(windows)]
{
// Non-lazy fromPipe path (Bun.spawn stdout/stderr): hold a
// ref across the pending uv_read_start so the source is not
// finalized while IOCP has a read queued on it.
if !self.started.get() && self.reader().source.is_some() && !self.reader().is_done()
{
self.waiting_for_on_reader_done.set(true);
// SAFETY: see `parent()`.
unsafe { (*self.parent()).increment_count() };
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

#[cfg(unix)]
Expand Down
145 changes: 110 additions & 35 deletions src/runtime/webcore/ReadableStream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -577,6 +577,10 @@ pub trait SourceContext: Sized {
fn js_on_drain_callback_set_cached(this: JSValue, global: &JSGlobalObject, value: JSValue);
/// `js_${NAME}InternalReadableStreamSource::on_drain_callback_get_cached`
fn js_on_drain_callback_get_cached(this: JSValue) -> Option<JSValue>;
/// `js_${NAME}InternalReadableStreamSource::on_close_callback_set_cached`
fn js_on_close_callback_set_cached(this: JSValue, global: &JSGlobalObject, value: JSValue);
/// `js_${NAME}InternalReadableStreamSource::on_close_callback_get_cached`
fn js_on_close_callback_get_cached(this: JSValue) -> Option<JSValue>;

fn on_start(&mut self) -> streams::Start;
fn on_pull(&mut self, buf: &mut [u8], view: JSValue) -> streams::Result;
Expand Down Expand Up @@ -646,7 +650,6 @@ pub struct NewSource<C: SourceContext> {
/// owned/freed here). The JS path stores
/// `on_js_close` and leaves this `None` — see [`Self::on_close`].
pub close_ctx: Option<NonNull<c_void>>,
pub close_jsvalue: bun_jsc::strong::Optional,
/// R-2: cleared via `&self` from `FetchTasklet::clear_stream_cancel_handler`
/// (through `ByteStream::parent_const`), so interior-mutable.
pub cancel_handler: Cell<Option<fn(Option<*mut c_void>)>>,
Expand All @@ -655,9 +658,15 @@ pub struct NewSource<C: SourceContext> {
// `start()` from a fresh `&JSGlobalObject`; `BackRef` gives a safe `Deref`
// projection without propagating a lifetime parameter into FFI codegen.
pub global_this: Option<bun_ptr::BackRef<JSGlobalObject>>,
// SAFETY: this is the self-wrapper JSValue (points at the JSCell that owns this m_ctx).
// Kept alive by the wrapper itself; zeroed in finalize() before sweep.
pub this_jsvalue: JSValue,
/// Back-reference to the owning `JS{Blob,Bytes,File}InternalReadableStreamSource`
/// wrapper. Starts `Weak` (set in [`Self::to_readable_stream`]), is
/// [`JsRef::upgrade`]d to `Strong` in [`Self::increment_count`] while a
/// native I/O ref is held (FileReader `waiting_for_on_reader_done`), and
/// [`JsRef::downgrade`]d back to `Weak` in [`Self::decrement_count`] when
/// only the wrapper's own ref remains. [`Self::finalize`] flips it to
/// `Finalized` so [`Self::on_js_close`] reads `None` instead of a
/// dead-but-unswept cell.
pub this_jsvalue: jsc::JsRef,
/// R-2: written by `&self` context methods (`ByteStream::to_any_blob`,
/// `ByteBlobLoader::to_any_blob`) via `parent_const()`, so interior-mutable.
pub is_closed: Cell<bool>,
Expand All @@ -672,11 +681,10 @@ impl<C: SourceContext + Default> Default for NewSource<C> {
pending_err: None,
close_handler: None,
close_ctx: None,
close_jsvalue: bun_jsc::strong::Optional::empty(),
cancel_handler: Cell::new(None),
cancel_ctx: Cell::new(None),
global_this: None,
this_jsvalue: JSValue::ZERO,
this_jsvalue: jsc::JsRef::empty(),
is_closed: Cell::new(false),
}
}
Expand All @@ -692,9 +700,11 @@ pub(crate) trait NewSourceCodegen {
fn pending_promise_set_cached(this: JSValue, global: &JSGlobalObject, value: JSValue);
fn on_drain_callback_set_cached(this: JSValue, global: &JSGlobalObject, value: JSValue);
fn on_drain_callback_get_cached(this: JSValue) -> Option<JSValue>;
fn on_close_callback_set_cached(this: JSValue, global: &JSGlobalObject, value: JSValue);
fn on_close_callback_get_cached(this: JSValue) -> Option<JSValue>;
}

/// Binds the four `SourceContext::js_*` accessors to the codegen'd
/// Binds the `SourceContext::js_*` accessors to the codegen'd
/// `crate::generated_classes::js_${Name}InternalReadableStreamSource` module
/// (one per `.classes.ts` entry: `Blob`, `File`, `Bytes`). The extern symbols
/// are declared exactly once inside that module — no local `extern "C"` block.
Expand Down Expand Up @@ -732,6 +742,20 @@ macro_rules! source_context_codegen {
) -> Option<$crate::webcore::jsc::JSValue> {
$crate::generated_classes::$gen::on_drain_callback_get_cached(this)
}
#[inline]
fn js_on_close_callback_set_cached(
this: $crate::webcore::jsc::JSValue,
global: &$crate::webcore::jsc::JSGlobalObject,
value: $crate::webcore::jsc::JSValue,
) {
$crate::generated_classes::$gen::on_close_callback_set_cached(this, global, value)
}
#[inline]
fn js_on_close_callback_get_cached(
this: $crate::webcore::jsc::JSValue,
) -> Option<$crate::webcore::jsc::JSValue> {
$crate::generated_classes::$gen::on_close_callback_get_cached(this)
}
};
}

Expand All @@ -754,6 +778,12 @@ impl<C: SourceContext> NewSourceCodegen for NewSource<C> {
fn on_drain_callback_get_cached(this: JSValue) -> Option<JSValue> {
C::js_on_drain_callback_get_cached(this)
}
fn on_close_callback_set_cached(this: JSValue, global: &JSGlobalObject, value: JSValue) {
C::js_on_close_callback_set_cached(this, global, value)
}
fn on_close_callback_get_cached(this: JSValue) -> Option<JSValue> {
C::js_on_close_callback_get_cached(this)
}
}

// Enforce the layout invariant `from_js`/`Source` rely on.
Expand Down Expand Up @@ -870,14 +900,37 @@ impl<C: SourceContext> NewSource<C> {
fn on_js_close(ptr: Option<*mut c_void>) {
// SAFETY: ptr was set to `self as *mut NewSource<C>` in on_close()/set_on_close_from_js.
let this = unsafe { &mut *(ptr.unwrap().cast::<NewSource<C>>()) };
if let Some(cb) = this.close_jsvalue.try_swap() {
this.global_this().queue_microtask(cb, &[]);
// Reached from `FileReader::on_reader_done` off the event loop. While
// the across-read ref is held (`increment_count` upgraded to Strong),
// the wrapper is rooted and `try_get()` is `Some`. If the wrapper was
// already finalized, `try_get()` is `None` and there is no callback.
let Some(this_jsvalue) = this.this_jsvalue.try_get() else {
return;
};
let global_this = this.global_this();
if let Some(cb) = <Self as NewSourceCodegen>::on_close_callback_get_cached(this_jsvalue) {
if !cb.is_undefined() {
global_this.queue_microtask(cb, &[]);
}
}
this.close_jsvalue.deinit();
<Self as NewSourceCodegen>::on_close_callback_set_cached(
this_jsvalue,
global_this,
JSValue::UNDEFINED,
);
}

pub fn increment_count(&mut self) {
self.ref_count += 1;
// A ref beyond the JS wrapper's own is held (in practice a FileReader
// `waiting_for_on_reader_done` I/O ref). Root the wrapper so
// `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 let Some(global) = self.global_this.as_deref() {
if self.this_jsvalue.is_not_empty() {
self.this_jsvalue.upgrade(global);
}
}
}

/// Release one reference. If the count hits zero, runs context teardown and
Expand All @@ -901,10 +954,15 @@ impl<C: SourceContext> NewSource<C> {
*r -= 1;
*r
};
if remaining == 1 {
// Only the JS wrapper's own ref remains: drop the Strong root so
// the wrapper becomes collectable again.
// SAFETY: caller contract — `this` is live while remaining > 0.
unsafe { (*this).this_jsvalue.downgrade() };
}
if remaining == 0 {
// SAFETY: still live; run side-effect teardown while fields are valid.
unsafe {
(*this).close_jsvalue.deinit();
(*this).context.deinit_fn();
}
// SAFETY: `this` originated from `Box::into_raw` in `Self::new`. No
Expand All @@ -925,13 +983,15 @@ impl<C: SourceContext> NewSource<C> {
}

pub fn to_readable_stream(&mut self, global_this: &JSGlobalObject) -> JsResult<JSValue> {
let out_value = if self.this_jsvalue != JSValue::ZERO {
self.this_jsvalue
let out_value = if let Some(v) = self.this_jsvalue.try_get() {
v
} else {
<Self as NewSourceCodegen>::to_js(self, global_this)
};
out_value.ensure_still_alive();
self.this_jsvalue = out_value;
if self.this_jsvalue.is_empty() {
self.this_jsvalue = jsc::JsRef::init_weak(out_value);
}
ReadableStream::from_native(global_this, out_value)
}

Expand Down Expand Up @@ -983,7 +1043,6 @@ impl<C: SourceContext> NewSource<C> {
let arguments = call_frame.arguments_old::<2>();
let view = arguments.ptr[0];
view.ensure_still_alive();
self.this_jsvalue = this_jsvalue;
let Some(mut buffer) = view.as_array_buffer(global_this) else {
return Ok(JSValue::UNDEFINED);
};
Expand All @@ -994,10 +1053,9 @@ impl<C: SourceContext> NewSource<C> {
pub fn start_from_js(
&mut self,
global_this: &JSGlobalObject,
call_frame: &CallFrame,
_call_frame: &CallFrame,
) -> JsResult<JSValue> {
self.global_this = Some(bun_ptr::BackRef::new(global_this));
self.this_jsvalue = call_frame.this();
match self.on_start_from_js() {
streams::Start::Empty => Ok(JSValue::js_number(0.0)),
streams::Start::Ready => Ok(JSValue::js_number(16384.0)),
Expand Down Expand Up @@ -1065,9 +1123,8 @@ impl<C: SourceContext> NewSource<C> {
pub fn cancel_from_js(
&mut self,
_global_object: &JSGlobalObject,
call_frame: &CallFrame,
_call_frame: &CallFrame,
) -> JsResult<JSValue> {
self.this_jsvalue = call_frame.this();
self.cancel();
Ok(JSValue::UNDEFINED)
}
Expand All @@ -1084,7 +1141,13 @@ impl<C: SourceContext> NewSource<C> {
self.global_this = Some(bun_ptr::BackRef::new(global_object));

if value.is_undefined() {
self.close_jsvalue.deinit();
if let Some(this_jsvalue) = self.this_jsvalue.try_get() {
<Self as NewSourceCodegen>::on_close_callback_set_cached(
this_jsvalue,
global_object,
JSValue::UNDEFINED,
);
}
return Ok(());
}

Expand All @@ -1096,7 +1159,13 @@ impl<C: SourceContext> NewSource<C> {
));
}
let cb = value.with_async_context_if_needed(global_object);
self.close_jsvalue.set(global_object, cb);
if let Some(this_jsvalue) = self.this_jsvalue.try_get() {
<Self as NewSourceCodegen>::on_close_callback_set_cached(
this_jsvalue,
global_object,
cb,
);
}
Ok(())
}

Expand All @@ -1107,9 +1176,13 @@ impl<C: SourceContext> NewSource<C> {
) -> JsResult<()> {
self.global_this = Some(bun_ptr::BackRef::new(global_object));

let Some(this_jsvalue) = self.this_jsvalue.try_get() else {
return Ok(());
};

if value.is_undefined() {
<Self as NewSourceCodegen>::on_drain_callback_set_cached(
self.this_jsvalue,
this_jsvalue,
global_object,
JSValue::UNDEFINED,
);
Expand All @@ -1124,20 +1197,25 @@ impl<C: SourceContext> NewSource<C> {
));
}
let cb = value.with_async_context_if_needed(global_object);
<Self as NewSourceCodegen>::on_drain_callback_set_cached(
self.this_jsvalue,
global_object,
cb,
);
<Self as NewSourceCodegen>::on_drain_callback_set_cached(this_jsvalue, global_object, cb);
Ok(())
}

pub fn get_on_close_from_js(&mut self, _global_object: &JSGlobalObject) -> JSValue {
self.close_jsvalue.get().unwrap_or(JSValue::UNDEFINED)
if let Some(this_jsvalue) = self.this_jsvalue.try_get() {
if let Some(val) =
<Self as NewSourceCodegen>::on_close_callback_get_cached(this_jsvalue)
{
return val;
}
}
JSValue::UNDEFINED
}

pub fn get_on_drain_from_js(&mut self, _global_object: &JSGlobalObject) -> JSValue {
<Self as NewSourceCodegen>::on_drain_callback_get_cached(self.this_jsvalue)
self.this_jsvalue
.try_get()
.and_then(<Self as NewSourceCodegen>::on_drain_callback_get_cached)
.unwrap_or(JSValue::UNDEFINED)
}

Expand All @@ -1146,7 +1224,6 @@ impl<C: SourceContext> NewSource<C> {
_global_object: &JSGlobalObject,
call_frame: &CallFrame,
) -> JsResult<JSValue> {
self.this_jsvalue = call_frame.this();
let ref_or_unref = call_frame.argument(0).to_boolean();
self.set_ref(ref_or_unref);
Ok(JSValue::UNDEFINED)
Expand All @@ -1158,7 +1235,7 @@ impl<C: SourceContext> NewSource<C> {
// the raw refcount via a raw pointer (the call may free `*this`).
let this = Box::into_raw(self);
// SAFETY: `this` is live — just unwrapped from `Box`.
unsafe { (*this).this_jsvalue = JSValue::ZERO };
unsafe { (*this).this_jsvalue.finalize() };
// SAFETY: `this` is live; the JS-wrapper ref below still pins the count.
if unsafe { (*this).context.finalize_detach() } {
// SAFETY: `this` is live; the JS-wrapper +1 (released below) keeps ref_count > 0.
Expand All @@ -1171,9 +1248,8 @@ impl<C: SourceContext> NewSource<C> {
pub fn drain_from_js(
&mut self,
global_this: &JSGlobalObject,
call_frame: &CallFrame,
_call_frame: &CallFrame,
) -> JsResult<JSValue> {
self.this_jsvalue = call_frame.this();
let mut list = self.drain();
if list.len() > 0 {
// Ownership of the buffer transfers to JSC: `to_js` installs
Expand All @@ -1192,10 +1268,9 @@ impl<C: SourceContext> NewSource<C> {
fn to_buffered_value_from_js(
&mut self,
global_this: &JSGlobalObject,
call_frame: &CallFrame,
_call_frame: &CallFrame,
action: streams::BufferActionTag,
) -> JsResult<JSValue> {
self.this_jsvalue = call_frame.this();
if let Some(r) = self.context.to_buffered_value(global_this, action) {
return r;
}
Expand Down
Loading
Loading