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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 10 additions & 6 deletions src/runtime/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -301,14 +301,15 @@ pub(crate) fn run_task(
.run();
}
task_tag::AsyncCpTask => {
// SAFETY: posted by `on_subtask_done` with the count at zero (exclusive).
unsafe { (*task.ptr.cast::<crate::node::fs::AsyncCpTask>()).run_from_js_thread()? };
// SAFETY: boxed by `schedule_new`; posted once by `on_subtask_done` with
// the count at zero; the arm consumes it.
unsafe { bun_core::heap::take(cast_ptr!(crate::node::fs::AsyncCpTask)) }
.run_from_js_thread()?;
}
task_tag::ShellAsyncCpTask => {
// SAFETY: as above.
unsafe {
(*task.ptr.cast::<crate::node::fs::ShellAsyncCpTask>()).run_from_js_thread()?
};
unsafe { bun_core::heap::take(cast_ptr!(crate::node::fs::ShellAsyncCpTask)) }
.run_from_js_thread()?;
}
task_tag::StatWatcherHop => {
// SAFETY: posted by `StatWatcher::post_to_js_thread` with a ref held.
Expand Down Expand Up @@ -429,7 +430,10 @@ pub(crate) fn run_task(
for_each_fs_uv_op!(__fs_pat) => {
macro_rules! __fs_run {
($($tag:ident $ty:ident;)*) => { match task.tag {
$(task_tag::$tag => cast!(fs_async::$ty).run_from_js_thread()?,)*
// SAFETY: §Dispatch — tag identifies the pointee: the box
// `UVFSRequest::create` leaked, enqueued once; the arm consumes it.
$(task_tag::$tag => unsafe { bun_core::heap::take(cast_ptr!(fs_async::$ty)) }
.run_from_js_thread()?,)*
// SAFETY: outer arm guard proves one of the table tags matched.
_ => unsafe { core::hint::unreachable_unchecked() },
}};
Expand Down
147 changes: 56 additions & 91 deletions src/runtime/node/node_fs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
// The top-level functions assume the arguments are already validated

use bun_paths::strings;
use core::ffi::{c_char, c_int, c_uint, c_void};
use core::ffi::{c_char, c_int, c_uint};
use core::ptr::NonNull;
use core::sync::atomic::{AtomicBool, AtomicUsize, Ordering};

Expand Down Expand Up @@ -648,6 +648,13 @@ mod _async_tasks {
pub(crate) tracker: AsyncTaskTracker,
}

#[cfg(windows)]
impl<R, A: Unprotect, const F: NodeFSFunctionEnum> Drop for UVFSRequest<R, A, F> {
fn drop(&mut self) {
self.r#ref.unref(bun_io::js_vm_ctx());
}
}

#[cfg(windows)]
impl<R: FsReturn, A: FsArgument, const F: NodeFSFunctionEnum> UVFSRequest<R, A, F>
where
Expand Down Expand Up @@ -681,10 +688,7 @@ mod _async_tasks {
r#ref: KeepAlive::default(),
tracker: AsyncTaskTracker::init(vm),
});
// Transfer ownership to libuv: the box outlives the async request and is
// reclaimed in `destroy()` (run_from_js_thread → scopeguard). `heap::release`
// names that hand-off — it is `Box::leak` under the hood; the reclaim
// happens in `destroy()`, not in this scope.
// Reclaimed by the task-queue arm (`run_from_js_thread`) or `release_unrun`.
let task: &mut Self = bun_core::heap::release(task);
// KeepAlive::ref_ now takes the type-erased aio EventLoopCtx; the JS
// event loop is the only one that owns AsyncFSTask/UVFSRequest.
Expand All @@ -693,7 +697,7 @@ mod _async_tasks {
task.tracker.did_schedule(global_object);

let loop_ = uv::Loop::get();
task.req.data = core::ptr::from_mut::<Self>(task).cast::<c_void>();
task.req.data = core::ptr::from_mut::<Self>(task).cast();

// The match resolves at compile time (`F` is a const generic), but
// each arm's body needs `A` re-asserted to its concrete `args::*`
Expand Down Expand Up @@ -942,15 +946,13 @@ mod _async_tasks {
.enqueue_task(bun_jsc::Task::init(this_ptr));
}

pub(crate) fn run_from_js_thread(&mut self) -> Result<(), bun_jsc::JsTerminated> {
// SAFETY: self was Box::leak'd in create(); destroy() runs exactly once on scope exit
let _deinit =
scopeguard::guard(core::ptr::from_mut(self), |p| unsafe { Self::destroy(p) });
// Move `result` out so the `global_object()` `&self` borrow can coexist
// with `&mut result` below; the sentinel left behind is dropped in `destroy()`.
/// Settles the promise; the request (and its keep-alive) is released on return.
#[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")]
pub(crate) fn run_from_js_thread(mut self: Box<Self>) -> Result<(), bun_jsc::JsTerminated> {
// Moved out: `fs_to_js` needs `&mut` while `global_object()` borrows `self`.
let mut result = core::mem::replace(&mut self.result, Err(sys::Error::default()));
let global_object = self.global_object();
let success = matches!(result, Ok(_));
let success = result.is_ok();
let promise_value = self.promise.value();
let promise = self.promise.get();
let result = match &mut result {
Expand Down Expand Up @@ -978,15 +980,6 @@ mod _async_tasks {
}
Ok(())
}

/// SAFETY: `this` must be the pointer Box::leak'd in `create()`; called exactly once.
pub(crate) unsafe fn destroy(this: *mut Self) {
// SAFETY: caller guarantees `this` is the live Box-leaked allocation;
// reclaim ownership (paired with the Box::leak in create()).
let mut task = unsafe { bun_core::heap::take(this) };
// `bun_sys::Error` frees its path on Drop.
task.r#ref.unref(bun_io::js_vm_ctx());
}
}

// ──────────────────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -1220,10 +1213,10 @@ mod _async_tasks {
{
const TAG: bun_event_loop::TaskTag = F.task_tag();
/// A libuv fs request that completed into the queue after the last
/// tick: destroy releases its promise handle and keep-alive.
/// tick: dropping it releases its promise handle and keep-alive.
unsafe fn release_unrun(this: *mut Self) {
// SAFETY: fn contract — `Box::leak`'d in `UVFSRequest::create`.
unsafe { Self::destroy(this) }
// SAFETY: fn contract — leaked in `UVFSRequest::create`, reclaimed once.
unsafe { bun_core::heap::destroy(this) }
}
}

Expand Down Expand Up @@ -1394,6 +1387,14 @@ mod _async_tasks {
pub(crate) shelltask: Option<bun_ptr::ParentRef<ShellCpTask, bun_ptr::Mut>>,
}

impl<const IS_SHELL: bool> Drop for NewAsyncCpTask<IS_SHELL> {
fn drop(&mut self) {
if !IS_SHELL {
self.r#ref.unref(event_loop_handle_to_ctx(self.evtloop));
}
}
}

bun_threading::intrusive_work_task!([const IS_SHELL: bool] NewAsyncCpTask<IS_SHELL>, task);

/// This task is used by `AsyncCpTask/fs.promises.cp` to copy a single file.
Expand All @@ -1403,7 +1404,7 @@ mod _async_tasks {
/// subtask via the `subtask_count` refcount (see `on_subtask_done`). Stored
/// as `ParentRef` (constructed from the `*mut` with `Box::leak` provenance)
/// so shared reads are safe-projected and `as_mut_ptr()` round-trips the
/// original write provenance for `on_subtask_done`'s `&mut` promotion.
/// original write provenance that `on_subtask_done`'s completion frees with.
pub(crate) cp_task: bun_ptr::ParentRef<NewAsyncCpTask<IS_SHELL>, bun_ptr::Mut>,
/// Single owned allocation laid out as `<src>\0<dest>\0`. Ownership is
/// encoded directly as `Box<[OSPathChar]>` and
Expand Down Expand Up @@ -1454,9 +1455,7 @@ mod _async_tasks {
}

fn run_owned(self: Box<Self>) {
// `ParentRef` preserves the `Box::leak` mutable provenance so
// `on_subtask_done` may later promote it to `&mut` via `as_mut_ptr()`
// once the refcount reaches zero.
// `ParentRef` keeps the `Box::leak` provenance the completion frees with.
let cp_task = self.cp_task;
// Shared borrow only — other workpool threads (and the directory-scan
// thread) may hold `&Self` to the same parent concurrently; `ParentRef`
Expand Down Expand Up @@ -1507,11 +1506,11 @@ mod _async_tasks {
} else {
bun_event_loop::task_tag::AsyncCpTask
};
/// A finished fs.cp whose completion will not run: destroy releases
/// A finished fs.cp whose completion will not run: dropping it releases
/// its promise handle, protected arguments and keep-alive.
unsafe fn release_unrun(this: *mut Self) {
// SAFETY: fn contract — posted by `on_subtask_done` with the count at zero.
unsafe { Self::destroy(this) }
unsafe { bun_core::heap::destroy(this) }
}
}

Expand Down Expand Up @@ -1550,7 +1549,9 @@ mod _async_tasks {
tracker,
core::ptr::null_mut(),
);
// SAFETY: `schedule_new` returns a Box::leak'd pointer; valid until destroy()
// SAFETY: `schedule_new` returns the leaked box, which only the completion
// posted to this thread reclaims, so it is live until we return to the
// event loop; the pool side never touches `promise`.
unsafe { &*task }.promise.value()
}

Expand Down Expand Up @@ -1644,10 +1645,7 @@ mod _async_tasks {
/// drops to zero) enqueues `runFromJSThread`, which resolves the promise
/// and destroys `this`.
///
/// Takes a raw `*mut Self` (not `&self`) so the pointer retains the
/// mutable provenance from the original `Box::leak`; the JS-thread
/// callback later materializes `&mut *this`, which would be UB if the
/// pointer were derived from a shared reference.
/// `*mut Self`, not `&self`: the completion frees the box through this pointer.
fn on_subtask_done(this: *mut Self) {
// SAFETY: `this` is a live Box-leaked task; shared access only here —
// other workpool threads may concurrently hold `&Self` until the
Expand All @@ -1666,9 +1664,7 @@ mod _async_tasks {
this_ref.result.set(Ok(()));
}

// Count reached zero ⇒ exclusive access. `this` carries mutable
// provenance from `Box::leak`, so the enqueued callback may safely
// form `&mut *this` on the JS thread.
// Count reached zero ⇒ exclusive access; the completion frees `*this`.
let poster = this_ref.poster.clone();
if poster.is_js() {
let ct = ConcurrentTask::ConcurrentTask::create(bun_jsc::Task::init(this));
Expand All @@ -1677,36 +1673,30 @@ mod _async_tasks {
unreachable!("VM handle closed with an fs.cp outstanding");
};
} else {
let at = AnyTaskWithExtraContext::from_callback_auto_deinit(
this,
|p: *mut Self, ctx| {
// SAFETY: subtask count hit zero ⇒ exclusive access to the leaked task.
unsafe { (*p).run_from_js_thread_mini(ctx) }
},
);
let at =
AnyTaskWithExtraContext::from_callback_auto_deinit(this, |p: *mut Self, _| {
debug_assert!(IS_SHELL, "only the shell posts cp tasks to a mini loop");
// SAFETY: count hit zero ⇒ sole pointer to the box `schedule_new` leaked.
let task = unsafe { bun_core::heap::take(p) };
let _ = task.run_from_js_thread(); // shell tasks settle no promise: always `Ok`
});
// `from_callback_auto_deinit` heap-allocates; never null.
poster.post_mini(core::ptr::NonNull::new(at).expect("heap task"));
}
// The pool side is done (`this` may already be freed by its loop).
poster.embedded_work_finished();
}

pub(crate) fn run_from_js_thread_mini(&mut self, _: *mut c_void) {
let _ = self.run_from_js_thread(); // TODO: properly propagate exception upwards
}

pub(crate) fn run_from_js_thread(&mut self) -> Result<(), bun_jsc::JsTerminated> {
/// Settles the promise (shell: continues the `ShellCpTask`); the task is freed on return.
#[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")]
pub(crate) fn run_from_js_thread(self: Box<Self>) -> Result<(), bun_jsc::JsTerminated> {
// `Maybe<ret::Cp>` (= `Maybe<()>`) has a cheap `Ok(())` placeholder.
let mut result = self.result.replace(Ok(()));
if IS_SHELL {
// SAFETY: shelltask is set by create_for_shell and outlives this task
// Move the result out — `Maybe<ret::Cp>` (= `Maybe<()>`) has a cheap
// `Ok(())` placeholder.
let result = core::mem::replace(self.result.get_mut(), Ok(()));
let shelltask = self.shelltask.expect("IS_SHELL ⇒ shelltask").as_mut_ptr();
// SAFETY: shelltask is non-null in the IS_SHELL specialization and
// outlives this task; `cp_on_finish` enqueues it concurrently.
// outlives this task; `cp_on_finish` continues it in place.
unsafe { ShellCpTask::cp_on_finish(shelltask, result) };
// SAFETY: self was Box::leak'd in create*(); destroyed exactly once here
unsafe { Self::destroy(std::ptr::from_mut::<Self>(self)) };
return Ok(());
}
let go_ptr = self.evtloop.global_object();
Expand All @@ -1717,60 +1707,35 @@ mod _async_tasks {
}
// SAFETY: non-null erased *mut JSGlobalObject from the JS event loop vtable.
let global_object: &JSGlobalObject = unsafe { &*go_ptr.cast::<JSGlobalObject>() };
let success = (*self.result.get_mut()).is_ok();
let success = result.is_ok();
let promise_value = self.promise.value();
// Captured as a raw pointer because `Self::destroy(self)` runs *before* the
// resolve/reject. The `JSPromise` itself lives on the JS heap
// and is kept alive past `destroy` by `promise_value.ensure_still_alive()`.
let promise: *mut bun_jsc::JSPromise = self.promise.get();
let result = match self.result.get_mut() {
// SAFETY: `promise` is the sole live reference to the heap `JSPromise`.
Err(err) => match err.to_js_with_async_stack(global_object, unsafe { &*promise }) {
let promise = self.promise.get();
let result = match &mut result {
Err(err) => match err.to_js_with_async_stack(global_object, promise) {
Ok(v) => v,
Err(e) => {
// SAFETY: `promise` points at a GC-rooted JS heap cell; sole live
// reference on this thread (see comment above `let promise`).
return unsafe { &mut *promise }.reject(global_object, Err(e));
return promise.reject(global_object, Err(e));
}
},
Ok(res) => match FsReturn::fs_to_js(res, global_object) {
Ok(v) => v,
Err(e) => {
// SAFETY: `promise` points at a GC-rooted JS heap cell; sole live
// reference on this thread (see comment above `let promise`).
return unsafe { &mut *promise }.reject(global_object, Err(e));
return promise.reject(global_object, Err(e));
}
},
};
promise_value.ensure_still_alive();

let _dispatch = self.tracker.dispatch(global_object);

// SAFETY: self was Box::leak'd in create*(); destroyed exactly once here
unsafe { Self::destroy(std::ptr::from_mut::<Self>(self)) };
if success {
bun_jsc::JSPromise::opaque_mut(promise).resolve(global_object, result)?;
promise.resolve(global_object, result)?;
} else {
bun_jsc::JSPromise::opaque_mut(promise).reject(global_object, Ok(result))?;
promise.reject(global_object, Ok(result))?;
}
Ok(())
}

/// SAFETY: `this` must be the pointer returned by Box::leak in
/// `schedule_new()`; called exactly once.
pub(crate) unsafe fn destroy(this: *mut Self) {
// SAFETY: caller guarantees `this` is the live Box-leaked allocation;
// reclaim ownership (paired with the Box::leak in
// schedule_new()).
let mut task = unsafe { bun_core::heap::take(this) };
if !IS_SHELL {
let ctx = event_loop_handle_to_ctx(task.evtloop);
task.r#ref.unref(ctx);
}
// `Drop for ThreadSafe<args::Cp>` releases the `protect()` taken by
// `to_thread_safe()` when `src`/`dest` are Buffers, so nothing leaks here.
}

/// Directory scanning + clonefile will block this thread, then each individual file copy (what the sync version
/// calls "copy_single_file_sync") will be dispatched as a separate task.
pub(crate) fn cp_async(nodefs: &mut NodeFS, this: *mut Self) {
Expand Down
Loading