diff --git a/src/jsc/AbortSignal.rs b/src/jsc/AbortSignal.rs index d714a0501133..78332ac48cf9 100644 --- a/src/jsc/AbortSignal.rs +++ b/src/jsc/AbortSignal.rs @@ -1,3 +1,4 @@ +use core::cell::Cell; use core::ffi::c_void; use core::ptr::NonNull; use core::sync::atomic::Ordering; @@ -47,6 +48,7 @@ unsafe extern "C" { safe fn WebCore__AbortSignal__fromJS(value0: JSValue) -> *mut AbortSignal; safe fn WebCore__AbortSignal__ref(arg0: &AbortSignal) -> *mut AbortSignal; safe fn WebCore__AbortSignal__toJS(arg0: &AbortSignal, arg1: &JSGlobalObject) -> JSValue; + // Only called from `ext_deref`; see `ref_`. safe fn WebCore__AbortSignal__unref(arg0: &AbortSignal); // `*mut Timeout` is round-tripped opaquely through C++ (stored from // `AbortSignal__Timeout__create`, never dereferenced on the C++ side), so @@ -159,19 +161,14 @@ impl AbortSignal { )) } + /// Takes a ref; release it by adopting it into an [`AbortSignalRef`]. There + /// is no `unref(&self)` on purpose: `&AbortSignal` does not prove ownership + /// of a ref (every `AbortSignalRef` derefs to one), so a safe release here + /// would let safe code double-release. pub fn ref_(&self) -> *mut AbortSignal { WebCore__AbortSignal__ref(self) } - pub fn unref(&self) { - WebCore__AbortSignal__unref(self) - } - - pub fn detach(&self, ctx: *mut c_void) { - self.clean_native_bindings(ctx); - self.unref(); - } - /// Lifetime: the returned pointer is borrowed from the JS wrapper and is /// valid only while `value` remains reachable. Use [`AbortSignal::ref_from_js`] /// to take refcounted ownership instead. @@ -188,9 +185,11 @@ impl AbortSignal { WebCore__AbortSignal__create(global) } - pub fn new(global: &JSGlobalObject) -> *mut AbortSignal { + pub fn new(global: &JSGlobalObject) -> AbortSignalRef { crate::mark_binding!(); - WebCore__AbortSignal__new(global) + // SAFETY: C++ returns `leakRef()` of a fresh signal, i.e. the `+1` + // adopted here. + unsafe { AbortSignalRef::adopt(WebCore__AbortSignal__new(global)) } } /// Returns a borrowed handle to the internal Timeout, or null. @@ -199,9 +198,9 @@ impl AbortSignal { /// /// Thread-safety: not thread-safe; call only on the owning thread/loop. /// - /// Usage: if you need to operate on the Timeout (run/cancel/deinit), hold a ref - /// to `this` for the duration (e.g., `this.ref_(); defer this.unref();`) and avoid - /// caching the pointer across turns. + /// Usage: if you need to operate on the Timeout (run/cancel/deinit), hold an + /// [`AbortSignalRef`] to `self` for the duration and avoid caching the + /// pointer across turns. pub fn get_timeout(&self) -> Option<&Timeout> { let ptr = WebCore__AbortSignal__getTimeout(self); // SAFETY: returned Timeout is owned by `self` and valid while `self` is held @@ -256,6 +255,66 @@ impl AbortSignal { } } +/// What a native operation holds on a signal while it is in flight: a ref, a +/// pending-activity count (keeps the JS wrapper and its `abort` listeners +/// alive), and at most one native listener. `Drop` gives all three back, so +/// holders keep an `Option` and release by taking it out. +pub struct PendingActivityRef { + signal: AbortSignalRef, + /// Null until [`Self::add_listener`]. + listener_ctx: Cell<*mut c_void>, +} + +impl PendingActivityRef { + pub fn new(signal: AbortSignalRef) -> Self { + signal.pending_activity_ref(); + Self { + signal, + listener_ctx: Cell::new(core::ptr::null_mut()), + } + } + + /// [`AbortSignal::add_listener`], but removed again when `self` drops. + /// Takes `&self` so it shadows the `Deref`'d original for every receiver. + pub fn add_listener( + &self, + ctx: *mut c_void, + callback: unsafe extern "C" fn(*mut c_void, JSValue), + ) { + debug_assert!( + self.listener_ctx.get().is_null(), + "PendingActivityRef already has a listener" + ); + self.listener_ctx.set(ctx); + self.signal.add_listener(ctx, callback); + } + + pub fn signal_ref(&self) -> AbortSignalRef { + self.signal.clone() + } +} + +impl core::ops::Deref for PendingActivityRef { + type Target = AbortSignal; + #[inline] + fn deref(&self) -> &AbortSignal { + &self.signal + } +} + +impl Drop for PendingActivityRef { + fn drop(&mut self) { + let ctx = self.listener_ctx.get(); + if !ctx.is_null() { + self.signal.clean_native_bindings(ctx); + } + // Second on purpose: `cleanNativeBindings` cancels an `AbortSignal.timeout` + // timer once nothing observes the signal, and pending activity counts + // as observing it. The ref itself goes when `self.signal` drops. + self.signal.pending_activity_unref(); + } +} + pub enum AbortReason { Common(CommonAbortReason), Js(JSValue), diff --git a/src/runtime/api/bun/h2_frame_parser.rs b/src/runtime/api/bun/h2_frame_parser.rs index 1be650bff0d6..6f3d5de20a13 100644 --- a/src/runtime/api/bun/h2_frame_parser.rs +++ b/src/runtime/api/bun/h2_frame_parser.rs @@ -32,7 +32,8 @@ use bun_jsc::abort_signal::AbortListener; use bun_jsc::array_buffer::BinaryType; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{ - CallFrame, GlobalRef, JSGlobalObject, JSValue, JsCell, JsClass, JsRef, JsResult, StrongOptional, + AbortSignalRef, CallFrame, GlobalRef, JSGlobalObject, JSValue, JsCell, JsClass, JsRef, + JsResult, StrongOptional, }; use bun_ptr::IntrusiveRc; @@ -1549,14 +1550,7 @@ pub struct Stream { } pub(crate) struct SignalRef { - // LIFETIMES.tsv: SHARED — AbortSignal is intrusively refcounted across FFI/codegen. - // `AbortSignal` is an opaque C++ type whose ref/unref go through - // `WebCore__AbortSignal__ref/unref`; it does not (and cannot) implement - // `bun_ptr::RefCounted`, so balance refs by hand in `attach_signal` / - // `Drop`. `BackRef` captures the backref invariant - // (signal is `ref_()`'d in `attach_signal` and outlives this struct until - // `Drop` calls `detach()`/`unref()`), so reads go through safe `Deref`. - signal: bun_ptr::BackRef, + signal: AbortSignalRef, // LIFETIMES.tsv: SHARED — H2FrameParser carries an intrusive RefCount and is // recovered via `from_field_ptr!` from the auto-flusher. It uses a hand-rolled // `Cell` ref count (not `bun_ptr::RefCount`), so `IntrusiveRc`'s @@ -1570,7 +1564,6 @@ pub(crate) struct SignalRef { impl SignalRef { pub(crate) fn is_aborted(&self) -> bool { - // BackRef invariant: signal kept alive via .ref_() in attach_signal. self.signal.aborted() } @@ -1594,12 +1587,9 @@ impl SignalRef { impl Drop for SignalRef { fn drop(&mut self) { - // BackRef invariant: `signal` is the C++-refcounted AbortSignal we - // ref_()'d in `attach_signal`; valid until this `detach` releases our - // listener and unrefs. Copy the `BackRef` out first so the `&mut self` - // taken by `from_mut` doesn't overlap the receiver borrow. - let signal = self.signal; - signal.detach(std::ptr::from_mut(self).cast::()); + // `attach_signal` registered the listener with `self`'s address as ctx. + let ctx = std::ptr::from_mut(self).cast::(); + self.signal.clean_native_bindings(ctx); // ParentRef backref — parser outlives every SignalRef (ref()'d in // `attach_signal`); release that ref now via the inherent `deref()`. H2FrameParser::deref(self.parser.get()); @@ -2120,20 +2110,19 @@ impl Stream { .unwrap_or_else(|| JSValue::js_number(self.id as f64)) } - pub fn attach_signal(&mut self, parser: &H2FrameParser, signal: &mut AbortSignal) { - // `ref_()` bumps the C++ intrusive refcount and returns the same live - // `self` pointer with FFI (wildcard) provenance — store *that* in the - // `BackRef` so its validity is tied to the refcount, not to the - // borrowed `&mut AbortSignal` parameter's lifetime. - let refed = core::ptr::NonNull::new(signal.ref_()).expect("AbortSignal::ref_"); + pub fn attach_signal(&mut self, parser: &H2FrameParser, signal: &AbortSignal) { + // SAFETY: `signal` is live (borrowed from the JS wrapper the caller is + // holding); `ref_()` returns the same pointer carrying a `+1`, which + // the `AbortSignalRef` now owns. + let owned = unsafe { AbortSignalRef::adopt(signal.ref_()) }; // we need a stable pointer to know what signal points to what stream_id + parser let mut signal_ref = Box::new(SignalRef { - signal: bun_ptr::BackRef::from(refed), + signal: owned, parser: bun_ptr::ParentRef::new(parser), stream_id: self.id, }); // `signal_ref` is heap-allocated and outlives the listener registration - // (cleared via `detach` in `Drop for SignalRef`). + // (removed in `Drop for SignalRef`). signal.listen(&raw mut *signal_ref); // TODO: We should not need this ref counting here, since Parser owns Stream parser.ref_(); @@ -9360,8 +9349,7 @@ impl H2FrameParser { if let Some(signal_arg) = options.get(global_object, "signal")? { if let Some(signal_ptr) = AbortSignal::from_js(signal_arg) { - // SAFETY: `from_js` returns a live *mut AbortSignal owned by JSC; rooted via `signal_arg` on the stack. - let signal_ = unsafe { &mut *signal_ptr }; + let signal_ = AbortSignal::opaque_ref(signal_ptr); if signal_.aborted() { stream.state = StreamState::IDLE; let wrapped = diff --git a/src/runtime/api/bun/js_bun_spawn_bindings.rs b/src/runtime/api/bun/js_bun_spawn_bindings.rs index 5397137e80f4..9f245bd4079c 100644 --- a/src/runtime/api/bun/js_bun_spawn_bindings.rs +++ b/src/runtime/api/bun/js_bun_spawn_bindings.rs @@ -367,22 +367,13 @@ fn spawn_maybe_sync( let mut windows_hide: bool = false; #[cfg(windows)] let mut windows_verbatim_arguments: bool = false; - let mut abort_signal: Option<*mut WebCore::AbortSignal> = None; + let mut abort_signal: Option = None; let mut terminal_info: Option = None; let mut existing_terminal: Option> = None; // Existing terminal passed by user let mut terminal_js_value: JSValue = JSValue::ZERO; let mut defer_guard = scopeguard::guard( - (&mut abort_signal, &mut terminal_info), - |(abort_signal, terminal_info): ( - &mut Option<*mut WebCore::AbortSignal>, - &mut Option, - )| { - if let Some(signal) = abort_signal.take() { - // signal was ref()'d when stored; unref releases that ref. - // `AbortSignal` is an `opaque_ffi!` ZST handle; `opaque_ref` is - // the centralised non-null deref proof. - WebCore::AbortSignal::opaque_ref(signal).unref(); - } + &mut terminal_info, + |terminal_info: &mut Option| { // If we created a new terminal but spawn failed, close it. The // writer/reader/finalize deref paths release the remaining refs. // Downgrade the JSRef so the wrapper is GC-eligible, and mark @@ -395,8 +386,8 @@ fn spawn_maybe_sync( } }, ); - // Note: reshaped for borrowck — re-borrow through the guard tuple. - let (abort_signal, terminal_info) = &mut *defer_guard; + // Note: reshaped for borrowck — re-borrow through the guard. + let terminal_info = &mut *defer_guard; // Owned ZBox for `cwd` held here so the `&[u8]` borrow stays valid until // `spawn_process` returns. @@ -511,15 +502,11 @@ fn spawn_maybe_sync( } if let Some(signal_val) = args.get_truthy(global_this, "signal")? { - if let Some(signal) = WebCore::AbortSignal::from_js(signal_val) { - // `from_js` returns a live FFI handle owned by JS. - // `AbortSignal` is an `opaque_ffi!` ZST handle; `opaque_ref` - // is the centralised non-null deref proof. - let sig = WebCore::AbortSignal::opaque_ref(signal); - if let Some(abort_error) = sig.node_abort_error_if_aborted(global_this) { + if let Some(signal) = WebCore::AbortSignal::ref_from_js(signal_val) { + if let Some(abort_error) = signal.node_abort_error_if_aborted(global_this) { return Err(global_this.throw_value(abort_error)); } - **abort_signal = Some(sig.ref_()); + abort_signal = Some(signal); } else { return Err(global_this.throw_invalid_argument_type_value( b"signal", @@ -1336,7 +1323,7 @@ fn spawn_maybe_sync( closed: Default::default(), this_value: Default::default(), weak_file_sink_stdin_ptr: Cell::new(None), - abort_signal: Cell::new(None), + abort_signal: JsCell::new(None), event_loop_timer_refd: Cell::new(false), event_loop_timer: JsCell::new(crate::timer::EventLoopTimer::init_paused( crate::timer::EventLoopTimerTag::SubprocessTimeout, @@ -1801,16 +1788,14 @@ fn spawn_maybe_sync( // Adding the abort listener may call the onAbortSignal callback immediately if it was already aborted // Therefore, we must do this at the very end. if let Some(signal) = abort_signal.take() { - // SAFETY: `signal` is a live *mut AbortSignal carrying the +1 ref taken - // above; ownership of that ref transfers to `subprocess.abort_signal`. + let signal = jsc::abort_signal::PendingActivityRef::new(signal); // `add_listener` may synchronously fire `on_abort_signal` (already - // aborted), which re-enters via `subprocess_ptr` — write through the - // raw pointer so no `&mut Subprocess` is held across the call. - unsafe { - (*signal).pending_activity_ref(); - let _ = (*signal).add_listener(subprocess_ptr.cast(), Subprocess::on_abort_signal); - (*subprocess_ptr).abort_signal.set(NonNull::new(signal)); - } + // aborted), which re-enters via `subprocess_ptr`, so the store below + // goes through the raw pointer rather than a `&mut Subprocess`. + signal.add_listener(subprocess_ptr.cast(), Subprocess::on_abort_signal); + // SAFETY: `subprocess_ptr` is the live Subprocess allocated above; + // `clear_abort_signal` drops the hold. + unsafe { (*subprocess_ptr).abort_signal.set(Some(signal)) }; } if !IS_SYNC { @@ -1841,24 +1826,8 @@ fn spawn_maybe_sync( // watchOrReap will handle the already exited case for us. } - match subprocess.process_mut().watch_or_reap() { - sys::Result::Ok(_) => { - // Once everything is set up, we can add the abort listener - // Adding the abort listener may call the onAbortSignal callback immediately if it was already aborted - // Therefore, we must do this at the very end. - if let Some(signal) = abort_signal.take() { - // SAFETY: see the matching block above. - unsafe { - (*signal).pending_activity_ref(); - let _ = - (*signal).add_listener(subprocess_ptr.cast(), Subprocess::on_abort_signal); - (*subprocess_ptr).abort_signal.set(NonNull::new(signal)); - } - } - } - sys::Result::Err(_) => { - subprocess.process_mut().wait(true); - } + if subprocess.process_mut().watch_or_reap().is_err() { + subprocess.process_mut().wait(true); } if !subprocess.has_exited() { diff --git a/src/runtime/api/bun/subprocess.rs b/src/runtime/api/bun/subprocess.rs index 85098e028179..a1def6cb6f63 100644 --- a/src/runtime/api/bun/subprocess.rs +++ b/src/runtime/api/bun/subprocess.rs @@ -29,7 +29,7 @@ use crate::api::bun_process::{Process, Rusage, Status}; use crate::ipc as IPC; use crate::node::node_cluster_binding; use crate::timer::{EventLoopTimer, EventLoopTimerState}; -use crate::webcore::{self, AbortSignal, FileSink}; +use crate::webcore::{self, FileSink}; #[cfg(windows)] use bun_libuv_sys::UvHandle as _; @@ -147,10 +147,8 @@ pub struct Subprocess<'a> { /// Weak observer of the stdin `FileSink` — holds no ownership/ref. `onStdinDestroyed` /// nulls this before the sink is freed, so it is never dereferenced after the sink dies. pub(crate) weak_file_sink_stdin_ptr: Cell>>, - /// +1 C++-intrusive ref held; released in `clear_abort_signal` via - /// `AbortSignal::unref()`. Not `Arc` — `AbortSignal` is an opaque FFI - /// handle whose refcount lives on the C++ side. - pub(crate) abort_signal: Cell>>, + /// The `signal` spawn option, with our abort listener registered on it. + pub(crate) abort_signal: JsCell>, pub(crate) event_loop_timer_refd: Cell, /// Intrusive timer node. `JsCell` so `&self` can hand `*mut EventLoopTimer` @@ -363,16 +361,10 @@ bun_spawn::link_impl_ProcessExit! { } impl Subprocess<'_> { - /// Shared borrow of the attached `AbortSignal`, if any. - /// - /// `abort_signal` holds a +1 C++-intrusive ref taken in - /// `spawn_maybe_sync`; the pointee is therefore live for as long as the - /// cell is `Some` (it is `take`n *before* `unref()` in - /// [`clear_abort_signal`](Self::clear_abort_signal)) — i.e. the - /// owner-outlives-holder `BackRef` invariant holds. + /// Owned, so it stays valid if `clear_abort_signal` runs meanwhile. #[inline] - pub(crate) fn abort_signal_ref(&self) -> Option> { - self.abort_signal.get().map(bun_ptr::BackRef::from) + pub(crate) fn abort_signal_ref(&self) -> Option { + self.abort_signal.get().as_ref().map(|s| s.signal_ref()) } #[bun_jsc::host_fn(method)] @@ -1279,14 +1271,7 @@ impl Subprocess<'_> { } fn clear_abort_signal(&self) { - if let Some(signal) = self.abort_signal.replace(None).map(bun_ptr::BackRef::from) { - // `signal` was stored with a +1 C++ intrusive ref (taken in - // `spawn_maybe_sync`); it stays live until `unref()` below, so the - // `BackRef` invariant (pointee outlives holder) holds for this scope. - signal.pending_activity_unref(); - signal.clean_native_bindings(self.as_ctx_ptr().cast::()); - signal.unref(); - } + drop(self.abort_signal.replace(None)); } pub fn finalize(self: Box) { diff --git a/src/runtime/server/RequestContext.rs b/src/runtime/server/RequestContext.rs index aceca771e0ce..cf2fda04815d 100644 --- a/src/runtime/server/RequestContext.rs +++ b/src/runtime/server/RequestContext.rs @@ -12,8 +12,8 @@ use bun_uws::{self as uws, WebSocketUpgradeContext}; use crate::server::jsc::{self, JSGlobalObject, JSValue, JsResult, VirtualMachine}; use crate::server::{RangeRequest, ServerLike}; use crate::webcore::{ - self as WebCore, AbortSignal, AnyBlob, ByteStream, CookieMap, CookieMapRef, FetchHeaders, - Request, Response, blob::SizeType as BlobSizeType, body, readable_stream, request, response, + self as WebCore, AnyBlob, ByteStream, CookieMap, CookieMapRef, FetchHeaders, Request, Response, + blob::SizeType as BlobSizeType, body, readable_stream, request, response, }; /// Q: Why is this needed? @@ -117,13 +117,10 @@ pub struct RequestContext< pub(crate) resp: Cell>, pub(crate) req: Cell>>, pub(crate) request_weakref: JsCell, - // NOTE: `Arc` was wrong — - // `AbortSignal` is an opaque ZST FFI handle; an `Arc` of a ZST never owns - // the C++ allocation. Store the raw pointer. The request holds TWO counts: - // the intrusive C++ `RefPtr` (+1 from `AbortSignal::new()`/`ref_()`) and a - // pending-activity count for GC visibility. Both are released together via - // `shim::signal_release` in `on_abort`/`finalize_without_deinit`. - pub(crate) signal: Cell>>, + /// `request.signal` (the `Request` holds its own ref to it). Taken out in + /// `on_abort` / `finalize_without_deinit`, or moved to the `ServerWebSocket` + /// on upgrade. + pub(crate) signal: JsCell>, pub method: Method, /// Owned `+1` ref on a C++ `CookieMap` (taken in `set_cookies`, released /// when the field is dropped/replaced — `CookieMapRef` handles the unref). @@ -369,36 +366,6 @@ mod shim { r.detach_readable_stream(g) } #[inline] - pub(super) fn signal_aborted(s: NonNull) -> bool { - // `signal` is kept alive by the intrusive C++ refcount (+1 from - // `AbortSignal::new()` / `ref_()`) plus `pending_activity_ref()` until - // `signal_release` drops both — satisfies the `BackRef` outlives-holder - // invariant for the duration of this call. - bun_ptr::BackRef::from(s).aborted() - } - #[inline] - pub(super) fn signal_fire( - s: NonNull, - g: &JSGlobalObject, - r: jsc::CommonAbortReason, - ) { - // See `signal_aborted` — counted ref keeps pointee live. - bun_ptr::BackRef::from(s).signal(g, r) - } - /// Release BOTH refcounts the request holds on its AbortSignal. - /// `pending_activity_unref()` drops the GC-visibility count and `unref()` - /// drops the intrusive C++ `RefPtr` count taken at creation. `s` must not - /// be dereferenced after this call. - #[inline] - pub(super) fn signal_release(s: NonNull) { - // See `signal_aborted`. Order: pending-activity first, - // then the owning intrusive ref (which may free). `BackRef` is dropped - // before `unref()` returns, so no dangling deref. - let signal = bun_ptr::BackRef::from(s); - signal.pending_activity_unref(); - signal.unref(); - } - #[inline] pub(super) fn iec_trigger( cb: &bun_jsc::JsCell, ev: request::EventType, @@ -653,12 +620,15 @@ where } pub(crate) fn set_signal_aborted(&self, reason: jsc::CommonAbortReason) { - if let Some(signal) = self.signal.get() { - if let Some(server) = self.server.get() { - // server is a BACKREF — valid while this RequestContext is alive - let global = server.global_this(); - shim::signal_fire(signal, global, reason); - } + // Own ref, not a borrow of the cell: abort listeners may end the request + // and empty it. + let Some(signal) = self.signal.get().as_ref().map(|s| s.signal_ref()) else { + return; + }; + if let Some(server) = self.server.get() { + // server is a BACKREF — valid while this RequestContext is alive + let global = server.global_this(); + signal.signal(global, reason); } } @@ -1331,7 +1301,7 @@ where defer_deinit_until_callback_completes: Cell::new(should_deinit_context), range: RangeRequest::raw_from_request(&Self::any_request(req)), request_weakref: JsCell::new(request::WeakRef::EMPTY), - signal: Cell::new(None), + signal: JsCell::new(None), cookies: JsCell::new(None), flags: Flags::::default(), upgrade_context: Cell::new(UpgradeState::None), @@ -1442,16 +1412,11 @@ where this.request_weakref.with_mut(|w| w.deref()); } // if signal is not aborted, abort the signal - if let Some(signal) = this.signal.take() { - if !shim::signal_aborted(signal) { - shim::signal_fire( - signal, - global_this, - jsc::CommonAbortReason::ConnectionClosed, - ); + if let Some(signal) = this.signal.replace(None) { + if !signal.aborted() { + signal.signal(global_this, jsc::CommonAbortReason::ConnectionClosed); any_js_calls.set(true); } - shim::signal_release(signal); } // if have sink, call onAborted on sink @@ -1534,15 +1499,10 @@ where } // if signal is not aborted, abort the signal - if let Some(signal) = self.signal.take() { - if self.flags.aborted() && !shim::signal_aborted(signal) { - shim::signal_fire( - signal, - global_this, - jsc::CommonAbortReason::ConnectionClosed, - ); + if let Some(signal) = self.signal.replace(None) { + if self.flags.aborted() && !signal.aborted() { + signal.signal(global_this, jsc::CommonAbortReason::ConnectionClosed); } - shim::signal_release(signal); } // Case 1: diff --git a/src/runtime/server/ServerWebSocket.rs b/src/runtime/server/ServerWebSocket.rs index 8c01305a3f4c..610fa9777f79 100644 --- a/src/runtime/server/ServerWebSocket.rs +++ b/src/runtime/server/ServerWebSocket.rs @@ -1,7 +1,6 @@ use core::cell::Cell; use core::ffi::c_void; use core::mem; -use core::ptr::NonNull; use bun_jsc::JsCell; use bun_uws::{self as uws, AnyWebSocket, WebSocketBehavior}; @@ -10,8 +9,8 @@ use bun_uws_sys::{Opcode, SendStatus}; use crate::server::WebSocketServerHandler; use crate::server::jsc::{ - self, AbortSignal, ArrayBuffer, BinaryType, CallFrame, CommonAbortReason, JSGlobalObject, - JSType, JSValue, JsError, JsRef, JsResult, ZigStringSlice, + self, ArrayBuffer, BinaryType, CallFrame, CommonAbortReason, JSGlobalObject, JSType, JSValue, + JsError, JsRef, JsResult, ZigStringSlice, abort_signal, }; use crate::server::web_socket_server_context::HandlerFlags; use crate::webcore::Blob; @@ -26,17 +25,15 @@ bun_output::declare_scope!(WebSocketServer, visible); // — `on_open` → `ws.cork(JS)` → `ws.close()` → `on_close` mutates `flags` / // `this_value` on the SAME `m_ctx`. A `&mut Self` receiver would alias under // Stacked Borrows. Receivers therefore take `&self`; per-field interior -// mutability (`Cell` for `Copy` flags/signal, `JsCell` for the non-`Copy` -// `JsRef`) carries the writes. +// mutability (`Cell`, `JsCell` for the non-`Copy` `JsRef`) carries the writes. #[bun_jsc::JsClass] pub struct ServerWebSocket { handler: bun_ptr::BackRef, this_value: JsCell, flags: Cell, - // `AbortSignal` is an opaque C++ type - // with intrusive WebCore ref-counting (ref/unref) — never `Arc`. The init - // caller transfers a +1 ref; `finalize`/`on_close` unref it. - signal: Cell>>, + /// The upgraded request's signal, moved here from its `RequestContext`; + /// `on_close` takes it out and fires it. + signal: Cell>, } // We pack the per-socket data into this struct below: @@ -350,12 +347,11 @@ impl ServerWebSocket { // pub const js = jsc.Codegen.JSServerWebSocket; — provided by #[bun_jsc::JsClass] // toJS / fromJS / fromJSDirect — provided by codegen (see `to_js_ptr` / `JsClass` impl) - /// Initialize a ServerWebSocket with the given handler, data value, and signal. - /// The signal will not be ref'd inside the ServerWebSocket init function, but will unref itself when the ServerWebSocket is destroyed. + /// `signal` is the upgraded request's; it is fired when the socket closes. pub(crate) fn init( handler: &WebSocketServerHandler, data_value: JSValue, - signal: Option>, + signal: Option, ) -> *mut ServerWebSocket { let global_object = handler.global_object(); let this = bun_core::heap::into_raw(Box::new(ServerWebSocket { @@ -699,15 +695,8 @@ impl ServerWebSocket { .try_get() .unwrap_or(JSValue::UNDEFINED); let this_value_cell: &JsCell = &self.this_value; - let _cleanup = scopeguard::guard(signal, move |sig| { - if let Some(sig) = sig { - // `sig` was stored with a +1 ref by the upgrade caller; it - // stays live until this paired `unref()`, so the transient - // `BackRef` (pointee-outlives-holder) is sound for both calls. - let sig = bun_ptr::BackRef::from(sig); - sig.pending_activity_unref(); - sig.unref(); - } + let cleanup = scopeguard::guard(signal, move |signal| { + drop(signal); if was_not_empty { // The `server` traced edge set in `init` must outlive this // wrapper: `self.handler` points into the server's allocation @@ -736,10 +725,7 @@ impl ServerWebSocket { let _loop_guard = vm.enter_event_loop_scope(); - if let Some(sig) = signal { - // `sig` is held alive by the +1 ref released in `_cleanup`; - // BackRef invariant (pointee outlives the temporary) holds. - let sig = bun_ptr::BackRef::from(sig); + if let Some(sig) = cleanup.as_deref() { if !sig.aborted() { sig.signal(handler.global_object(), CommonAbortReason::ConnectionClosed); } @@ -766,12 +752,9 @@ impl ServerWebSocket { handler.run_error_callback(on_error, global_object, err); return; } - } else if let Some(sig) = signal { + } else if let Some(sig) = cleanup.as_deref() { let _loop_guard = vm.enter_event_loop_scope(); - // `sig` is held alive by the +1 ref released in `_cleanup`; - // BackRef invariant (pointee outlives the temporary) holds. - let sig = bun_ptr::BackRef::from(sig); if !sig.aborted() { sig.signal(handler.global_object(), CommonAbortReason::ConnectionClosed); } @@ -804,15 +787,7 @@ impl ServerWebSocket { pub fn finalize(self: Box) { bun_output::scoped_log!(WebSocketServer, "finalize"); self.this_value.with_mut(|v| v.finalize()); - if let Some(signal) = self.signal.take() { - // `signal` was stored with a +1 ref by the upgrade caller; it - // stays live until this paired `unref()`, so the transient - // `BackRef` (pointee-outlives-holder) is sound for both calls — - // same pattern as `on_close()`'s `_cleanup` guard. - let sig = bun_ptr::BackRef::from(signal); - sig.pending_activity_unref(); - sig.unref(); - } + // Dropping the box releases `signal` if `on_close` never ran. } #[bun_jsc::host_fn(method)] diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index 564d7dc4c855..bc6b25e6f067 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -806,14 +806,9 @@ impl NewServer { ctx_ref.request_body.set(Some(body_hive.clone())); let global = server.global_this(); - let signal = jsc::AbortSignal::new(global); - // S008: `AbortSignal` is an `opaque_ffi!` ZST — safe deref. - ctx_ref.signal.set(core::ptr::NonNull::new(signal)); - bun_opaque::opaque_deref_mut(signal).pending_activity_ref(); - - // SAFETY: `signal.ref_()` bumps the intrusive count and returns +1. - let signal_ref = - unsafe { jsc::AbortSignalRef::adopt(bun_opaque::opaque_deref_mut(signal).ref_()) }; + let signal = jsc::abort_signal::PendingActivityRef::new(jsc::AbortSignal::new(global)); + let signal_ref = signal.signal_ref(); + ctx_ref.signal.set(Some(signal)); // ownership: `Request::new` is `bun.TrivialNew` — the heap // allocation is handed to the JS GC via `to_js`/`to_js_for_bake` (C++ // wrapper finalizer frees it), or, for `CreateJsRequest::No`, retained diff --git a/src/runtime/server/server_body.rs b/src/runtime/server/server_body.rs index eebed3c859dc..a091beebc618 100644 --- a/src/runtime/server/server_body.rs +++ b/src/runtime/server/server_body.rs @@ -120,7 +120,7 @@ trait RequestCtxOps: RequestCtx { reason = "the body slot is a separate pooled allocation, not a field of *self (R-2)" )] fn request_body_mut(&self) -> Option<&mut BodyValue>; - fn set_signal(&self, sig: *mut AbortSignal); + fn set_signal(&self, sig: jsc::abort_signal::PendingActivityRef); fn set_request_weakref(&self, req: *mut Request); fn clear_req(&self); fn set_is_web_browser_navigation(&self, v: bool); @@ -209,12 +209,8 @@ where .map(|h| unsafe { &mut (*h.as_ptr()).value }) } #[inline] - fn set_signal(&self, sig: *mut AbortSignal) { - // `AbortSignal::new` returns a raw +1 ref to a C++-refcounted opaque; - // `RequestContext.signal` stores it as `Option>` - // and pairs the unref in RequestContext cleanup (`shim::signal_release`, - // which drops both the pending-activity count and the intrusive ref). - self.signal.set(NonNull::new(sig)); + fn set_signal(&self, sig: jsc::abort_signal::PendingActivityRef) { + self.signal.set(Some(sig)); } #[inline] fn set_request_weakref(&self, req: *mut Request) { @@ -2183,7 +2179,7 @@ where // --- After this point, do not throw an exception // See https://github.com/oven-sh/bun/issues/1339 upgrader.upgrade_context.set(UpgradeState::Upgraded); - let signal = upgrader.signal.take(); + let signal = upgrader.signal.replace(None); upgrader.resp.set(None); // Snapshot lazy url/headers before detaching (mirrors to_async_without_abort_handler). @@ -3229,15 +3225,9 @@ where // same slot. Paired drop in `RequestContext::deinit` / `Request::finalize`. ctx.set_request_body(Some(body_hive.clone())); - let signal = AbortSignal::new(&server.global()); + let signal = jsc::abort_signal::PendingActivityRef::new(AbortSignal::new(&server.global())); + let signal_for_req = signal.signal_ref(); ctx.set_signal(signal); - // S008: `AbortSignal` is an `opaque_ffi!` ZST — safe deref. - bun_opaque::opaque_deref_mut(signal).pending_activity_ref(); - - // Bump once for the Request's owned - // copy and adopt into RAII so it pairs with `Request::Drop`'s unref. - // SAFETY: `signal` is live; `ref_()` returns the same non-null ptr +1. - let signal_for_req = unsafe { jsc::AbortSignalRef::adopt((*signal).ref_()) }; let request_object_box = Request::new(Request::init( ctx.ctx_method(), AnyRequestContext::init(std::ptr::from_ref::(ctx)), @@ -3497,17 +3487,11 @@ where // same slot. Paired drop in `RequestContext::deinit` / `Request::finalize`. ctx.request_body.set(Some(body_hive.clone())); - let signal = AbortSignal::new(&this.global()); - // The - // RequestContext owns one ref so aborts during the WS-upgrade fallback - // fetch path propagate. - ctx.signal.set(NonNull::new(signal)); - // S008: `AbortSignal` is an `opaque_ffi!` ZST — safe deref. - bun_opaque::opaque_deref_mut(signal).pending_activity_ref(); - // Bump once for the Request's copy and - // adopt into RAII so it pairs with `Request::Drop`'s unref. - // SAFETY: `signal` is live; `ref_()` returns the same non-null ptr +1. - let signal_for_req = unsafe { jsc::AbortSignalRef::adopt((*signal).ref_()) }; + // The RequestContext holds the signal too, so aborts during the + // WS-upgrade fallback fetch path propagate. + let signal = jsc::abort_signal::PendingActivityRef::new(AbortSignal::new(&this.global())); + let signal_for_req = signal.signal_ref(); + ctx.signal.set(Some(signal)); let request_object_box = Request::new(Request::init( ctx.method, AnyRequestContext::init(std::ptr::from_ref(ctx)), diff --git a/src/runtime/webcore/Response.rs b/src/runtime/webcore/Response.rs index c80be4ab57f6..1b23610eb775 100644 --- a/src/runtime/webcore/Response.rs +++ b/src/runtime/webcore/Response.rs @@ -115,7 +115,8 @@ impl Drop for HeadersRef { /// Errors the owning fetch `Response`'s body on abort (Fetch spec "abort a fetch" step 4). pub(crate) struct BodyAbortListener { - signal: AbortSignalRef, + /// `on_abort` is registered on it with this box's address as ctx. + signal: bun_jsc::abort_signal::PendingActivityRef, /// `Response` owns `Box`, so a ref-counted pointer here would cycle. response: bun_ptr::ParentRef, global: GlobalRef, @@ -155,14 +156,6 @@ impl BodyAbortListener { } } -impl Drop for BodyAbortListener { - fn drop(&mut self) { - let ctx = core::ptr::from_mut(self).cast::(); - self.signal.clean_native_bindings(ctx); - self.signal.pending_activity_unref(); - } -} - // `jsc.Codegen.JSResponse` — generated by `.classes.ts`. The Rust bindings // live in `bun_jsc::generated::JSResponse` (emitted by `js_class_module!`): // `from_js` / `from_js_direct` / `to_js` / `get_constructor` / @@ -519,19 +512,19 @@ impl Response { global: &JSGlobalObject, signal: &AbortSignal, ) { - // SAFETY: `signal` is live; `ref_()` bumps the intrusive refcount. + // SAFETY: `signal` is live; `ref_()` returns it carrying the `+1` + // adopted here. let signal_ref = unsafe { AbortSignalRef::adopt(signal.ref_()) }; - signal.pending_activity_ref(); let mut listener = Box::new(BodyAbortListener { - signal: signal_ref, + signal: bun_jsc::abort_signal::PendingActivityRef::new(signal_ref), // SAFETY: caller contract; `this` is live and owns the box. response: unsafe { bun_ptr::ParentRef::from_raw_mut(this) }, global: GlobalRef::new(global), }); - signal.add_listener( - core::ptr::from_mut(&mut *listener).cast::(), - BodyAbortListener::on_abort, - ); + let ctx = core::ptr::from_mut(&mut *listener).cast::(); + listener + .signal + .add_listener(ctx, BodyAbortListener::on_abort); // SAFETY: caller contract; `this` is live. unsafe { (*this).abort_listener.set(Some(listener)) }; } diff --git a/src/runtime/webcore/fetch.rs b/src/runtime/webcore/fetch.rs index f9d20ddba40c..c1818b29c074 100644 --- a/src/runtime/webcore/fetch.rs +++ b/src/runtime/webcore/fetch.rs @@ -55,7 +55,7 @@ use bun_core::{String as BunString, Tag as BunStringTag, ZigStringSlice}; use bun_http::{self as http, FetchRedirect, Headers, HeadersExt as _, MimeType}; use bun_http_jsc::method_jsc; use bun_http_types::Method::Method; -use bun_jsc::{HTTPHeaderName, StringJsc as _, SysErrorJsc as _}; +use bun_jsc::{AbortSignalRef, HTTPHeaderName, StringJsc as _, SysErrorJsc as _}; use bun_paths::{self, PathBuffer}; use bun_sys::FdExt as _; // `FromJsEnum for FetchRedirect` lives in bun_http_jsc; importing the impl crate @@ -129,10 +129,10 @@ impl SignalRef { impl Drop for SignalRef { fn drop(&mut self) { if let Some(sig) = self.0.take() { - // `sig` was obtained from `AbortSignal::ref_()` which bumped the - // C++ intrusive refcount; the pointee outlives this `BackRef` - // until `unref()` releases that +1. - bun_ptr::BackRef::from(sig).unref(); + // SAFETY: `sig` is the `+1` taken by `ref_()` in `extract_signal` + // and was not handed to `FetchOptions` (`take()` would have cleared + // it), so adopting it here releases that ref exactly once. + drop(unsafe { AbortSignalRef::adopt(sig.as_ptr()) }); } } } diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 1e8ea8278f0c..ab3982734a28 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -19,7 +19,8 @@ use bun_io::KeepAlive; use bun_jsc::debugger::AsyncTaskTracker; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{ - self as jsc, GlobalRef, JSGlobalObject, JSValue, JsResult, StringJsc, StrongOptional, + self as jsc, AbortSignalRef, GlobalRef, JSGlobalObject, JSValue, JsResult, StringJsc, + StrongOptional, }; use bun_sys::FdExt; use bun_threading::Mutex; @@ -1303,14 +1304,13 @@ impl FetchTasklet { let Some(signal) = self.signal.take() else { return; }; - // `signal` is a live C++-owned WebCore::AbortSignal*; we hold one ref - // (taken in `fetch.rs` before populating FetchOptions). Order matters: - // cleanNativeBindings first, then unref + pending_activity_unref. - // S008: `AbortSignal` is an `opaque_ffi!` ZST — safe `*const → &`. - let signal = bun_opaque::opaque_deref(signal); + // SAFETY: `signal` carries the `+1` taken in `fetch.rs` before it was + // moved into `FetchOptions`; it was just taken out of `self.signal`, so + // dropping the adopted ref below is the only release of that `+1`. + let signal = unsafe { AbortSignalRef::adopt(signal) }; + // Unregister our listener while we still hold the ref. signal.clean_native_bindings(std::ptr::from_mut(self).cast::()); signal.pending_activity_unref(); - signal.unref(); } fn on_reject(&mut self) -> BodyValueError { diff --git a/test/internal/source-lints/refcount-release-owner.test.ts b/test/internal/source-lints/refcount-release-owner.test.ts new file mode 100644 index 000000000000..5026d12bb6f6 --- /dev/null +++ b/test/internal/source-lints/refcount-release-owner.test.ts @@ -0,0 +1,114 @@ +import { file } from "bun"; +import { expect, test } from "bun:test"; +import { realpathSync } from "fs"; +import path from "path"; +import { globAllSources } from "../../../scripts/glob-sources.ts"; + +// Intrusively refcounted foreign objects (`WTF::RefCounted` on the C++ side, +// `napi_env`, ...) are owned from Rust through an RAII handle: `Clone` takes a +// ref, `Drop` releases one, and adopting a raw `+1` into the handle is an +// `unsafe fn` that states the ownership transfer. The FFI shim that releases a +// ref frees the object when the count hits zero, so it may only be called from +// that handle's release hook. Wrapping it in a safe method on the pointee +// (`AbortSignal::unref(&self)`, `AbortSignal::detach(&self, ..)`) let safe code +// release a ref it did not own: the pointee is an `opaque_ffi!` ZST, so +// `&AbortSignal` is obtainable from any pointer (and from `AbortSignalRef`'s +// `Deref`), and `AbortSignal::ref_from_js(v)?.unref()` followed by the +// `AbortSignalRef`'s own `Drop` was a double release without any `unsafe`. +// +// Each entry below is a release shim, the file that declares it, and the one +// function allowed to call it. Every other mention in `src/**/*.rs` is an +// escape hatch that needs to become an owned handle (or an `unsafe` adopt of +// the raw ref into one). Counts are asserted exactly so a renamed shim shows +// up here instead of silently dropping out of the lint. +const RELEASE_SHIMS: Record = { + // `WebCore::AbortSignal`; the owner is `AbortSignalRef` (`ExternalShared`). + WebCore__AbortSignal__unref: { file: "src/jsc/AbortSignal.rs", owner: "ext_deref" }, + // `JSC::ArrayBuffer`; the owner is `ExternalShared`. + JSC__ArrayBuffer__deref: { file: "src/jsc/array_buffer.rs", owner: "ext_deref" }, + // `napi_env`; the owner is `NapiEnvRef` (`ExternalShared`). + NapiEnv__deref: { file: "src/runtime/napi/napi_body.rs", owner: "ext_deref" }, + // `WebCore::CookieMap`; the owner is the hand-rolled `CookieMapRef`. + CookieMap__deref: { file: "src/runtime/webcore/CookieMap.rs", owner: "drop" }, +}; + +const root = path.resolve(import.meta.dir, "..", "..", ".."); +const rustSources = globAllSources().rust.filter(p => p.endsWith(".rs")); + +// Only scan files tracked in HEAD (a `git stash` round-trip can leave stray +// `.rs` files in the working tree; CI runs on a clean checkout). Same guard as +// dead-code-escapes.test.ts. +const tracked: Set | null = (() => { + const r = Bun.spawnSync({ + cmd: ["git", "-C", root, "ls-tree", "-r", "--name-only", "-z", "HEAD"], + stdout: "pipe", + stderr: "ignore", + }); + if (!r.success) return null; + return new Set(r.stdout.toString().split("\0").filter(Boolean)); +})(); + +interface Usage { + declarations: number; + ownerCalls: number; + other: string[]; +} + +const usage = new Map(); +for (const shim of Object.keys(RELEASE_SHIMS)) { + usage.set(shim, { declarations: 0, ownerCalls: 0, other: [] }); +} + +const FN_HEADER = /\bfn\s+([A-Za-z_][A-Za-z0-9_]*)\s*[(<]/; + +let scanned = 0; +for (const abs of rustSources) { + const source = path.relative(root, abs).replaceAll(path.sep, "/"); + // `src/cli` is a symlink into `src/runtime/cli`; count each file once under + // its canonical path. + if (path.relative(root, realpathSync(abs)).replaceAll(path.sep, "/") !== source) continue; + if (tracked !== null && !tracked.has(source)) continue; + scanned++; + const content = await file(abs).text(); + // Every shim name is a `__`-joined identifier that never appears in normal + // prose, so a whole-file substring check skips almost everything. + if (!Object.keys(RELEASE_SHIMS).some(shim => content.includes(shim))) continue; + + // Strip `//` comments (full-line and trailing) so prose mentions don't count. + const lines = content.split("\n").map(line => line.replace(/\/\/.*$/, "")); + for (const [shim, { file: declFile, owner }] of Object.entries(RELEASE_SHIMS)) { + const ident = new RegExp(`\\b${shim}\\b`); + const record = usage.get(shim)!; + for (let i = 0; i < lines.length; i++) { + const line = lines[i]; + if (!ident.test(line)) continue; + const declared = new RegExp(`\\bfn\\s+${shim}\\s*\\(`).test(line); + // The enclosing function is the nearest `fn` header above the call. + let enclosing: string | undefined; + for (let j = i - 1; j >= 0 && enclosing === undefined; j--) { + enclosing = FN_HEADER.exec(lines[j])?.[1]; + } + if (source === declFile && declared) { + record.declarations++; + } else if (source === declFile && enclosing === owner) { + record.ownerCalls++; + } else { + record.other.push(`${source}:${i + 1} (in fn ${enclosing ?? ""}): ${line.trim()}`); + } + } + } +} + +test("scans a non-empty set of tracked Rust sources", () => { + // Guards against the tracked/realpath filters above over-firing (e.g. a + // symlinked checkout root) and leaving nothing to scan, so a failure below + // is attributable to the sources rather than to the scan. Same guard as + // unsound-erased-box.test.ts. + expect(scanned).toBeGreaterThan(0); +}); + +for (const [shim, { owner }] of Object.entries(RELEASE_SHIMS)) { + test(`${shim} is declared once and called only from ${owner}`, () => { + expect(usage.get(shim)).toEqual({ declarations: 1, ownerCalls: 1, other: [] }); + }); +} diff --git a/test/js/bun/spawn/spawn-signal.test.ts b/test/js/bun/spawn/spawn-signal.test.ts index 708789b30edf..7130d09c5747 100644 --- a/test/js/bun/spawn/spawn-signal.test.ts +++ b/test/js/bun/spawn/spawn-signal.test.ts @@ -1,5 +1,5 @@ import { describe, expect, test } from "bun:test"; -import { bunEnv, bunExe } from "harness"; +import { bunEnv, bunExe, isASAN, isDebug } from "harness"; test("spawn AbortSignal works after spawning", async () => { const controller = new AbortController(); @@ -134,6 +134,37 @@ test("spawnSync AbortSignal works as timeout", async () => { expect(end - start).toBeLessThan(100); }); +// The subprocess lets go of the signal when the child exits. The caller still +// holds it, so its timer has to keep running. +describe.concurrent("AbortSignal.timeout() still fires after the child has exited", () => { + // The child must exit before the timer fires, so it skips booting a JS VM + // (`--version`) and slow builds get a wider window (the passing path waits + // it out, so it has to stay inside the per-test ceiling). + const timeoutMs = isDebug || isASAN ? 3000 : 1000; + const cmd = [bunExe(), "--version"]; + + async function expectTimeoutAfterExit(signal: AbortSignal, exitCode: number | null) { + // Proves nothing unless the child was gone before the timer fired. + expect(exitCode).toBe(0); + expect(signal.aborted).toBe(false); + await new Promise(resolve => signal.addEventListener("abort", () => resolve(), { once: true })); + expect(signal.reason).toBeInstanceOf(DOMException); + expect(signal.reason.name).toBe("TimeoutError"); + } + + test("Bun.spawn", async () => { + const signal = AbortSignal.timeout(timeoutMs); + await using proc = Bun.spawn({ cmd, env: bunEnv, stdio: ["ignore", "ignore", "ignore"], signal }); + await expectTimeoutAfterExit(signal, await proc.exited); + }); + + test("Bun.spawnSync", async () => { + const signal = AbortSignal.timeout(timeoutMs); + const { exitCode } = Bun.spawnSync({ cmd, env: bunEnv, stdio: ["ignore", "ignore", "ignore"], signal }); + await expectTimeoutAfterExit(signal, exitCode); + }); +}); + describe("Bun.spawn option validation", () => { const spawners = [ ["Bun.spawn", (opts: any) => Bun.spawn(opts)],