From ff3bdd216bc7d0c721581355ff3e587436fa24ee Mon Sep 17 00:00:00 2001 From: Dylan Conway Date: Mon, 21 Sep 2026 09:25:48 +0000 Subject: [PATCH 1/4] event loop: every queued callback is its own task type; ManagedTask is gone ManagedTask was a heap box holding a context pointer and a function pointer, queued under one shared tag. Each of its 19 users now queues its own pointer under its own tag, and the dispatch arm calls the function directly. That removes one allocation, one free and one indirect call per task. Each type's Taskable impl now says how it is freed when its VM stops before the task runs. The boxed payloads are dropped, and the fetch request-drain hop and the HTMLRewriter background pull give back their ref. Before, all of these leaked. --- src/bundler/ParseTask.rs | 29 +-- src/bundler/ServerComponentParseTask.rs | 9 +- src/bundler/bundle_v2.rs | 100 +++++++---- src/event_loop/ConcurrentTask.rs | 59 +++--- src/event_loop/ManagedTask.rs | 83 --------- src/event_loop/lib.rs | 1 - src/jsc/VirtualMachine.rs | 123 +++++++++---- src/jsc/event_loop.rs | 1 - src/jsc/lib.rs | 4 +- src/jsc/virtual_machine_exports.rs | 42 +++-- src/runtime/api/JSBundler.rs | 15 +- src/runtime/api/html_rewriter.rs | 33 +++- src/runtime/dispatch.rs | 188 +++++++++++++++++++- src/runtime/dns_jsc/cares_jsc.rs | 84 +++++---- src/runtime/server/mod.rs | 75 ++++++-- src/runtime/test_runner/bun_test.rs | 32 ++-- src/runtime/valkey_jsc/valkey.rs | 26 +-- src/runtime/webcore/blob/copy_file.rs | 40 +++-- src/runtime/webcore/blob/write_file.rs | 51 ++++-- src/runtime/webcore/fetch.rs | 2 +- src/runtime/webcore/fetch/FetchTasklet.rs | 29 ++- src/runtime/webview/ChromeProcess.rs | 32 ++-- test/js/bun/http/bun-serve-html-405.test.ts | 2 +- test/js/node/watch/fs.watchFile.test.ts | 2 +- 24 files changed, 676 insertions(+), 386 deletions(-) delete mode 100644 src/event_loop/ManagedTask.rs diff --git a/src/bundler/ParseTask.rs b/src/bundler/ParseTask.rs index b70ee99f16f5..35ae7a7c825e 100644 --- a/src/bundler/ParseTask.rs +++ b/src/bundler/ParseTask.rs @@ -128,7 +128,7 @@ pub enum ParseTaskStage { // ─────────────────────────────────────────────────────────────────────────── /// The information returned to the Bundler thread when a parse finishes. -pub(crate) struct Result { +pub struct Result { pub(crate) task: EventLoop::Task, pub(crate) ctx: bun_ptr::ParentRef, bun_ptr::Mut>, pub(crate) value: ResultValue, @@ -138,6 +138,19 @@ pub(crate) struct Result { /// returned source code by the plugin. pub(crate) external: ExternalFreeFunction, } +impl bun_event_loop::Taskable for Result { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::BundleV2ParseTaskResult; + /// The VM that runs the bundle is going: the result goes to nobody. + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: fn contract — the box `run_from_thread_pool_impl` (or + // `ServerComponentParseTask`) leaked. + drop(unsafe { bun_core::heap::take(this) }); + } + /// A step of the bundle. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} // `Result` lives in a bump arena (no Drop on free); boxing the large arm // would leak the heap allocation. The size diff is acceptable. #[allow(clippy::large_enum_variant)] @@ -2889,13 +2902,7 @@ pub mod parse_worker { .expect("BundleV2.linker.loop must be set before scheduling ParseTask") { bun_event_loop::AnyEventLoop::Js { .. } => { - let ct = - bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback(result, |p| { - // SAFETY: `p` is the `result` Box leaked above; ownership - // transfers to `on_complete`, which deallocates it. - unsafe { on_complete(p) }; - Ok(()) - }); + let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from(result); let poster = worker .ctx .js_poster @@ -2906,7 +2913,7 @@ pub mod parse_worker { // SAFETY: refused ⇒ we own the task box and the leaked result. unsafe { bun_event_loop::ConcurrentTask::ConcurrentTask::release_refused(ct); - drop(bun_core::heap::take(result)); + ::release_unrun(result); } } } @@ -2961,7 +2968,7 @@ pub mod parse_worker { /// (or `ServerComponentParseTask`'s equivalent). Ownership transfers to /// this fn, which deallocates `result` before returning. Must run on the /// main/bundler thread (it dereferences `result.ctx` mutably). - pub(crate) unsafe fn on_complete(result: *mut Result) { + pub unsafe fn on_complete(result: *mut Result) { // SAFETY: result allocated via heap::alloc above; uniquely owned here. let r = unsafe { &mut *result }; let ctx = r.ctx; @@ -2978,4 +2985,4 @@ pub mod parse_worker { } } // end mod parse_worker -pub(crate) use parse_worker::on_complete; +pub use parse_worker::on_complete; diff --git a/src/bundler/ServerComponentParseTask.rs b/src/bundler/ServerComponentParseTask.rs index b2f10b316297..02e76cd73c8c 100644 --- a/src/bundler/ServerComponentParseTask.rs +++ b/src/bundler/ServerComponentParseTask.rs @@ -117,12 +117,7 @@ fn task_callback_wrap(thread_pool_task: *mut ThreadPoolTask) { .expect("BundleV2.linker.loop must be set before scheduling ServerComponentParseTask") { bun_event_loop::AnyEventLoop::Js { .. } => { - let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback(result, |p| { - // SAFETY: `p` is the `result` Box leaked above; ownership - // transfers to `on_complete`, which deallocates it. - unsafe { on_complete(p) }; - Ok(()) - }); + let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from(result); let poster = worker .ctx .js_poster @@ -133,7 +128,7 @@ fn task_callback_wrap(thread_pool_task: *mut ThreadPoolTask) { // SAFETY: refused ⇒ we own the task box and the leaked result. unsafe { bun_event_loop::ConcurrentTask::ConcurrentTask::release_refused(ct); - drop(bun_core::heap::take(result)); + ::release_unrun(result); } } } diff --git a/src/bundler/bundle_v2.rs b/src/bundler/bundle_v2.rs index 37994f596d5c..5c17a111a4a1 100644 --- a/src/bundler/bundle_v2.rs +++ b/src/bundler/bundle_v2.rs @@ -1143,6 +1143,27 @@ pub mod bv2_impl { unsafe { (*(*this).bv2).plugin_context } } } + /// The plugins answered a [`Resolve`]: the hop back to the loop that runs the bundle, when + /// that is a JS loop. Same pointer as the request, its own tag. + #[repr(transparent)] + pub struct ResolveAnswered(Resolve); + impl bun_event_loop::Taskable for ResolveAnswered { + const TAG: bun_event_loop::TaskTag = + bun_event_loop::task_tag::BundleV2PluginResolveAnswered; + /// The VM that runs the bundle is going: the answer goes to nobody. + unsafe fn release_unrun(_: *mut Self) {} + /// A step of the bundle. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } + } + impl ResolveAnswered { + pub fn run(&mut self) { + // SAFETY: `bv2` is a live backref set in `Resolve::init`. + let bv2 = unsafe { &mut *self.0.bv2 }; + BundleV2::on_resolve(&mut self.0, bv2); + } + } impl Resolve { pub(crate) fn init(bv2: &mut BundleV2<'_>, record: MiniImportRecord) -> Self { Self { @@ -1352,6 +1373,47 @@ pub mod bv2_impl { unsafe { (*(*this).bv2).plugin_context } } } + /// The plugins answered a [`Load`]: as [`ResolveAnswered`]. + #[repr(transparent)] + pub struct LoadAnswered(Load); + impl bun_event_loop::Taskable for LoadAnswered { + const TAG: bun_event_loop::TaskTag = + bun_event_loop::task_tag::BundleV2PluginLoadAnswered; + /// As `ResolveAnswered::release_unrun`. + unsafe fn release_unrun(_: *mut Self) {} + /// A step of the bundle. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } + } + impl LoadAnswered { + pub fn run(&mut self) { + // SAFETY: `bv2` is a live backref set in `Load::init`. + let bv2 = unsafe { &mut *self.0.bv2 }; + BundleV2::on_load(&mut self.0, bv2); + } + } + /// A plugin `.defer()`red a [`Load`]: the notice to the loop that runs the bundle, when + /// that is a JS loop. Same pointer as the request, its own tag. + #[repr(transparent)] + pub struct LoadDeferred(Load); + impl bun_event_loop::Taskable for LoadDeferred { + const TAG: bun_event_loop::TaskTag = + bun_event_loop::task_tag::BundleV2PluginLoadDeferred; + /// As `ResolveAnswered::release_unrun`. + unsafe fn release_unrun(_: *mut Self) {} + /// A step of the bundle. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } + } + impl LoadDeferred { + pub fn run(&mut self) { + // SAFETY: `bv2` is a live backref set in `Load::init`. + let bv2 = unsafe { &mut *self.0.bv2 }; + BundleV2::on_notify_defer(&mut self.0, bv2); + } + } } } @@ -4509,9 +4571,8 @@ pub mod bv2_impl { // mutate `graph` / allocate from `graph.heap` off-thread. match self.any_loop_mut() { bun_event_loop::AnyEventLoop::Js { .. } => { - let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback( - std::ptr::from_mut(load), - on_load_from_js_loop_raw, + let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from( + std::ptr::from_mut(load).cast::(), ); let poster = self .js_poster @@ -4543,9 +4604,8 @@ pub mod bv2_impl { // See `on_load_async` — must dispatch on the bundler's own loop. match self.any_loop_mut() { bun_event_loop::AnyEventLoop::Js { .. } => { - let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback( - std::ptr::from_mut(resolve), - on_resolve_from_js_loop_raw, + let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from( + std::ptr::from_mut(resolve).cast::(), ); let poster = self .js_poster @@ -4586,20 +4646,6 @@ pub mod bv2_impl { BundleV2::on_resolve(unsafe { &mut *resolve }, unsafe { &mut *this }); } - fn on_load_from_js_loop(load: &mut jsc_api::JSBundler::Load) { - // SAFETY: `bv2` is a live backref set in `Load::init`. - let bv2 = unsafe { &mut *load.bv2 }; - BundleV2::on_load(load, bv2); - } - - fn on_load_from_js_loop_raw( - load: *mut jsc_api::JSBundler::Load, - ) -> bun_event_loop::JsResult<()> { - // SAFETY: `load` is a valid pointer set up by `from_callback`. - on_load_from_js_loop(unsafe { &mut *load }); - Ok(()) - } - impl<'a> BundleV2<'a> { pub(crate) fn on_load(load: &mut jsc_api::JSBundler::Load, this: &mut BundleV2) { if load.deferred_in.take() == Some(this.graph.defer_epoch) { @@ -4787,20 +4833,6 @@ pub mod bv2_impl { } } - fn on_resolve_from_js_loop(resolve: &mut jsc_api::JSBundler::Resolve) { - // SAFETY: `bv2` is a live backref set in `Resolve::init`. - let bv2 = unsafe { &mut *resolve.bv2 }; - BundleV2::on_resolve(resolve, bv2); - } - - fn on_resolve_from_js_loop_raw( - resolve: *mut jsc_api::JSBundler::Resolve, - ) -> bun_event_loop::JsResult<()> { - // SAFETY: `resolve` is a valid pointer set up by `from_callback`. - on_resolve_from_js_loop(unsafe { &mut *resolve }); - Ok(()) - } - impl<'a> BundleV2<'a> { /// Re-run the idempotent barrel seeding pass after a plugin `onResolve` result patches a record that the importer's parse-completion pass saw unresolved and skipped (#40606). fn schedule_barrel_imports_after_plugin_resolve( diff --git a/src/event_loop/ConcurrentTask.rs b/src/event_loop/ConcurrentTask.rs index 20baa0470e6d..bc82a778b9d0 100644 --- a/src/event_loop/ConcurrentTask.rs +++ b/src/event_loop/ConcurrentTask.rs @@ -8,7 +8,6 @@ //! If `auto_delete` is true, the task is automatically deallocated when it's finished. //! Otherwise, it's expected that the containing struct will deallocate the task. -use crate::ManagedTask; use bun_threading::UnboundedQueue; use bun_threading::unbounded_queue::{Link, Linked}; @@ -68,6 +67,14 @@ pub mod task_tag { BundleV2DeferredBatchTask, // bun.bundle_v2.DeferredBatchTask BundleV2PluginResolve, // bun.bundle_v2.Resolve (JS-thread hop) BundleV2PluginLoad, // bun.bundle_v2.Load (JS-thread hop) + BundleV2PluginResolveAnswered, // bun.bundle_v2.Resolve (hop back to the bundle's loop) + BundleV2PluginLoadAnswered, // bun.bundle_v2.Load (hop back to the bundle's loop) + BundleV2PluginLoadDeferred, // bun.bundle_v2.Load (`.defer()` notice to the bundle's loop) + BundleV2ParseTaskResult, // bun.bundle_v2.ParseTask.Result + ChromePipeEvent, + CopyFileWindowsMkdirp, + WriteFileWindowsMkdirp, + DnsErrorDeferred, ShellYesTask, // shell.Interpreter.Builtin.Yes.YesTask Close, CppTask, @@ -75,13 +82,18 @@ pub mod task_tag { FetchTasklet, FetchTaskletDeinit, FetchTaskletPromiseSettle, + FetchTaskletRequestDrain, FSWatchTask, GetAddrInfoLibuvComplete, + GraphContextStopAgain, + GraphContextStopAndFree, + DeadContextStopAgain, + HandledPromise, HotReloadTask, + HTMLRewriterBackgroundPull, WatchReloadTask, JSBundleCompletionTask, JSCDeferredWorkTask, - ManagedTask, NapiAsyncWork, // napi_async_work NapiFinalizerTask, NativePromiseContextDeferredDerefTask, @@ -96,11 +108,18 @@ pub mod task_tag { Read, Readv, FlushPendingFileSinkTask, + RunTestsTask, RuntimeTranspilerStore, S3HttpDownloadStreamingTask, S3HttpSimpleTask, SendQueueDeferred, // bun_runtime::ipc::SendQueue (close / after-close hop) ServerAllConnectionsClosedTask, + HTTPServerDeinit, + HTTPSServerDeinit, + DebugHTTPServerDeinit, + DebugHTTPSServerDeinit, + HTTPAppClose, + HTTPSAppClose, ShellAsync, ShellCondExprStatTask, ShellCpTask, @@ -120,6 +139,7 @@ pub mod task_tag { StreamPending, ThreadSafeFunction, ValkeyDeferredClose, + ValkeyDeferredFailure, WindowsNamedPipeContext, Write, Writev, @@ -250,19 +270,6 @@ impl Task { } } -// Taskable impls for the low-tier task wrappers defined in this crate. -impl Taskable for crate::ManagedTask::ManagedTask { - const TAG: TaskTag = task_tag::ManagedTask; - unsafe fn release_unrun(this: *mut Self) { - // SAFETY: fn contract — a queued ManagedTask is the heap box `new*` made. - unsafe { crate::ManagedTask::ManagedTask::release(this) } - } - /// A callback task always runs: a callback that continues some script enters that script's - /// context itself, so what it reports goes to nobody once the context has stopped. - unsafe fn context(_this: *const Self) -> ContextId { - ContextId::NONE - } -} // ──────────────────────────────────────────────────────────────────────────── #[repr(C)] @@ -347,16 +354,6 @@ impl ConcurrentTask { Self::create(Task::init(task)) } - // callback returns `JsResult<()>` to match `ManagedTask::new`'s stored ABI; - // callers that have a `fn(*mut T)` should wrap it as `|p| { f(p); Ok(()) }` at the call site. - pub fn from_callback( - ptr: *mut T, - callback: fn(*mut T) -> crate::JsResult<()>, - ) -> core::ptr::NonNull { - bun_core::mark_binding!(); - Self::create(ManagedTask::ManagedTask::new(ptr, callback)) - } - pub fn from( &mut self, of: *mut T, @@ -388,20 +385,14 @@ impl ConcurrentTask { } /// A weak poster got `task` back because the target VM has closed: free - /// it if it is a heap task (`create*`); an intrusive one belongs to its - /// container. + /// the carrier if it is a heap one (`create*`); an intrusive one belongs to + /// its container. What the task points at stays the poster's. /// /// # Safety /// `task` was just refused and is not queued anywhere. pub unsafe fn release_refused(task: core::ptr::NonNull) { // SAFETY: fn contract. - let inner = unsafe { Self::into_task(task) }; - // A callback task (`from_callback`, `ManagedTask::new*`) owns a heap - // `ManagedTask` behind `task.ptr` as well. - if inner.tag == crate::task_tag::ManagedTask { - // SAFETY: as above; refused ⇒ ours. - unsafe { crate::ManagedTask::ManagedTask::release(inner.ptr.cast()) }; - } + let _ = unsafe { Self::into_task(task) }; } /// Returns whether this task should be automatically deallocated after execution. diff --git a/src/event_loop/ManagedTask.rs b/src/event_loop/ManagedTask.rs deleted file mode 100644 index 67fff2477689..000000000000 --- a/src/event_loop/ManagedTask.rs +++ /dev/null @@ -1,83 +0,0 @@ -//! This is a slow, dynamically-allocated one-off task -//! Use it when you can't add to jsc.Task directly and managing the lifetime of the Task struct is overly complex - -use core::ffi::c_void; -use core::ptr::NonNull; - -use crate::{JsResult, Task}; - -pub struct ManagedTask { - // Opaque userdata pointer round-tripped through `new`/`run`; raw by design. - pub ctx: Option>, - pub(crate) callback: fn(*mut c_void) -> JsResult<()>, - pub cleanup: Option, -} - -impl ManagedTask { - pub(crate) fn task(this: *mut ManagedTask) -> Task { - // Per §Dispatch (tag+ptr), name the tag explicitly. - Task::init(this) - } - - /// # Safety - /// `this` must be the live `*mut ManagedTask` embedded in a `Task` returned - /// by `new()`/`new_owned()`; ownership transfers — `this` is freed (via - /// `heap::take`) before return on both Ok and Err paths. - pub unsafe fn run(this: *mut ManagedTask) -> JsResult<()> { - // SAFETY: `this` was produced by `heap::into_raw` in `new`/`new_owned` - // (caller contract). Reconstituting the Box here frees it at scope - // exit on both the Ok and Err paths. - let this = unsafe { bun_core::heap::take(this) }; - let callback = this.callback; - let ctx = this.ctx; - callback(ctx.unwrap().as_ptr()) - } - - /// Free without running: the owned context (if `new_owned`) is dropped. - /// - /// # Safety - /// As [`run`](Self::run); the task is not queued anywhere. - pub unsafe fn release(this: *mut ManagedTask) { - // SAFETY: fn contract. - let this = unsafe { bun_core::heap::take(this) }; - if let (Some(cleanup), Some(ctx)) = (this.cleanup, this.ctx) { - cleanup(ctx.as_ptr()); - } - } - - // A per-(Type, Callback) trampoline is folded away by storing - // the type-erased fn pointer directly — `fn(*mut T)` and `fn(*mut c_void)` share ABI. - pub fn new(ctx: *mut T, callback: fn(*mut T) -> JsResult<()>) -> Task { - let managed = bun_core::heap::into_raw(Box::new(ManagedTask { - // SAFETY: `fn(*mut T) -> R` and `fn(*mut c_void) -> R` have identical - // ABI for all `T: Sized`; `run` passes back the exact pointer stored - // in `ctx` below, so the callee observes its original `*mut T`. - callback: unsafe { - bun_ptr::cast_fn_ptr:: JsResult<()>, fn(*mut c_void) -> JsResult<()>>( - callback, - ) - }, - ctx: NonNull::new(ctx.cast::()), - cleanup: None, - })); - ManagedTask::task(managed) - } - - pub fn new_owned(ctx: *mut T, callback: fn(*mut T) -> JsResult<()>) -> Task { - fn drop_ctx(p: *mut c_void) { - // SAFETY: `p` is the `heap::into_raw(Box)` stored in `ctx` by `new_owned`. - unsafe { bun_core::heap::destroy(p.cast::()) }; - } - let managed = bun_core::heap::into_raw(Box::new(ManagedTask { - // SAFETY: same fn-pointer ABI cast as `new`. - callback: unsafe { - bun_ptr::cast_fn_ptr:: JsResult<()>, fn(*mut c_void) -> JsResult<()>>( - callback, - ) - }, - ctx: NonNull::new(ctx.cast::()), - cleanup: Some(drop_ctx::), - })); - ManagedTask::task(managed) - } -} diff --git a/src/event_loop/lib.rs b/src/event_loop/lib.rs index bace43ec6ef8..30bd185cf130 100644 --- a/src/event_loop/lib.rs +++ b/src/event_loop/lib.rs @@ -4,7 +4,6 @@ pub mod AnyTaskWithExtraContext; pub mod ConcurrentTask; pub mod DeferredTaskQueue; pub mod EventLoopTimer; -pub mod ManagedTask; // ──────────────────────────────────────────────────────────────────────────── // AnyEventLoop / SpawnSyncEventLoop / MiniEventLoop. diff --git a/src/jsc/VirtualMachine.rs b/src/jsc/VirtualMachine.rs index 52461eaf0526..673f7c77a779 100644 --- a/src/jsc/VirtualMachine.rs +++ b/src/jsc/VirtualMachine.rs @@ -1285,30 +1285,10 @@ impl VirtualMachine { /// stopped context `id`: they go on the next turn of the loop, before the /// context can be freed. pub(crate) fn stop_graph_context_again(&mut self, id: crate::ContextId) { - fn stop_again(id: *mut crate::ContextId) -> crate::JsResult<()> { - // SAFETY: boxed below for this task. - let id = *unsafe { Box::from_raw(id) }; - let vm = VirtualMachine::get().as_mut(); - if let Some(context) = vm.graph_context(id).map(NonNull::from) { - // SAFETY: registered ⇒ not freed. - let _ = unsafe { vm.stop_graph_context(context, crate::StopReason::Disposed) }; - } - Ok(()) - } if id == self.dead_context.id() { - fn stop_dead(vm: *mut VirtualMachine) -> crate::JsResult<()> { - // SAFETY: the VM that queued this task on its own loop. - let dead_context = &unsafe { &*vm }.dead_context; - let _ = dead_context.stop(crate::StopReason::Disposed); - if let Some(hooks) = runtime_hooks() { - // SAFETY: live per-thread VM on the JS thread. - unsafe { (hooks.cancel_timers)(vm, Some(dead_context.id())) }; - } - Ok(()) - } if !self.dead_context.stop_again_is_queued() { - let vm = std::ptr::from_mut(self); - self.enqueue_task(bun_event_loop::ManagedTask::ManagedTask::new(vm, stop_dead)); + let vm = std::ptr::from_mut(self).cast::(); + self.enqueue_task(bun_event_loop::Task::init(vm)); } return; } @@ -1316,10 +1296,8 @@ impl VirtualMachine { .graph_context(id) .is_some_and(|context| !context.stop_again_is_queued()) { - // (Owned: released with the task if the VM goes before it runs.) - self.enqueue_task(bun_event_loop::ManagedTask::ManagedTask::new_owned( - Box::into_raw(Box::new(id)), - stop_again, + self.enqueue_task(bun_event_loop::Task::init( + id.raw() as usize as *mut GraphContextStopAgain )); } } @@ -1402,19 +1380,8 @@ impl VirtualMachine { unsafe { self.free_graph_context(context) }; return; } - fn stop_and_free(context: *mut crate::ScriptExecutionContext) -> crate::JsResult<()> { - let vm = VirtualMachine::get().as_mut(); - // SAFETY: still registered (`destroy` frees what teardown left), so not freed. - unsafe { - let context = NonNull::new_unchecked(context); - let _ = vm.stop_graph_context(context, crate::StopReason::Disposed); - vm.free_graph_context(context); - } - Ok(()) - } - self.enqueue_task(bun_event_loop::ManagedTask::ManagedTask::new( - context.as_ptr(), - stop_and_free, + self.enqueue_task(bun_event_loop::Task::init( + context.as_ptr().cast::(), )); } @@ -2613,6 +2580,84 @@ impl VirtualMachine { } } +/// [`VirtualMachine::stop_graph_context_again`]'s task for a `Bun.ModuleGraph` context; `ptr` +/// packs the [`ContextId`](crate::ContextId), nothing is owned. +pub struct GraphContextStopAgain; + +impl bun_event_loop::Taskable for GraphContextStopAgain { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::GraphContextStopAgain; + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + +impl GraphContextStopAgain { + pub fn run(vm: &mut VirtualMachine, id: crate::ContextId) { + if let Some(context) = vm.graph_context(id).map(NonNull::from) { + // SAFETY: registered ⇒ not freed. + let _ = unsafe { vm.stop_graph_context(context, crate::StopReason::Disposed) }; + } + } +} + +/// [`VirtualMachine::stop_graph_context_again`]'s task for the dead context: same pointer as the +/// VM, its own tag. +#[repr(transparent)] +pub struct DeadContextStopAgain(VirtualMachine); + +impl bun_event_loop::Taskable for DeadContextStopAgain { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::DeadContextStopAgain; + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + +impl DeadContextStopAgain { + /// # Safety + /// `this` is the VM that queued this task on its own loop. + pub unsafe fn run(this: *mut Self) { + let vm = this.cast::(); + // SAFETY: fn contract. + let dead_context = &unsafe { &*vm }.dead_context; + let _ = dead_context.stop(crate::StopReason::Disposed); + if let Some(hooks) = runtime_hooks() { + // SAFETY: live per-thread VM on the JS thread. + unsafe { (hooks.cancel_timers)(vm, Some(dead_context.id())) }; + } + } +} + +/// [`VirtualMachine::release_graph_context`]'s task: same pointer as the context, its own tag. +#[repr(transparent)] +pub struct GraphContextStopAndFree(crate::ScriptExecutionContext); + +impl bun_event_loop::Taskable for GraphContextStopAndFree { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::GraphContextStopAndFree; + /// Still registered: `destroy` frees what teardown left. + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + +impl GraphContextStopAndFree { + /// # Safety + /// `this` is the context `release_graph_context` queued. + pub unsafe fn run(vm: &mut VirtualMachine, this: *mut Self) { + // SAFETY: still registered (`destroy` frees what teardown left), so not freed. + unsafe { + let context = NonNull::new_unchecked(this.cast::()); + let _ = vm.stop_graph_context(context, crate::StopReason::Disposed); + vm.free_graph_context(context); + } + } +} + /// Which exit funnel is running [`VirtualMachine::teardown`]. #[derive(Clone, Copy, PartialEq, Eq)] pub(crate) enum Teardown { diff --git a/src/jsc/event_loop.rs b/src/jsc/event_loop.rs index a44ba7575b59..145bb9efc2af 100644 --- a/src/jsc/event_loop.rs +++ b/src/jsc/event_loop.rs @@ -28,7 +28,6 @@ pub use bun_event_loop::ConcurrentTask::{ self, ConcurrentTask as ConcurrentTaskItem, Queue as ConcurrentQueue, }; pub use bun_event_loop::DeferredTaskQueue::{self, DeferredRepeatingTask}; -pub use bun_event_loop::ManagedTask; pub use bun_event_loop::MiniEventLoop; pub use bun_event_loop::Task; pub use bun_event_loop::any_event_loop::{AnyEventLoop, EventLoopHandle, EventLoopTask}; diff --git a/src/jsc/lib.rs b/src/jsc/lib.rs index 72100756b9d9..64bbac5457e5 100644 --- a/src/jsc/lib.rs +++ b/src/jsc/lib.rs @@ -1240,8 +1240,8 @@ pub use self::event_loop as EventLoop; pub mod job; pub use self::event_loop::{ AnyEventLoop, AnyTaskWithExtraContext, ConcurrentCppTask, ConcurrentTask, CppTask, - DeferredTaskQueue, EventLoopHandle, EventLoopTask, GarbageCollectionController, ManagedTask, - MiniEventLoop, PosixSignalHandle, PosixSignalTask, Stopped, Task, WorkPool, WorkPoolTask, + DeferredTaskQueue, EventLoopHandle, EventLoopTask, GarbageCollectionController, MiniEventLoop, + PosixSignalHandle, PosixSignalTask, Stopped, Task, WorkPool, WorkPoolTask, }; pub use self::job::{Completion, Job, JobContext, JsPtr, JsThread, Protected}; #[cfg(unix)] diff --git a/src/jsc/virtual_machine_exports.rs b/src/jsc/virtual_machine_exports.rs index 19b5e70cb913..15e371b99ce9 100644 --- a/src/jsc/virtual_machine_exports.rs +++ b/src/jsc/virtual_machine_exports.rs @@ -7,7 +7,6 @@ use crate::{ }; use bun_bundler::transpiler::PluginResolver; use bun_core::String as BunString; -use bun_event_loop::ManagedTask::ManagedTask; use bun_sourcemap::SourceProviderMap; use bun_sourcemap::parsed_source_map::AnySourceProvider; @@ -131,7 +130,8 @@ pub fn handle_rejected_promise( jsc_vm.auto_garbage_collect(); } -struct HandledPromiseContext { +/// `Bun__handleHandledPromise`'s hop to the next turn of the loop. +pub struct HandledPromiseTask { // VM-lifetime backref (JSC_BORROW) — `GlobalRef` encapsulates the deref. global_this: crate::GlobalRef, // PORTING.md forbids bare JSValue fields on heap-allocated structs; @@ -140,20 +140,27 @@ struct HandledPromiseContext { promise: Strong, } -impl HandledPromiseContext { - fn callback(context: *mut Self) -> bun_event_loop::JsResult<()> { - // SAFETY: `context` was produced by `heap::alloc` below; we are the - // sole owner and reconstitute the Box to drop it at end of scope. - let context = unsafe { bun_core::heap::take(context) }; - let global: &JSGlobalObject = &context.global_this; +impl HandledPromiseTask { + #[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")] + pub fn run(self: Box) { + let global: &JSGlobalObject = &self.global_this; // JSGlobalObject::bun_vm contract. let _ = global .bun_vm() .as_mut() - .handled_promise(global, context.promise.get()); - // drop(context) — Box freed at scope exit (replaces `default_allocator.destroy`); - // Strong's Drop replaces the explicit `.unprotect()`. - Ok(()) + .handled_promise(global, self.promise.get()); + } +} + +impl bun_event_loop::Taskable for HandledPromiseTask { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::HandledPromise; + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: fn contract — boxed in `handle_handled_promise`. + drop(unsafe { bun_core::heap::take(this) }); + } + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE } } @@ -161,14 +168,15 @@ impl HandledPromiseContext { pub fn handle_handled_promise(global: &JSGlobalObject, promise: &JSPromise) { crate::mark_binding!(); let promise_js = promise.to_js(); - let context = bun_core::heap::into_raw(Box::new(HandledPromiseContext { - global_this: global.into(), - promise: Strong::create(promise_js, global), - })); global .bun_vm() .event_loop_mut() - .enqueue_task(ManagedTask::new(context, HandledPromiseContext::callback)); + .enqueue_task(bun_event_loop::Task::from_boxed(Box::new( + HandledPromiseTask { + global_this: global.into(), + promise: Strong::create(promise_js, global), + }, + ))); } // HOST_EXPORT(Bun__onDidAppendPlugin, c) diff --git a/src/runtime/api/JSBundler.rs b/src/runtime/api/JSBundler.rs index 0f332bee4b83..c44d240f3717 100644 --- a/src/runtime/api/JSBundler.rs +++ b/src/runtime/api/JSBundler.rs @@ -1474,7 +1474,7 @@ pub mod js_bundler { // dependency. Only the JSC-aware bits (`on_defer`, `JSBundlerPlugin__*` // C-ABI exports) live here. pub use bun_bundler::bundle_v2::api::JSBundler::{ - Load, LoadSuccess, LoadValue, Resolve, ResolveSuccess, ResolveValue, + Load, LoadDeferred, LoadSuccess, LoadValue, Resolve, ResolveSuccess, ResolveValue, }; /// `&mut BundleV2` for the live backref stored on `Resolve`/`Load`. @@ -1588,9 +1588,8 @@ pub mod js_bundler { .expect("BundleV2.linker.loop must be set before plugins run"); match &mut *any_loop.as_ptr() { bun_event_loop::AnyEventLoop::Js { .. } => { - let ct = ConcurrentTask::from_callback( - std::ptr::from_mut::(self), - on_notify_defer_js, + let ct = ConcurrentTask::create_from( + std::ptr::from_mut::(self).cast::(), ); let poster = (*ctx.as_mut_ptr()) .js_poster @@ -1615,14 +1614,6 @@ pub mod js_bundler { } } - fn on_notify_defer_js(load: *mut Load) -> bun_event_loop::JsResult<()> { - // SAFETY: task contract — `load` is the live request `on_defer` posted; this runs on the loop - // that runs the bundle (bake: the plugins' own), so it is the bundle thread here. - let load = unsafe { &mut *load }; - BundleV2::on_notify_defer(load, bv2_mut(load.bv2)); - Ok(()) - } - fn on_notify_defer_mini_wrap(load: *mut Load, ctx: *mut BundleV2<'static>) { // SAFETY: callback contract — `load` was passed as the `Context` arg to // `enqueue_task_concurrent_with_extra_ctx`; `ctx` is the bundle-thread diff --git a/src/runtime/api/html_rewriter.rs b/src/runtime/api/html_rewriter.rs index cbcb16857ebc..1ae27cfd784d 100644 --- a/src/runtime/api/html_rewriter.rs +++ b/src/runtime/api/html_rewriter.rs @@ -696,6 +696,27 @@ fn active_sink(global: &JSGlobalObject) -> Option> { /// promise. pub type HTMLRewriterTransform = RewriterPipe; +/// [`RewriterPipe::schedule_background_pull`]'s task: same pointer as the pipe, its own tag. +#[repr(transparent)] +pub struct RewriterPipeBackgroundPull(RewriterPipe); +impl bun_event_loop::Taskable for RewriterPipeBackgroundPull { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::HTMLRewriterBackgroundPull; + /// Give back what `schedule_background_pull` took: the cell's protect and the pipe's ref. + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: the task's ref keeps the allocation live until the `deref_nn` below. + let pipe = BackRef::from(unsafe { NonNull::new_unchecked(this.cast::()) }); + let cell = pipe.cell.get(); + if cell.is_cell() { + cell.unprotect(); + } + RewriterPipe::deref_nn(pipe.into()); + } + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + /// Streaming pipe for one `HTMLRewriter::transform()`: receives input bytes /// via [`SinkHandle::HTMLRewriter`], feeds them through lol-html (suspending /// when a content handler returns a pending Promise), and emits output either @@ -1349,14 +1370,14 @@ impl RewriterPipe { cell.protect(); } self.ref_(); - vm.as_mut() - .enqueue_task(bun_jsc::ManagedTask::ManagedTask::new( - core::ptr::from_ref(self).cast_mut(), - Self::run_background_pull, - )); + vm.as_mut().enqueue_task(bun_event_loop::Task::init( + core::ptr::from_ref(self) + .cast_mut() + .cast::(), + )); } - fn run_background_pull(pipe: *mut RewriterPipe) -> bun_event_loop::JsResult<()> { + pub(crate) fn run_background_pull(pipe: *mut RewriterPipe) -> bun_event_loop::JsResult<()> { // SAFETY: the task's ref (taken in `schedule_background_pull`) keeps // the allocation live until the `deref_nn` below. let this = BackRef::from(unsafe { NonNull::new_unchecked(pipe) }); diff --git a/src/runtime/dispatch.rs b/src/runtime/dispatch.rs index 19cd418548b4..12486f2be83d 100644 --- a/src/runtime/dispatch.rs +++ b/src/runtime/dispatch.rs @@ -25,7 +25,6 @@ #[path = "dispatch_js2native.rs"] pub mod js2native; -use bun_event_loop::ManagedTask::ManagedTask; use bun_event_loop::{Task, task_tag}; // `FilePoll::on_update` dispatch is POSIX-only (the symbol is declared @@ -309,10 +308,135 @@ pub(crate) fn run_task( )) }?; } - task_tag::ManagedTask => { - // SAFETY: `task.ptr` was produced by `heap::alloc` in `ManagedTask::new` - // and enqueued under `task_tag::ManagedTask`; `run` consumes/frees it. - unsafe { ManagedTask::run(cast_ptr!(ManagedTask)) }?; + task_tag::BundleV2PluginResolveAnswered => { + cast!(bun_bundler::bundle_v2::api::JSBundler::ResolveAnswered).run(); + } + task_tag::BundleV2PluginLoadAnswered => { + cast!(bun_bundler::bundle_v2::api::JSBundler::LoadAnswered).run(); + } + task_tag::BundleV2PluginLoadDeferred => { + cast!(bun_bundler::bundle_v2::api::JSBundler::LoadDeferred).run(); + } + task_tag::BundleV2ParseTaskResult => { + // SAFETY: boxed by the parse worker; `on_complete` consumes it. + unsafe { + bun_bundler::parse_task::on_complete(cast_ptr!(bun_bundler::parse_task::Result)) + }; + } + task_tag::FetchTaskletRequestDrain => { + crate::webcore::fetch::FetchTaskletRequestDrain::run(cast_ptr!( + crate::webcore::fetch::FetchTaskletRequestDrain + )); + } + task_tag::HTMLRewriterBackgroundPull => { + crate::api::html_rewriter::RewriterPipe::run_background_pull(cast_ptr!( + crate::api::html_rewriter::RewriterPipe + ))?; + } + task_tag::RunTestsTask => { + // SAFETY: boxed in `run_next_tick`; the arm consumes it. + crate::test_runner::bun_test::RunTestsTask::call(unsafe { + bun_core::heap::take(cast_ptr!(crate::test_runner::bun_test::RunTestsTask)) + })?; + } + task_tag::ValkeyDeferredFailure => { + // SAFETY: boxed at the enqueue site; the arm consumes it. + unsafe { bun_core::heap::take(cast_ptr!(crate::valkey_jsc::valkey::DeferredFailure)) } + .run()?; + } + task_tag::DnsErrorDeferred => { + // SAFETY: boxed in `reject_later`; the arm consumes it. + unsafe { + bun_core::heap::take(cast_ptr!(crate::dns_jsc::cares_jsc::ErrorDeferredTask)) + } + .run()?; + } + task_tag::HandledPromise => { + // SAFETY: boxed in `handle_handled_promise`; the arm consumes it. + unsafe { + bun_core::heap::take(cast_ptr!( + bun_jsc::virtual_machine_exports::HandledPromiseTask + )) + } + .run(); + } + task_tag::GraphContextStopAgain => { + // `ptr` packs the `ContextId`, not a pointer. + bun_jsc::virtual_machine::GraphContextStopAgain::run( + vm, + bun_jsc::ContextId::from_raw(task.ptr as usize as u32), + ); + } + task_tag::DeadContextStopAgain => { + // SAFETY: the VM that queued this task on its own loop. + unsafe { + bun_jsc::virtual_machine::DeadContextStopAgain::run(cast_ptr!( + bun_jsc::virtual_machine::DeadContextStopAgain + )) + }; + } + task_tag::GraphContextStopAndFree => { + // SAFETY: queued by `release_graph_context`; still registered. + unsafe { + bun_jsc::virtual_machine::GraphContextStopAndFree::run( + vm, + cast_ptr!(bun_jsc::virtual_machine::GraphContextStopAndFree), + ) + }; + } + task_tag::HTTPAppClose => { + crate::server::AppCloseTask::::run(cast_ptr!( + crate::server::AppCloseTask + )); + } + task_tag::HTTPSAppClose => { + crate::server::AppCloseTask::::run(cast_ptr!(crate::server::AppCloseTask)); + } + // SAFETY (all four): the unique owning server pointer `schedule_deinit` queued. + task_tag::HTTPServerDeinit => unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )); + }, + task_tag::HTTPSServerDeinit => unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )); + }, + task_tag::DebugHTTPServerDeinit => unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )); + }, + task_tag::DebugHTTPSServerDeinit => unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )); + }, + #[cfg(windows)] + task_tag::CopyFileWindowsMkdirp => { + // SAFETY: the live copy `on_mkdirp_complete_concurrent` posted. + unsafe { + crate::webcore::blob::copy_file::CopyFileWindowsMkdirp::run(cast_ptr!( + crate::webcore::blob::copy_file::CopyFileWindowsMkdirp<'_> + )) + }; + } + #[cfg(windows)] + task_tag::WriteFileWindowsMkdirp => { + // SAFETY: the live write `on_mkdirp_complete_concurrent` posted. + unsafe { + crate::webcore::blob::write_file::WriteFileWindowsMkdirp::run(cast_ptr!( + crate::webcore::blob::write_file::WriteFileWindowsMkdirp + )) + }; + } + #[cfg(windows)] + task_tag::ChromePipeEvent => { + // SAFETY: boxed in `PipeEvent::post`; the arm consumes it. + crate::webview::chrome_process::QueuedEvent::deliver(unsafe { + bun_core::heap::take(cast_ptr!(crate::webview::chrome_process::QueuedEvent)) + })?; } task_tag::CppTask => { cast!(CppTask).run(global)?; @@ -589,7 +713,7 @@ fn run_task_cold(task: Task) { /// `release_task_unrun` track `bun_event_loop::task_tag::COUNT`. Bump when /// adding a variant — and give it an arm in both. const _: () = assert!( - task_tag::COUNT == 61, + task_tag::COUNT == 82, "dispatch::run_task / release_task_unrun arm count out of sync with bun_event_loop::task_tag", ); @@ -1267,7 +1391,39 @@ fn __bun_release_task_unrun(task: bun_event_loop::Task) { release!(crate::api::js_bundle_completion_task::JSBundleCompletionTask) } task_tag::JSCDeferredWorkTask => release!(JSCDeferredWorkTask), - task_tag::ManagedTask => release!(ManagedTask), + task_tag::BundleV2PluginResolveAnswered => { + release!(bun_bundler::bundle_v2::api::JSBundler::ResolveAnswered) + } + task_tag::BundleV2PluginLoadAnswered => { + release!(bun_bundler::bundle_v2::api::JSBundler::LoadAnswered) + } + task_tag::BundleV2PluginLoadDeferred => { + release!(bun_bundler::bundle_v2::api::JSBundler::LoadDeferred) + } + task_tag::BundleV2ParseTaskResult => release!(bun_bundler::parse_task::Result), + task_tag::FetchTaskletRequestDrain => { + release!(crate::webcore::fetch::FetchTaskletRequestDrain) + } + task_tag::HTMLRewriterBackgroundPull => { + release!(crate::api::html_rewriter::RewriterPipeBackgroundPull) + } + task_tag::RunTestsTask => release!(crate::test_runner::bun_test::RunTestsTask), + task_tag::ValkeyDeferredFailure => release!(crate::valkey_jsc::valkey::DeferredFailure), + task_tag::DnsErrorDeferred => release!(crate::dns_jsc::cares_jsc::ErrorDeferredTask), + task_tag::HandledPromise => release!(bun_jsc::virtual_machine_exports::HandledPromiseTask), + task_tag::GraphContextStopAgain => { + release!(bun_jsc::virtual_machine::GraphContextStopAgain) + } + task_tag::DeadContextStopAgain => release!(bun_jsc::virtual_machine::DeadContextStopAgain), + task_tag::GraphContextStopAndFree => { + release!(bun_jsc::virtual_machine::GraphContextStopAndFree) + } + task_tag::HTTPAppClose => release!(crate::server::AppCloseTask), + task_tag::HTTPSAppClose => release!(crate::server::AppCloseTask), + task_tag::HTTPServerDeinit => release!(crate::server::ServerDeinitTask), + task_tag::HTTPSServerDeinit => release!(crate::server::ServerDeinitTask), + task_tag::DebugHTTPServerDeinit => release!(crate::server::ServerDeinitTask), + task_tag::DebugHTTPSServerDeinit => release!(crate::server::ServerDeinitTask), task_tag::NapiAsyncWork => release!(napi_async_work), task_tag::NapiFinalizerTask => release!(NapiFinalizerTask), task_tag::NativePromiseContextDeferredDerefTask => { @@ -1326,6 +1482,24 @@ fn __bun_release_task_unrun(task: bun_event_loop::Task) { #[cfg(not(windows))] unreachable!("windows-only tag"); } + task_tag::CopyFileWindowsMkdirp => { + #[cfg(windows)] + release!(crate::webcore::blob::copy_file::CopyFileWindowsMkdirp<'_>); + #[cfg(not(windows))] + unreachable!("windows-only tag"); + } + task_tag::WriteFileWindowsMkdirp => { + #[cfg(windows)] + release!(crate::webcore::blob::write_file::WriteFileWindowsMkdirp); + #[cfg(not(windows))] + unreachable!("windows-only tag"); + } + task_tag::ChromePipeEvent => { + #[cfg(windows)] + release!(crate::webview::chrome_process::QueuedEvent); + #[cfg(not(windows))] + unreachable!("windows-only tag"); + } task_tag::Open | task_tag::Close | task_tag::Read diff --git a/src/runtime/dns_jsc/cares_jsc.rs b/src/runtime/dns_jsc/cares_jsc.rs index bb1d1dc995b0..53c71c072f45 100644 --- a/src/runtime/dns_jsc/cares_jsc.rs +++ b/src/runtime/dns_jsc/cares_jsc.rs @@ -709,51 +709,57 @@ impl ErrorDeferred { global_this: &JSGlobalObject, context: bun_jsc::ContextId, ) { - struct Context { - /// The context of the script that asked; it may stop before the task runs. - asking: bun_jsc::ContextId, - deferred: Box, - // LIFETIMES.tsv row 1403: JSC_BORROW — the global outlives the - // enqueued task (VM-owned), so a `BackRef` captures the invariant. - global_this: bun_ptr::BackRef, - } - impl Context { - // `bun_event_loop::ManagedTask::new` expects - // `fn(*mut T) -> bun_event_loop::JsResult<()>` (tier-0 `bun_core::JsError`). - fn callback(this: *mut Context) -> bun_event_loop::JsResult<()> { - // SAFETY: `this` is the heap-allocated pointer passed to ManagedTask::new - // below; ManagedTask::run calls us exactly once with that pointer. - let this = unsafe { bun_core::heap::take(this) }; - let global = this.global_this.get(); - // For the script that asked: once its context has stopped the error goes to nobody. - let _context = global.bun_vm().enter_context(this.asking); - this.deferred.reject(global) - } - } - let vm = global_this.bun_vm(); // Worker terminate's `stop_dns_for_vm_teardown` fires EDESTRUCTION with - // `is_shutting_down` already set; the task queue is about to be - // drained-without-run and ManagedTask has no cleanup here, so enqueuing - // would leak the `Context` and its `JSPromiseStrong` box. Drop now while - // JSC is still live so the Strong handle releases cleanly. + // `is_shutting_down` already set: there is nobody to reject for. if vm.is_shutting_down() { return; } - - let context = bun_core::heap::into_raw(Box::new(Context { - asking: context, - deferred: self, - global_this: bun_ptr::BackRef::new(global_this), - })); - // TODO(@heimskr): new custom Task type - // SAFETY: `bun_vm()` returns a non-null VM pointer (VM-owned for the lifetime of - // the JSGlobalObject). vm.as_mut() - .enqueue_task(bun_jsc::ManagedTask::ManagedTask::new_owned( - context, - Context::callback, - )); + .enqueue_task(bun_event_loop::Task::from_boxed(Box::new( + ErrorDeferredTask { + asking: context, + deferred: self, + global_this: bun_ptr::BackRef::new(global_this), + }, + ))); + } +} + +/// [`ErrorDeferred::reject_later`]'s hop to the next turn of the loop. +pub(crate) struct ErrorDeferredTask { + /// The context of the script that asked; it may stop before the task runs. + asking: bun_jsc::ContextId, + deferred: Box, + // LIFETIMES.tsv row 1403: JSC_BORROW — the global outlives the + // enqueued task (VM-owned), so a `BackRef` captures the invariant. + global_this: bun_ptr::BackRef, +} + +impl ErrorDeferredTask { + #[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")] + pub(crate) fn run(self: Box) -> JsResult<()> { + let Self { + asking, + deferred, + global_this, + } = *self; + let global = global_this.get(); + // For the script that asked: once its context has stopped the error goes to nobody. + let _context = global.bun_vm().enter_context(asking); + deferred.reject(global) + } +} + +impl bun_event_loop::Taskable for ErrorDeferredTask { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::DnsErrorDeferred; + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: fn contract — boxed in `reject_later`. + drop(unsafe { bun_core::heap::take(this) }); + } + /// `run` enters the asking script's context itself. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE } } diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index da89a189c0f4..df8bb2eacf01 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1965,27 +1965,15 @@ impl NewServer { // (non-null for the server's lifetime); single-threaded JS // context, `&mut` scoped to this call. unsafe { - (*self.vm_mut()).enqueue_task(bun_event_loop::ManagedTask::ManagedTask::new( - app, - |app| { - // S008: `NewApp` is a ZST opaque — safe `*mut → &mut` deref. - bun_opaque::opaque_deref_mut(app).close(); - Ok(()) - }, - )); + (*self.vm_mut()) + .enqueue_task(bun_event_loop::Task::init(app.cast::>())); } } // SAFETY: as above — `&mut` scoped to this call. unsafe { - (*self.vm_mut()).enqueue_task(bun_event_loop::ManagedTask::ManagedTask::new( - std::ptr::from_mut::(self), - |this| { - // SAFETY: `this` is the unique owning server pointer enqueued - // above; the task runs once on the JS thread. - Self::deinit(this); - Ok(()) - }, + (*self.vm_mut()).enqueue_task(bun_event_loop::Task::init( + std::ptr::from_mut::(self).cast::>(), )); } } @@ -4274,6 +4262,61 @@ pub(crate) enum SavedRequestUnion<'a> { Saved(SavedRequest), } +// ─── schedule_deinit's tasks ───────────────────────────────────────────────── +/// `schedule_deinit`'s first task, `app.close()`: same pointer as the app, one tag per `SSL`. +#[repr(transparent)] +pub struct AppCloseTask(uws_sys::NewApp); + +impl bun_event_loop::Taskable for AppCloseTask { + const TAG: bun_event_loop::TaskTag = if SSL { + bun_event_loop::task_tag::HTTPSAppClose + } else { + bun_event_loop::task_tag::HTTPAppClose + }; + /// The app goes with its server. + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + +impl AppCloseTask { + pub(crate) fn run(this: *mut Self) { + // S008: `NewApp` is a ZST opaque — safe `*mut → &mut` deref. + bun_opaque::opaque_deref_mut(this.cast::>()).close(); + } +} + +/// `schedule_deinit`'s second task, `deinit()`: same pointer as the server, one tag per +/// monomorphization. +#[repr(transparent)] +pub struct ServerDeinitTask(NewServer); + +impl bun_event_loop::Taskable for ServerDeinitTask { + const TAG: bun_event_loop::TaskTag = match (SSL, DEBUG) { + (false, false) => bun_event_loop::task_tag::HTTPServerDeinit, + (true, false) => bun_event_loop::task_tag::HTTPSServerDeinit, + (false, true) => bun_event_loop::task_tag::DebugHTTPServerDeinit, + (true, true) => bun_event_loop::task_tag::DebugHTTPSServerDeinit, + }; + /// Not freed here: `finalize()` frees a server whose deinit was scheduled once the VM is + /// shutting down (`DEINIT_SCHEDULED`). + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + +impl ServerDeinitTask { + /// # Safety + /// `this` is the unique owning server pointer `schedule_deinit` queued. + pub(crate) unsafe fn run(this: *mut Self) { + NewServer::::deinit(this.cast()); + } +} + // ─── ServerAllConnectionsClosedTask ────────────────────────────────────────── pub struct ServerAllConnectionsClosedTask { pub(crate) global_object: *const jsc::JSGlobalObject, diff --git a/src/runtime/test_runner/bun_test.rs b/src/runtime/test_runner/bun_test.rs index f56ef96413ab..fb9ba43720f9 100644 --- a/src/runtime/test_runner/bun_test.rs +++ b/src/runtime/test_runner/bun_test.rs @@ -888,18 +888,11 @@ impl BunTest { debug_assert!(false); // shouldn't be calling runNextTick after moving on to the next file return; // but just in case }; - let done_callback_test = bun_core::heap::into_raw(Box::new(RunTestsTask { + let task = jsc::Task::from_boxed(Box::new(RunTestsTask { weak: Weak::clone(weak), global_this: GlobalRef::from(global_this), phase, })); - fn call_erased(this: *mut RunTestsTask) -> bun_event_loop::JsResult<()> { - // `this` was `heap::into_raw`'d above (always non-null) and is - // invoked exactly once by `ManagedTask`. - RunTestsTask::call(NonNull::new(this).unwrap()) - } - // `new_owned`: if the task never runs (VM teardown), the queue drainer frees `done_callback_test`. - let task = jsc::ManagedTask::ManagedTask::new_owned::(done_callback_test, call_erased); // SAFETY: single field write through `UnsafeCell`; no other `&mut` live. strong.get().wants_wakeup = true; // we need to wake up the event loop so autoTick() doesn't wait for 16-100ms because we just enqueued a task @@ -1519,15 +1512,8 @@ pub struct RunTestsTask { pub(crate) phase: RefDataValue, } impl RunTestsTask { - /// `ManagedTask` callback ABI: `fn(*mut T) -> JsResult<()>`. The pointer - /// was `heap::alloc`'d in `run_next_tick`; reconstitute and drop here. - /// - /// `this` must be the pointer produced by `heap::into_raw` in - /// `run_next_tick`; ownership is consumed (the box is dropped on return). - pub fn call(this: NonNull) -> JsResult<()> { - // SAFETY: `this` was produced by `heap::into_raw` in `run_next_tick` and - // is invoked exactly once by `ManagedTask`; ownership is reclaimed here. - let this = unsafe { bun_core::heap::take(this.as_ptr()) }; + #[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")] + pub fn call(this: Box) -> JsResult<()> { // Box drops at end of scope; the Weak drops with it. let Some(strong) = this.weak.upgrade() else { return Ok(()) }; if let Err(e) = BunTest::run(&strong, &this.global_this) { @@ -1549,6 +1535,18 @@ impl RunTestsTask { } } +impl bun_event_loop::Taskable for RunTestsTask { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::RunTestsTask; + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: fn contract — boxed in `run_next_tick`. + drop(unsafe { bun_core::heap::take(this) }); + } + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + #[derive(Copy, Clone, PartialEq, Eq, strum::IntoStaticStr)] pub enum HandleUncaughtExceptionResult { #[strum(serialize = "hide_error")] diff --git a/src/runtime/valkey_jsc/valkey.rs b/src/runtime/valkey_jsc/valkey.rs index 3460e1b28a80..c4f1faf831b7 100644 --- a/src/runtime/valkey_jsc/valkey.rs +++ b/src/runtime/valkey_jsc/valkey.rs @@ -287,7 +287,7 @@ enum SubscribeHandled { Fallthrough, } -struct DeferredFailure { +pub(crate) struct DeferredFailure { message: Box<[u8]>, err: RedisError, global_this: GlobalRef, @@ -296,7 +296,7 @@ struct DeferredFailure { } impl DeferredFailure { - fn run(self) -> JsResult<()> { + pub(crate) fn run(self) -> JsResult<()> { debug!("running deferred failure"); let mut this = self; let err = valkey_error_to_js(&this.global_this, &*this.message, this.err); @@ -310,17 +310,21 @@ impl DeferredFailure { fn enqueue(self: Box) { debug!("enqueueing deferred failure"); - // The Box is leaked into a raw pointer here and reconstituted inside the trampoline. - fn run_raw(ptr: *mut DeferredFailure) -> bun_event_loop::JsResult<()> { - // SAFETY: `ptr` was produced by `heap::alloc` below; we are the sole owner. - let this = unsafe { bun_core::heap::take(ptr) }; - DeferredFailure::run(*this) - } - let managed_task = - bun_jsc::ManagedTask::ManagedTask::new(bun_core::heap::into_raw(self), run_raw); VirtualMachine::get() .event_loop_mut() - .enqueue_task(managed_task); + .enqueue_task(bun_event_loop::Task::from_boxed(self)); + } +} + +impl bun_event_loop::Taskable for DeferredFailure { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::ValkeyDeferredFailure; + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: fn contract — boxed at the enqueue site. + drop(unsafe { bun_core::heap::take(this) }); + } + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE } } diff --git a/src/runtime/webcore/blob/copy_file.rs b/src/runtime/webcore/blob/copy_file.rs index 312095662e56..ad481a865236 100644 --- a/src/runtime/webcore/blob/copy_file.rs +++ b/src/runtime/webcore/blob/copy_file.rs @@ -1926,19 +1926,39 @@ fn on_mkdirp_complete_concurrent(ctx: *mut (), err_: bun_sys::Maybe<()>, ticket: bun_sys::Result::Err(e) => Some(e), bun_sys::Result::Ok(()) => None, }; - // callback signature to match `ManagedTask::new`'s `fn(*mut T) -> jsc::JsResult<()>`. - fn call_erased(this: *mut CopyFileWindows<'_>) -> bun_event_loop::JsResult<()> { - // SAFETY: `this` is the heap-allocated `CopyFileWindows` passed to - // `ManagedTask::new` below; `on_mkdirp_complete` may free it via `throw`, so we - // do not touch `this` afterward. - unsafe { (*this).on_mkdirp_complete() }; - Ok(()) - } - ticket.post(jsc::ConcurrentTask::create( - jsc::ManagedTask::ManagedTask::new::(this, call_erased), + ticket.post(jsc::ConcurrentTask::create_from( + std::ptr::from_mut(this).cast::>(), )); } +/// `mkdirp` finished on the work pool: the hop back to the JS thread. Same pointer as the copy, its +/// own tag. +#[cfg(windows)] +#[repr(transparent)] +pub struct CopyFileWindowsMkdirp<'a>(CopyFileWindows<'a>); + +#[cfg(windows)] +impl bun_event_loop::Taskable for CopyFileWindowsMkdirp<'_> { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::CopyFileWindowsMkdirp; + /// Frees nothing: the copy is not this task's. + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} + +#[cfg(windows)] +impl CopyFileWindowsMkdirp<'_> { + /// # Safety + /// `this` is the live `CopyFileWindows` `on_mkdirp_complete_concurrent` posted; + /// `on_mkdirp_complete` may free it via `throw`. + pub(crate) unsafe fn run(this: *mut Self) { + // SAFETY: fn contract. + unsafe { (*this).0.on_mkdirp_complete() }; + } +} + // ─────────────────────────────────────────────────────────────────────────── // IOWhich + module-level constants // ─────────────────────────────────────────────────────────────────────────── diff --git a/src/runtime/webcore/blob/write_file.rs b/src/runtime/webcore/blob/write_file.rs index 6cd32d81f56f..25c84de4e073 100644 --- a/src/runtime/webcore/blob/write_file.rs +++ b/src/runtime/webcore/blob/write_file.rs @@ -538,7 +538,9 @@ impl WriteFile { // ────────────────────────────────────────────────────────────────────────── #[cfg(windows)] -pub(crate) use self::windows_impl::{WriteFileWindows, WriteFileWindowsError}; +pub(crate) use self::windows_impl::{ + WriteFileWindows, WriteFileWindowsError, WriteFileWindowsMkdirp, +}; #[cfg(windows)] mod windows_impl { @@ -546,9 +548,9 @@ mod windows_impl { use core::ptr::null_mut; use bun_io::{self as aio, IntrusiveUvFs as _, KeepAlive}; - // `bun_jsc::EventLoop`/`ManagedTask` are *modules* (namespace - // re-exports); the structs live one level deeper. - use bun_jsc::{ConcurrentTask, ManagedTask::ManagedTask, event_loop::EventLoop}; + // `bun_jsc::EventLoop` is a *module* (namespace re-export); the struct + // lives one level deeper. + use bun_jsc::{ConcurrentTask, event_loop::EventLoop}; use bun_sys::ReturnCodeExt as _; use bun_sys::windows::libuv as uv; @@ -589,6 +591,31 @@ mod windows_impl { } } + /// `mkdirp` finished on the work pool: the hop back to the JS thread. Same pointer as the + /// write, its own tag. + #[repr(transparent)] + pub(crate) struct WriteFileWindowsMkdirp(WriteFileWindows); + + impl bun_event_loop::Taskable for WriteFileWindowsMkdirp { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::WriteFileWindowsMkdirp; + /// Frees nothing: the write is not this task's. + unsafe fn release_unrun(_: *mut Self) {} + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } + } + + impl WriteFileWindowsMkdirp { + /// # Safety + /// `this` is the live `WriteFileWindows` `on_mkdirp_complete_concurrent` posted; + /// `on_mkdirp_complete` may free it. + pub(crate) unsafe fn run(this: *mut Self) { + // SAFETY: fn contract. + unsafe { WriteFileWindows::on_mkdirp_complete(this.cast::()) }; + } + } + impl WriteFileWindows { pub(crate) fn create_with_ctx( file_blob: Blob, @@ -919,18 +946,6 @@ mod windows_impl { } } - /// `ManagedTask`-shaped trampoline for [`on_mkdirp_complete`]: takes - /// `*mut Self` and returns the event-loop `jsc::JsResult<()>` (always `Ok`: the inner body - /// reports a delivery exception itself). - fn on_mkdirp_complete_task(this: *mut WriteFileWindows) -> bun_event_loop::JsResult<()> { - // SAFETY: `this` is the live Box-allocated `WriteFileWindows` whose - // pointer was stashed in `on_mkdirp_complete_concurrent` below; - // the JS thread is the sole accessor at this point. `*this` may be - // freed inside; not accessed afterward. - unsafe { Self::on_mkdirp_complete(this) }; - Ok(()) - } - fn on_mkdirp_complete_concurrent( ctx: *mut (), err_: bun_sys::Result<()>, @@ -945,8 +960,8 @@ mod windows_impl { bun_sys::Result::Err(e) => Some(e), bun_sys::Result::Ok(()) => None, }; - ticket.post(ConcurrentTask::create( - ManagedTask::new::(this, Self::on_mkdirp_complete_task), + ticket.post(ConcurrentTask::create_from( + std::ptr::from_mut(this).cast::(), )); } diff --git a/src/runtime/webcore/fetch.rs b/src/runtime/webcore/fetch.rs index ce129c78bee5..20ad92b232a6 100644 --- a/src/runtime/webcore/fetch.rs +++ b/src/runtime/webcore/fetch.rs @@ -81,7 +81,7 @@ use bun_url::PercentEncoding; use bun_url::URL as ZigURL; use self::fetch_tasklet::{FetchOptions, HTTPRequestBody}; -pub use self::fetch_tasklet::{FetchTasklet, FetchTaskletDeinitHop}; +pub use self::fetch_tasklet::{FetchTasklet, FetchTaskletDeinitHop, FetchTaskletRequestDrain}; // ────────────────────────────────────────────────────────────────────────── // Local extension shims (upstream methods not yet ported / not in scope) diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index ae77f31fa4bd..9bc6937fbdb0 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -37,7 +37,6 @@ use bun_jsc::AbortSignalRef; // `bun_event_loop::JsResult` (cycle-broken erased error) — used by // ConcurrentTask callbacks at the tier-3 layer. -type ElJsResult = bun_event_loop::JsResult; use http::signals::BODY_HIGH_WATER_MARK; @@ -71,6 +70,27 @@ impl FetchTaskletDeinitHop { } } +/// The HTTP thread drained the request body's buffer: the hop that tells the sink, on the JS +/// thread. Same pointer, its own tag. +#[repr(transparent)] +pub struct FetchTaskletRequestDrain(FetchTasklet); +impl Taskable for FetchTaskletRequestDrain { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::FetchTaskletRequestDrain; + /// Carries the +1 `on_write_request_data_drain` took. + unsafe fn release_unrun(this: *mut Self) { + FetchTasklet::deref(this.cast::()); + } + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE + } +} +impl FetchTaskletRequestDrain { + pub(crate) fn run(this: *mut Self) { + FetchTasklet::resume_request_data_stream(this.cast::()); + } +} + impl Taskable for FetchTasklet { const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::FetchTasklet; /// A progress hop the HTTP thread posted: it carries the +1 that @@ -2226,8 +2246,7 @@ impl FetchTasklet { let this_ref = Self::from_raw_ref(this); // ref until the main thread callback is called this_ref.ref_(); - // `from_callback` heap-allocates a fresh `ConcurrentTaskItem`. - let task = ConcurrentTask::from_callback(this, FetchTasklet::resume_request_data_stream); + let task = ConcurrentTask::create_from(this.cast::()); this_ref .http_ticket .as_ref() @@ -2236,8 +2255,7 @@ impl FetchTasklet { } /// This is ALWAYS called from the main thread - // ConcurrentTask::from_callback expects `fn(*mut T) -> bun_event_loop::JsResult<()>`. - fn resume_request_data_stream(this: *mut FetchTasklet) -> ElJsResult<()> { + fn resume_request_data_stream(this: *mut FetchTasklet) { let this_ref = Self::from_raw_mut(this); bun_output::scoped_log!(FetchTasklet, "resumeRequestDataStream"); if !this_ref.signal_aborted() { @@ -2249,7 +2267,6 @@ impl FetchTasklet { // deref when done because we ref inside onWriteRequestDataDrain // SAFETY: `this` is the live heap tasklet; we hold a ref. FetchTasklet::deref(this); - Ok(()) } /// True for upgraded connections, HTTP/2 (DATA frames) and `Content-Length` framing. diff --git a/src/runtime/webview/ChromeProcess.rs b/src/runtime/webview/ChromeProcess.rs index d74c5833ab0f..8d575ea305f1 100644 --- a/src/runtime/webview/ChromeProcess.rs +++ b/src/runtime/webview/ChromeProcess.rs @@ -923,7 +923,7 @@ enum PipeEvent { } #[cfg(windows)] -struct QueuedEvent { +pub(crate) struct QueuedEvent { generation: u32, event: PipeEvent, } @@ -931,25 +931,33 @@ struct QueuedEvent { #[cfg(windows)] impl PipeEvent { fn post(self, generation: u32) { - let queued = bun_core::heap::into_raw(Box::new(QueuedEvent { - generation, - event: self, - })); // Not dispatched from the read callback: C++ runs JS that may spin a nested event loop (bun:test does), and libuv re-arms the read only after the callback returns. VirtualMachine::get() .as_mut() - .enqueue_task(bun_jsc::ManagedTask::ManagedTask::new_owned( - queued, - QueuedEvent::deliver, - )); + .enqueue_task(bun_jsc::Task::from_boxed(Box::new(QueuedEvent { + generation, + event: self, + }))); + } +} + +#[cfg(windows)] +impl bun_event_loop::Taskable for QueuedEvent { + const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::ChromePipeEvent; + unsafe fn release_unrun(this: *mut Self) { + // SAFETY: fn contract — boxed in `PipeEvent::post`. + drop(unsafe { bun_core::heap::take(this) }); + } + /// Enters no context. + unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { + bun_event_loop::ContextId::NONE } } #[cfg(windows)] impl QueuedEvent { - fn deliver(this: *mut QueuedEvent) -> bun_jsc::JsResult<()> { - // SAFETY: the box leaked by `post`; ManagedTask hands it over once. - let queued = unsafe { bun_core::heap::take(this) }; + #[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")] + pub(crate) fn deliver(queued: Box) -> bun_jsc::JsResult<()> { if queued.generation != GENERATION.load(Ordering::Relaxed) { scoped_log!( Chrome, diff --git a/test/js/bun/http/bun-serve-html-405.test.ts b/test/js/bun/http/bun-serve-html-405.test.ts index b94d76f72441..6b9faa95a92a 100644 --- a/test/js/bun/http/bun-serve-html-405.test.ts +++ b/test/js/bun/http/bun-serve-html-405.test.ts @@ -45,7 +45,7 @@ test("dev server html route: non-GET/HEAD requests complete without hanging", as }); // A server whose JS wrapper survives to lastChanceToFinalize used to leak its -// NewServer Box: finalize() -> schedule_deinit() enqueued a ManagedTask that +// NewServer Box: finalize() -> schedule_deinit() enqueued a task that // the now-exiting event loop never ran. The html_bundle::Route.server // back-pointer makes the orphaned Box a pointer cycle, so LSan reports every // interior allocation as an indirect leak. Only observable via LeakSanitizer. diff --git a/test/js/node/watch/fs.watchFile.test.ts b/test/js/node/watch/fs.watchFile.test.ts index 60f4320e61de..dd181f361f3d 100644 --- a/test/js/node/watch/fs.watchFile.test.ts +++ b/test/js/node/watch/fs.watchFile.test.ts @@ -333,7 +333,7 @@ describe("fs.watchFile", () => { cmd: [bunExe(), "-e", fixture], env: { ...bunEnv, - // detect_leaks=0: ConcurrentTask/ManagedTask nodes left in a + // detect_leaks=0: ConcurrentTask nodes left in a // terminated worker's undrained concurrent queue are a known // pre-existing leak (see #32071); this test asserts no crash, not // no leaks. symbolize=0 so a pre-fix ASAN abort exits promptly From 0bcc65d9c634b4c87bbdb454a05f0a1206908312 Mon Sep 17 00:00:00 2001 From: Dylan Conway Date: Mon, 21 Sep 2026 11:05:20 +0000 Subject: [PATCH 2/4] event loop: review fixes for the typed tasks - dispatch: each server-deinit arm carries its own SAFETY comment, and the boxed test-runner and Chrome tasks take `self: Box` (clippy). - A parse result released unrun calls its native plugin's free function, which the bundle's finalizers would have called. - The server-deinit task's release comment said `finalize()` frees the server later. It does not: the task is only queued after the wrapper is finalized. - Test: a worker that exits with an unobserved HTMLRewriter transform's pull still queued frees the pipe. --- src/bundler/ParseTask.rs | 7 ++- src/runtime/dispatch.rs | 63 ++++++++++++---------- src/runtime/server/mod.rs | 3 +- src/runtime/test_runner/bun_test.rs | 3 +- src/runtime/webview/ChromeProcess.rs | 3 +- test/js/workerd/html-rewriter-leak.test.ts | 39 ++++++++++++++ 6 files changed, 86 insertions(+), 32 deletions(-) diff --git a/src/bundler/ParseTask.rs b/src/bundler/ParseTask.rs index 35ae7a7c825e..823ebfc005fe 100644 --- a/src/bundler/ParseTask.rs +++ b/src/bundler/ParseTask.rs @@ -144,7 +144,12 @@ impl bun_event_loop::Taskable for Result { unsafe fn release_unrun(this: *mut Self) { // SAFETY: fn contract — the box `run_from_thread_pool_impl` (or // `ServerComponentParseTask`) leaked. - drop(unsafe { bun_core::heap::take(this) }); + let mut result = unsafe { bun_core::heap::take(this) }; + // A native plugin's source buffer: `on_parse_task_complete` would have handed this to + // the bundle's finalizers. The source may borrow the buffer, so it goes first. + let external = core::mem::take(&mut result.external); + drop(result); + external.call(); } /// A step of the bundle. unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { diff --git a/src/runtime/dispatch.rs b/src/runtime/dispatch.rs index 12486f2be83d..a48ad9a85f65 100644 --- a/src/runtime/dispatch.rs +++ b/src/runtime/dispatch.rs @@ -335,9 +335,8 @@ pub(crate) fn run_task( } task_tag::RunTestsTask => { // SAFETY: boxed in `run_next_tick`; the arm consumes it. - crate::test_runner::bun_test::RunTestsTask::call(unsafe { - bun_core::heap::take(cast_ptr!(crate::test_runner::bun_test::RunTestsTask)) - })?; + unsafe { bun_core::heap::take(cast_ptr!(crate::test_runner::bun_test::RunTestsTask)) } + .call()?; } task_tag::ValkeyDeferredFailure => { // SAFETY: boxed at the enqueue site; the arm consumes it. @@ -392,27 +391,38 @@ pub(crate) fn run_task( task_tag::HTTPSAppClose => { crate::server::AppCloseTask::::run(cast_ptr!(crate::server::AppCloseTask)); } - // SAFETY (all four): the unique owning server pointer `schedule_deinit` queued. - task_tag::HTTPServerDeinit => unsafe { - crate::server::ServerDeinitTask::::run(cast_ptr!( - crate::server::ServerDeinitTask - )); - }, - task_tag::HTTPSServerDeinit => unsafe { - crate::server::ServerDeinitTask::::run(cast_ptr!( - crate::server::ServerDeinitTask - )); - }, - task_tag::DebugHTTPServerDeinit => unsafe { - crate::server::ServerDeinitTask::::run(cast_ptr!( - crate::server::ServerDeinitTask - )); - }, - task_tag::DebugHTTPSServerDeinit => unsafe { - crate::server::ServerDeinitTask::::run(cast_ptr!( - crate::server::ServerDeinitTask - )); - }, + task_tag::HTTPServerDeinit => { + // SAFETY: the unique owning server pointer `schedule_deinit` queued. + unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )) + }; + } + task_tag::HTTPSServerDeinit => { + // SAFETY: the unique owning server pointer `schedule_deinit` queued. + unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )) + }; + } + task_tag::DebugHTTPServerDeinit => { + // SAFETY: the unique owning server pointer `schedule_deinit` queued. + unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )) + }; + } + task_tag::DebugHTTPSServerDeinit => { + // SAFETY: the unique owning server pointer `schedule_deinit` queued. + unsafe { + crate::server::ServerDeinitTask::::run(cast_ptr!( + crate::server::ServerDeinitTask + )) + }; + } #[cfg(windows)] task_tag::CopyFileWindowsMkdirp => { // SAFETY: the live copy `on_mkdirp_complete_concurrent` posted. @@ -434,9 +444,8 @@ pub(crate) fn run_task( #[cfg(windows)] task_tag::ChromePipeEvent => { // SAFETY: boxed in `PipeEvent::post`; the arm consumes it. - crate::webview::chrome_process::QueuedEvent::deliver(unsafe { - bun_core::heap::take(cast_ptr!(crate::webview::chrome_process::QueuedEvent)) - })?; + unsafe { bun_core::heap::take(cast_ptr!(crate::webview::chrome_process::QueuedEvent)) } + .deliver()?; } task_tag::CppTask => { cast!(CppTask).run(global)?; diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index df8bb2eacf01..4bfd0cfb5edf 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -4300,8 +4300,7 @@ impl bun_event_loop::Taskable for ServerDein (false, true) => bun_event_loop::task_tag::DebugHTTPServerDeinit, (true, true) => bun_event_loop::task_tag::DebugHTTPSServerDeinit, }; - /// Not freed here: `finalize()` frees a server whose deinit was scheduled once the VM is - /// shutting down (`DEINIT_SCHEDULED`). + /// Frees nothing: a server whose deinit is still queued when its VM stops stays allocated. unsafe fn release_unrun(_: *mut Self) {} /// Enters no context. unsafe fn context(_: *const Self) -> bun_event_loop::ContextId { diff --git a/src/runtime/test_runner/bun_test.rs b/src/runtime/test_runner/bun_test.rs index fb9ba43720f9..f65e2c1f883d 100644 --- a/src/runtime/test_runner/bun_test.rs +++ b/src/runtime/test_runner/bun_test.rs @@ -1513,7 +1513,8 @@ pub struct RunTestsTask { } impl RunTestsTask { #[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")] - pub fn call(this: Box) -> JsResult<()> { + pub fn call(self: Box) -> JsResult<()> { + let this = self; // Box drops at end of scope; the Weak drops with it. let Some(strong) = this.weak.upgrade() else { return Ok(()) }; if let Err(e) = BunTest::run(&strong, &this.global_this) { diff --git a/src/runtime/webview/ChromeProcess.rs b/src/runtime/webview/ChromeProcess.rs index 8d575ea305f1..0d3f1a729325 100644 --- a/src/runtime/webview/ChromeProcess.rs +++ b/src/runtime/webview/ChromeProcess.rs @@ -957,7 +957,8 @@ impl bun_event_loop::Taskable for QueuedEvent { #[cfg(windows)] impl QueuedEvent { #[allow(clippy::boxed_local, reason = "reclaim point for the boxed task")] - pub(crate) fn deliver(queued: Box) -> bun_jsc::JsResult<()> { + pub(crate) fn deliver(self: Box) -> bun_jsc::JsResult<()> { + let queued = self; if queued.generation != GENERATION.load(Ordering::Relaxed) { scoped_log!( Chrome, diff --git a/test/js/workerd/html-rewriter-leak.test.ts b/test/js/workerd/html-rewriter-leak.test.ts index 66136af0da4d..ed0a5dbbc5cd 100644 --- a/test/js/workerd/html-rewriter-leak.test.ts +++ b/test/js/workerd/html-rewriter-leak.test.ts @@ -1,6 +1,7 @@ import { heapStats } from "bun:jsc"; import { describe, expect, test } from "bun:test"; import { bunEnv, bunExe, expectRssDeltaBelow, isASAN, isDebug, tempDir } from "harness"; +import { join } from "node:path"; // `wire_input`'s materialized-body path transfers the body's `+1` (a // `WTFStringImpl` for an all-ASCII `new Response("...")`) into an `AnyBlob` @@ -711,6 +712,44 @@ test("a direct-stream pull parked on flush(true) is released when the handler pr expect(msg).toContain("will never settle"); }); +// An unobserved transform reads one upstream chunk per event-loop turn: after the first chunk it +// queues a task for the next one, and that task holds a ref on the pipe (and so the chunks the pipe +// is holding) and a protect() on its cell. A worker that exits with the task queued must give both +// back. +test("a worker exiting with an unobserved transform's pull queued frees the pipe", async () => { + using dir = tempDir("html-rewriter-queued-pull", { + "worker.js": /* js */ ` + const chunk = new Uint8Array(1024 * 1024).fill(0x61); + const body = new ReadableStream({ + type: "direct", + pull(c) { + for (let i = 0; i < 8; i++) c.write(chunk); + c.end(); + }, + }); + // One chunk in, seven held: the next pull is queued, not run. + globalThis.keep = new HTMLRewriter().on("p", { element() {} }).transform(new Response(body)); + process.exit(0); + `, + "main.js": /* js */ ` + async function round(n) { + for (let i = 0; i < n; i++) { + const worker = new Worker(new URL("./worker.js", import.meta.url).href); + await new Promise(resolve => worker.addEventListener("close", resolve)); + } + Bun.gc(true); + return process.memoryUsage.rss(); + } + const before = await round(5); + const after = await round(15); + console.log(JSON.stringify({ deltaMiB: (after - before) / 1024 / 1024 })); + `, + }); + + // Unfixed: ~130 MiB. Fixed: allocator slack only. + await expectRssDeltaBelow([join(String(dir), "main.js")], { release: 50, debug: 60 }); +}, 20_000); + test("element.attributes iterator does not leak names/values", async () => { const code = /* js */ ` const big = Buffer.alloc(256 * 1024, "a").toString(); From 365c6307be284ed02df1cd3d0d696ad4d03f8f47 Mon Sep 17 00:00:00 2001 From: Dylan Conway Date: Mon, 21 Sep 2026 11:15:12 +0000 Subject: [PATCH 3/4] test: a failing worker fails the HTMLRewriter queued-pull leak test A worker that threw or failed to load leaked nothing, so the run exited 0 with a small RSS delta and the test passed without exercising the release path. The worker's error event now rejects. --- test/js/workerd/html-rewriter-leak.test.ts | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/test/js/workerd/html-rewriter-leak.test.ts b/test/js/workerd/html-rewriter-leak.test.ts index ed0a5dbbc5cd..4166d3b8fb5f 100644 --- a/test/js/workerd/html-rewriter-leak.test.ts +++ b/test/js/workerd/html-rewriter-leak.test.ts @@ -735,7 +735,11 @@ test("a worker exiting with an unobserved transform's pull queued frees the pipe async function round(n) { for (let i = 0; i < n; i++) { const worker = new Worker(new URL("./worker.js", import.meta.url).href); - await new Promise(resolve => worker.addEventListener("close", resolve)); + // A worker that fails would leak nothing and pass: fail the run instead. + await new Promise((resolve, reject) => { + worker.addEventListener("close", resolve); + worker.addEventListener("error", event => reject(new Error(event.message))); + }); } Bun.gc(true); return process.memoryUsage.rss(); From 305fb677b385d54780256656fd871b5b9895aca2 Mon Sep 17 00:00:00 2001 From: Dylan Conway Date: Mon, 21 Sep 2026 12:52:47 +0000 Subject: [PATCH 4/4] test: the HTMLRewriter queued-pull leak test fits the default timeout Eight workers holding 28 MiB each instead of twenty holding 7 MiB: about 1.7 s on a debug build, ~145 MiB of growth without the fix. The per-test timeout is gone. --- test/js/workerd/html-rewriter-leak.test.ts | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/test/js/workerd/html-rewriter-leak.test.ts b/test/js/workerd/html-rewriter-leak.test.ts index 4166d3b8fb5f..7c6ff001fc86 100644 --- a/test/js/workerd/html-rewriter-leak.test.ts +++ b/test/js/workerd/html-rewriter-leak.test.ts @@ -719,7 +719,7 @@ test("a direct-stream pull parked on flush(true) is released when the handler pr test("a worker exiting with an unobserved transform's pull queued frees the pipe", async () => { using dir = tempDir("html-rewriter-queued-pull", { "worker.js": /* js */ ` - const chunk = new Uint8Array(1024 * 1024).fill(0x61); + const chunk = new Uint8Array(4 * 1024 * 1024).fill(0x61); const body = new ReadableStream({ type: "direct", pull(c) { @@ -744,15 +744,15 @@ test("a worker exiting with an unobserved transform's pull queued frees the pipe Bun.gc(true); return process.memoryUsage.rss(); } - const before = await round(5); - const after = await round(15); + const before = await round(2); + const after = await round(6); console.log(JSON.stringify({ deltaMiB: (after - before) / 1024 / 1024 })); `, }); - // Unfixed: ~130 MiB. Fixed: allocator slack only. + // Unfixed: ~145 MiB. Fixed: allocator slack only. await expectRssDeltaBelow([join(String(dir), "main.js")], { release: 50, debug: 60 }); -}, 20_000); +}); test("element.attributes iterator does not leak names/values", async () => { const code = /* js */ `