diff --git a/docs/ROOT-B-SHUTDOWN-FENCE.md b/docs/ROOT-B-SHUTDOWN-FENCE.md new file mode 100644 index 000000000000..0f4de9099f7f --- /dev/null +++ b/docs/ROOT-B-SHUTDOWN-FENCE.md @@ -0,0 +1,148 @@ +# Root B: worker teardown vs. cross-thread completions + +## The failure class + +`WebWorker::shutdown` (`src/jsc/web_worker.rs:1216`) tears down a worker's +`VirtualMachine` in five ordered steps. Step 2 sets `is_shutting_down`, drains +timers, force-closes sockets and the c-ares channel, calls +`ScriptExecutionContext::markTerminating()` (so C++ `postTaskTo` posters are +fenced), and drains the concurrent queue via +`release_queued_tasks_for_shutdown()`. Step 5 `std::alloc::dealloc`'s the raw +`VirtualMachine` box (≈`web_worker.rs:1383`) and frees the uws loop. + +Nothing in that sequence cancels or awaits work already handed to a +process-global thread (WorkPool, the HTTP thread, the bundle thread). Those +jobs complete later and post back through a pointer captured at schedule time: +a `BackRef`, a `&'static VirtualMachine`, or a `*const +JSGlobalObject`. Every one of those pointers is into the freed VM box or the +freed JSC heap. + +All reproduced cross-thread UAFs in this class converge on +**`EventLoop::enqueue_task_concurrent`** (`src/jsc/event_loop.rs:997`). Putting +a check _inside_ the funnel does not help: `&self` there is already a pointer +into the freed box. Several callers also "guard" with an off-thread +`vm.is_shutting_down()` read; that read is itself a UAF face (the flag lives in +the freed box). + +## Shutdown map (reference) + +| step | action | fences | +| ---- | -------------------------------------------------------------- | ----------------------------- | +| 1 | `self.vm = null` under `vm_lock` | parent-thread readers | +| 2a | `is_shutting_down = true`, `on_exit`, drain timers/sockets/DNS | on-thread re-entry | +| 2b | `ScriptExecutionContext::markTerminating()` | C++ `postTaskTo` posters | +| 2c | `Bun__JSCTaskScheduler__markShuttingDown()` | `Atomics.notify` posters | +| 2d | `release_queued_tasks_for_shutdown()` | tasks already queued | +| 3 | `WebWorker__teardownJSCVM` | GC finalizers, JSC heap freed | +| 4 | `WebWorker__dispatchExit` | parent releases its ref | +| 5 | `vm.destroy()`; `dealloc(vm_ptr)`; free uws loop | VM box freed | + +Rust-side cross-thread posters were not serialized with any of 2b/2c/2d. + +## The fence: enqueue by identifier + +Off-thread jobs now carry the worker's `ScriptExecutionContextIdentifier` (a +`u32`) instead of a `BackRef` / `&VirtualMachine`, and post through + +```rust +ScriptExecutionContextIdentifier::post_concurrent_task(id, task) -> bool +``` + +backed by the same locked-registry + `isTerminating()` gate that +`ScriptExecutionContext::postTaskTo` already uses: + +```cpp +extern "C" bool ScriptExecutionContext__postConcurrentTask(Identifier id, void* task) { + Locker locker { allScriptExecutionContextsMapLock }; + auto* ctx = allScriptExecutionContextsMap().get(id); + if (!ctx || ctx->isTerminating()) return false; + Bun__EventLoop__enqueueConcurrentTask(ctx->globalObject(), task); + return true; +} +``` + +`markTerminating()` (shutdown step 2b) takes the same lock to set the flag, so +every poster serializes into exactly one of two cases: + +1. The poster's whole critical section ran before `markTerminating()`: the task + is in the concurrent queue, and step 2d's drain observes and reclaims it. +2. `markTerminating()` ran first: the poster sees `isTerminating()` and returns + `false` without touching the VM. The caller owns the task and runs its + abandon path. + +A `u32` identifier cannot dangle. The funnel's `bool` return replaces every +stale off-thread `is_shutting_down()` read. + +Companion helpers on the same lock: + +- `ScriptExecutionContextIdentifier::is_alive()` — "should I even start?" check + for work bodies that write into JSC-heap buffers (e.g. `Scrypt`'s output + `ArrayBuffer`). A best-effort fast drop; the authoritative gate is + `post_concurrent_task`. +- `ScriptExecutionContextIdentifier::unref_event_loop_concurrently()` — the + `concurrent_ref` decrement that `ConcurrentCppTask` (WebCrypto) needs after + its body ran on the pool thread, without dereferencing the VM. + +## Abandon path + +On `post_concurrent_task` → `false` the target VM and its JSC heap are gone (or +about to be). The abandon path: + +- **must not** touch `Strong`/`Weak`/`JSPromiseStrong`/`JSGlobalObject`/ + `VirtualMachine`/`KeepAlive::unref` — the HandleSet is freed, the loop is + freed; +- **may** free any pure-Rust heap it owns (body buffers, `Vec`s, `Box<[u8]>`); +- **must** free a freshly heap-allocated `ConcurrentTask` node (ownership was + not transferred); +- **may** leak the job box when it holds JSC handles. Bounded: one per + terminated worker per in-flight op. + +## Coverage + +Three generic helpers carry most of the surface: + +| helper | users | +| -------------------------- | ----------------------------------------------------------- | +| `WorkTask` | `ReadFile`, `WriteFile`, `GetAddrInfoRequest` | +| `ConcurrentPromiseTask` | `CopyFile`, `TransformTask`, `WalkTask`, `PipelineTask` | +| `AnyTaskJob` | `Pbkdf2Ctx`, `CryptoJob`, `ZstdCtx`, `SecretsCtx` | +| `ConcurrentCppTask` | WebCrypto (`PhonyWorkQueue::dispatch`) | + +Direct callers converted alongside: `FetchTasklet`, `PasswordJob`, +`CompressionStream` (zlib/brotli/zstd), `AsyncFSTask` / `NewAsyncCpTask` / +`AsyncReaddirRecursiveTask`, `S3HttpSimpleTask` / `S3HttpDownloadStreamingTask`, +`Archive::AsyncTask`, `JSBundleCompletionTask`, `TranspilerJob`. + +Same enqueue shape, deferred to follow-up (not in the verify harness): +`napi_async_work` / `ThreadSafeFunction`, `fs.watch` / `fs.watchFile` +(PathWatcherManager reader thread), the shell WorkPool tasks, `AsyncModule` +package-manager wake, and the Windows-only `WriteFileWindows` / `CopyFileWindows` +`AsyncMkdirp` completion (the POSIX `WriteFile`/`CopyFile` paths above go +through the converted `WorkTask`/`ConcurrentPromiseTask`). + +Explicitly **not** enqueue-shaped and left for their own fixes: nested-worker +child-init reading a freed parent VM, `node:quic` finalizer ordering, +`RedisClient::finalize` free-then-read, `Bun.SQL` handle crashes during +terminating-VM JS execution, the `serve.listen`/JS-re-entry assert zone. + +## Reserve alternative (not taken) + +Refcount-deferred VM dealloc + a closed-flag at the funnel: each off-thread job +takes an `Arc` clone of a per-VM gate and brackets its enqueue with a read +lock; `shutdown` takes the write lock before dealloc. Same sweep, but: + +- `shutdown` then _waits_ on in-flight pool jobs (a slow argon2/RSA can stall + terminate for seconds); +- more atomics on the hot enqueue path; +- the gate must be `Arc`'d so it outlives the VM box anyway. + +The identifier route reuses an existing lock, adds no wait, and matches what +the C++ side already does. + +## Post-fix gate + +`repro/rootB-verify/verify.mjs`: one worker per iteration arms one in-flight op +of every cross-thread source above (self-contained, loopback-only, public API), +the parent `terminate()`s mid-flight, ×100. PASS = rc 0, `ROOT-B VERIFY: PASS`, +zero ASan/assert/panic. Baseline on current canary: SIGSEGV on teardown 1 +(release), heap-use-after-free on teardown 1-3 (debug+ASAN). diff --git a/repro/rootB-verify/verify.mjs b/repro/rootB-verify/verify.mjs new file mode 100644 index 000000000000..354683a7fbc1 --- /dev/null +++ b/repro/rootB-verify/verify.mjs @@ -0,0 +1,45 @@ +// Root-B post-fix gate: one worker per iteration arms one in-flight op of +// every cross-thread completion source (public API, loopback-only), parent +// terminates mid-flight, x ITERATIONS. PASS = rc 0 + the PASS line + zero +// ASan/assert/panic. On stock canary this SIGSEGVs on the first teardown. +// +// Run under the debug+ASAN build so the fail-before is deterministic: +// bun bd repro/rootB-verify/verify.mjs + +import { Worker } from "node:worker_threads"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +const ITERATIONS = Number(process.env.ROOTB_ITER ?? 100); +const body = new URL("./worker-body.mjs", import.meta.url); +const scratch = fs.mkdtempSync(path.join(os.tmpdir(), "rootB-verify-")); +process.env.ROOTB_SCRATCH = scratch; +process.on("exit", () => fs.rmSync(scratch, { recursive: true, force: true })); + +for (let i = 0; i < ITERATIONS; i++) { + const w = new Worker(body); + // Wait for the worker to finish arming (or give up after a bound). + const armed = await new Promise((resolve) => { + const t = setTimeout(() => resolve("timeout"), 5000); + w.once("message", (m) => { + clearTimeout(t); + resolve(m); + }); + w.once("error", (e) => { + clearTimeout(t); + resolve(e); + }); + }); + if (armed instanceof Error) { + console.error(`iter ${i}: worker error before terminate:`, armed); + process.exit(1); + } + // Small jitter so terminate lands at varying points in the in-flight work. + const jitter = (i * 2654435761 >>> 0) % 5; + if (jitter) await new Promise((r) => setTimeout(r, jitter)); + await w.terminate(); + if (i % 10 === 9) console.log(`… ${i + 1}/${ITERATIONS} teardowns clean`); +} + +console.log(`ROOT-B VERIFY: PASS (${ITERATIONS} teardowns)`); diff --git a/repro/rootB-verify/worker-body.mjs b/repro/rootB-verify/worker-body.mjs new file mode 100644 index 000000000000..89fd53236c34 --- /dev/null +++ b/repro/rootB-verify/worker-body.mjs @@ -0,0 +1,107 @@ +// One worker instance: arm one in-flight op of every cross-thread completion +// source, then sit. Parent terminates us mid-flight. Every op is self-contained +// and loopback-only (no public internet). Catch-and-ignore everywhere: the +// point is to have work airborne on a process-global thread when the VM dies, +// not to observe the result. + +import { parentPort, threadId } from "node:worker_threads"; +import fs from "node:fs"; +import fsp from "node:fs/promises"; +import zlib from "node:zlib"; +import crypto from "node:crypto"; +import dns from "node:dns/promises"; +import os from "node:os"; +import path from "node:path"; + +const sink = () => {}; +const swallow = (p) => Promise.resolve(p).then(sink, sink); + +// One scratch root for all iterations (verify.mjs removes it on exit). +const root = process.env.ROOTB_SCRATCH ?? path.join(os.tmpdir(), "rootB-verify"); +const tmp = path.join(root, String(threadId)); +fs.mkdirSync(tmp, { recursive: true }); +const tmpFile = path.join(tmp, "a.txt"); +fs.writeFileSync(tmpFile, Buffer.alloc(1 << 16, "x").toString()); + +// fetch / HTMLRewriter / TLS (HTTP thread -> FetchTasklet) +{ + const srv = Bun.serve({ + port: 0, + fetch: () => new Response("x".repeat(1 << 14)), + }); + swallow(fetch(`http://127.0.0.1:${srv.port}/`).then((r) => r.text())); + swallow( + fetch(`http://127.0.0.1:${srv.port}/`).then((r) => + new HTMLRewriter().on("*", { text() {} }).transform(r).text(), + ), + ); +} + +// Bun.file / Bun.write (WorkPool -> WorkTask/) +swallow(Bun.write(path.join(tmp, "b.txt"), "y".repeat(1 << 16))); +swallow(Bun.file(tmpFile).text()); + +// node:fs promises (WorkPool -> AsyncFSTask) +swallow(fsp.readFile(tmpFile)); +swallow(fsp.stat(tmpFile)); +swallow(fsp.readdir(tmp, { recursive: true })); + +// pbkdf2 / scrypt / generateKeyPair (WorkPool -> AnyTaskJob) +crypto.pbkdf2("p", "s", 100000, 64, "sha512", sink); +crypto.scrypt("p", "saltsalt", 64, sink); +crypto.generateKeyPair("rsa", { modulusLength: 2048 }, sink); + +// Bun.password (WorkPool -> PasswordJob) +swallow(Bun.password.hash("hunter2", { algorithm: "bcrypt", cost: 8 })); + +// zlib (WorkPool -> CompressionStream async_job_run) +zlib.deflate(Buffer.alloc(1 << 18), sink); +zlib.gzip(Buffer.alloc(1 << 18), sink); + +// S3 (HTTP thread -> S3HttpSimpleTask); loopback endpoint that never answers. +{ + const srv = Bun.serve({ port: 0, fetch: () => new Promise(sink) }); + const s3 = new Bun.S3Client({ + accessKeyId: "x", + secretAccessKey: "y", + endpoint: `http://127.0.0.1:${srv.port}`, + bucket: "b", + }); + swallow(s3.file("k").text()); +} + +// Bun.build (BundleThread -> JSBundleCompletionTask) +{ + const entry = path.join(tmp, "entry.ts"); + fs.writeFileSync(entry, `export const x: number = 1;\n`); + swallow(Bun.build({ entrypoints: [entry], target: "bun" })); +} + +// Transpiler.transform (WorkPool -> ConcurrentPromiseTask) +swallow(new Bun.Transpiler({ loader: "tsx" }).transform("const x: number = 1;")); + +// dns.lookup (c-ares; clean-by-construction via close_dns_for_terminate) +swallow(dns.lookup("localhost")); + +// WebCrypto (WorkPool -> ConcurrentCppTask) +swallow( + crypto.subtle.digest("SHA-256", new Uint8Array(1 << 16)), +); +swallow( + crypto.subtle.generateKey({ name: "AES-GCM", length: 256 }, true, [ + "encrypt", + "decrypt", + ]), +); + +// Glob (WorkPool -> ConcurrentPromiseTask) +swallow(Array.fromAsync(new Bun.Glob("**/*").scan(tmp))); + +// zstd (WorkPool -> AnyTaskJob) +swallow(Bun.zstdCompress(Buffer.alloc(1 << 16))); + +// Signal parent that everything is airborne. +parentPort.postMessage("armed"); + +// Keep the loop alive so terminate() lands mid-flight. +setInterval(sink, 1 << 30); diff --git a/src/event_loop/ConcurrentTask.rs b/src/event_loop/ConcurrentTask.rs index f88145557f66..94a32b3458cf 100644 --- a/src/event_loop/ConcurrentTask.rs +++ b/src/event_loop/ConcurrentTask.rs @@ -309,6 +309,21 @@ impl ConcurrentTask { Self::create(ManagedTask::ManagedTask::new(ptr, callback)) } + /// Reclaim a node produced by [`Self::from_callback`] that was never + /// enqueued. Frees both the outer `ConcurrentTask` and the inner + /// `ManagedTask` box; does not call the callback. + /// + /// # Safety + /// `node` must be a `from_callback`-allocated node whose ownership was not + /// transferred to a queue. + pub unsafe fn destroy_from_callback(node: core::ptr::NonNull) { + // SAFETY: caller contract — `from_callback` wraps a heap `ManagedTask`. + let outer = unsafe { bun_core::heap::take(node.as_ptr()) }; + debug_assert!(outer.task.tag == task_tag::ManagedTask); + // SAFETY: `ManagedTask::new` produced the inner box via `heap::into_raw`. + drop(unsafe { bun_core::heap::take(outer.task.ptr.cast::()) }); + } + pub fn from( &mut self, of: *mut T, diff --git a/src/jsc/ConcurrentPromiseTask.rs b/src/jsc/ConcurrentPromiseTask.rs index 62aba68d025b..5ee3af114db9 100644 --- a/src/jsc/ConcurrentPromiseTask.rs +++ b/src/jsc/ConcurrentPromiseTask.rs @@ -2,11 +2,9 @@ use bun_event_loop::ConcurrentTask::{AutoDeinit, ConcurrentTask, TaskTag, Taskab use bun_io::{self as Async, KeepAlive}; use bun_threading::{IntrusiveWorkTask as _, WorkPoolTask, work_pool::WorkPool}; -use crate::event_loop::EventLoop; +use crate::js_global_object::ScriptExecutionContextIdentifier; use crate::js_promise::{JSPromise, Strong as JSPromiseStrong}; -use crate::virtual_machine::VirtualMachine; use crate::{JSGlobalObject, JsTerminated}; -use bun_ptr::BackRef; /// The `Context` type parameter for [`ConcurrentPromiseTask`] must implement this trait: /// - `run(&mut self)` — performs the work on the thread pool @@ -31,9 +29,8 @@ pub struct ConcurrentPromiseTask<'a, Context: ConcurrentPromiseTaskContext> { // Owned here so dropping the task frees the context. pub ctx: Box, pub task: WorkPoolTask, - /// BACKREF — captured from the JS-thread VM at create time; the VM (and its - /// `EventLoop`) outlives every task scheduled on it. - pub event_loop: BackRef, + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + pub context_id: ScriptExecutionContextIdentifier, pub promise: JSPromiseStrong, pub global_this: &'a JSGlobalObject, pub concurrent_task: ConcurrentTask, @@ -57,11 +54,8 @@ impl Taskable for ConcurrentPromiseTask<' impl<'a, Context: ConcurrentPromiseTaskContext> ConcurrentPromiseTask<'a, Context> { pub fn create_on_js_thread(global_this: &'a JSGlobalObject, value: Box) -> Box { - // `VirtualMachine::get()` returns the JS-thread singleton; the VM and - // its `EventLoop` outlive every task scheduled on it. - let event_loop = BackRef::new(VirtualMachine::get().as_mut().event_loop_shared()); let mut this = Box::new(Self { - event_loop, + context_id: global_this.script_execution_context_identifier(), ctx: value, task: WorkPoolTask { node: Default::default(), @@ -108,15 +102,14 @@ impl<'a, Context: ConcurrentPromiseTaskContext> ConcurrentPromiseTask<'a, Contex // `this` while holding `&mut *this` is sound because `from` only stores // the pointer (does not dereference it). let this_ref = unsafe { &mut *this }; - let event_loop = this_ref.event_loop; + let context_id = this_ref.context_id; let task = core::ptr::NonNull::from( this_ref .concurrent_task .from(this, AutoDeinit::ManualDeinit), ); - // `task` is the live `concurrent_task` field of the heap-allocated - // job; the queue takes ownership of its intrusive `next` link. - event_loop.enqueue_task_concurrent(task); + // Abandon: JSC handles cannot drop off-thread, leak the box (task is intrusive). + let _ = context_id.post_concurrent_task(task); } /// Frees the heap allocation backing this task. diff --git a/src/jsc/CppTask.rs b/src/jsc/CppTask.rs index 130402469c25..e3c9e508a278 100644 --- a/src/jsc/CppTask.rs +++ b/src/jsc/CppTask.rs @@ -1,5 +1,6 @@ use core::ptr::NonNull; +use crate::js_global_object::ScriptExecutionContextIdentifier; use crate::{JSGlobalObject, JsResult, VirtualMachineRef as VirtualMachine}; use bun_event_loop::{TaskTag, Taskable, task_tag}; use bun_threading::work_pool::{Task as WorkPoolTask, WorkPool}; @@ -10,6 +11,7 @@ unsafe extern "C" { safe fn Bun__EventLoopTaskNoContext__createdInBunVm( task: &EventLoopTaskNoContext, ) -> *mut VirtualMachine; + safe fn Bun__EventLoopTaskNoContext__contextId(task: &EventLoopTaskNoContext) -> u32; } bun_opaque::opaque_ffi! { @@ -55,6 +57,11 @@ impl EventLoopTaskNoContext { pub fn get_vm(&self) -> Option> { NonNull::new(Bun__EventLoopTaskNoContext__createdInBunVm(self)).map(bun_ptr::BackRef::from) } + + /// See [`ScriptExecutionContextIdentifier::unref_event_loop_concurrently`]. + pub fn context_id(&self) -> ScriptExecutionContextIdentifier { + ScriptExecutionContextIdentifier(Bun__EventLoopTaskNoContext__contextId(self)) + } } /// A task created from C++ code that runs inside the workpool, usually via ScriptExecutionContext. @@ -73,14 +80,14 @@ impl ConcurrentCppTask { let cpp_task = self.cpp_task; // `EventLoopTaskNoContext` is an `opaque_ffi!` ZST handle; `opaque_ref` // is the centralised non-null deref proof. Valid until `run` consumes it. - let maybe_vm = EventLoopTaskNoContext::opaque_ref(cpp_task).get_vm(); + let context_id = EventLoopTaskNoContext::opaque_ref(cpp_task).context_id(); drop(self); // SAFETY: `cpp_task` is the valid C++ handle stored by `ConcurrentCppTask__createAndRun`; // `opaque_ref` above proved it non-null and it has not yet been freed — `run` consumes it here. unsafe { EventLoopTaskNoContext::run(cpp_task) }; - if let Some(vm) = maybe_vm { - vm.event_loop_shared().unref_concurrently(); - } + // Task body posts its own result via `postTaskTo`; route the trailing + // `concurrent_ref` decrement through the same lock. + context_id.unref_event_loop_concurrently(); } } diff --git a/src/jsc/JSGlobalObject.rs b/src/jsc/JSGlobalObject.rs index 896ca52e75e3..79810bccaf69 100644 --- a/src/jsc/JSGlobalObject.rs +++ b/src/jsc/JSGlobalObject.rs @@ -1702,9 +1702,61 @@ impl ScriptExecutionContextIdentifier { pub fn valid(self) -> bool { self.global_object().is_some() } + + /// `true` while the context exists and has not been marked terminating. + /// Serializes with `markTerminating()` on the contexts-map lock. + #[inline] + pub fn is_alive(self) -> bool { + ScriptExecutionContext__isAlive(self.0) + } + + /// Enqueue a heap-allocated [`ConcurrentTaskItem`] onto this context's + /// event loop from any thread. Runs under the contexts-map lock (same lock + /// as `markTerminating()`), so the target VM is dereferenced only while + /// known live. + /// + /// Returns `true` if enqueued (ownership transferred). Returns `false` if + /// the context is gone or terminating; caller retains ownership of `task` + /// and must free it without touching the target VM/JSC heap. + /// + /// [`ConcurrentTaskItem`]: crate::event_loop::ConcurrentTaskItem + #[inline] + pub fn post_concurrent_task( + self, + task: core::ptr::NonNull, + ) -> bool { + ScriptExecutionContext__postConcurrentTask(self.0, task.as_ptr().cast::()) + } + + /// Decrement this context's event-loop `concurrent_ref` from off-thread + /// under the contexts-map lock; no-op when the context is gone or + /// terminating. + #[inline] + pub fn unref_event_loop_concurrently(self) { + ScriptExecutionContext__unrefEventLoopConcurrently(self.0); + } } unsafe extern "C" { // safe: by-value `u32` in, raw nullable pointer out (caller checks before deref). safe fn ScriptExecutionContextIdentifier__getGlobalObject(id: u32) -> *mut JSGlobalObject; + // safe: by-value `u32` in, opaque pointer C++ only forwards back to Rust + // under the contexts-map lock; bool out. + safe fn ScriptExecutionContext__postConcurrentTask(id: u32, task: *mut c_void) -> bool; + // safe: by-value `u32` in; bool out. + safe fn ScriptExecutionContext__isAlive(id: u32) -> bool; + // safe: by-value `u32` in; void out. + safe fn ScriptExecutionContext__unrefEventLoopConcurrently(id: u32); + // safe: `&JSGlobalObject` is a live opaque handle; returns the context's id. + safe fn ScriptExecutionContextIdentifier__forGlobalObject(global: &JSGlobalObject) -> u32; +} + +impl JSGlobalObject { + /// The stable id for this global's `ScriptExecutionContext`. Capture on the + /// JS thread; post completions back via + /// [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + #[inline] + pub fn script_execution_context_identifier(&self) -> ScriptExecutionContextIdentifier { + ScriptExecutionContextIdentifier(ScriptExecutionContextIdentifier__forGlobalObject(self)) + } } diff --git a/src/jsc/RuntimeTranspilerStore.rs b/src/jsc/RuntimeTranspilerStore.rs index acd5bb4d9438..b869d722cdfc 100644 --- a/src/jsc/RuntimeTranspilerStore.rs +++ b/src/jsc/RuntimeTranspilerStore.rs @@ -34,6 +34,7 @@ use bun_watcher::Watcher; use crate::async_module::AsyncModule; use crate::event_loop::{ConcurrentTask, EventLoop}; use crate::hot_reloader::ImportWatcher; +use crate::js_global_object::ScriptExecutionContextIdentifier; use crate::resolved_source::OwnedResolvedSource; use crate::resolved_source_tag::ResolvedSourceTag; use crate::runtime_transpiler_cache::{ @@ -319,6 +320,7 @@ impl RuntimeTranspilerStore { non_threadsafe_input_specifier: OwnedString::new(input_specifier), path: owned_path, global_this: BackRef::new(global_object), + context_id: global_object.script_execution_context_identifier(), non_threadsafe_referrer: OwnedString::new(referrer), vm, log: bun_ast::Log::init(), @@ -378,6 +380,7 @@ pub struct TranspilerJob { // store and outlives every job). pub vm: *mut VirtualMachine, pub global_this: BackRef, + pub context_id: ScriptExecutionContextIdentifier, pub fetcher: Fetcher, pub poll_ref: KeepAlive, pub generation_number: u32, @@ -485,17 +488,39 @@ impl TranspilerJob { } pub(crate) fn dispatch_to_main_thread(&mut self) { + let context_id = self.context_id; + // VM owns both `transpiler_store.queue` and the event loop; touch neither after teardown. + if !context_id.is_alive() { + // Reclaim pure-Rust heap; leak JSC-heap handles. + let old_path = core::mem::take(&mut self.path); + if !old_path.text.is_empty() { + // SAFETY: `text` is the `Box<[u8]>` produced by + // `heap::into_raw` in `transpile()`; this is the unique owner. + drop(unsafe { + bun_core::heap::take(ptr::from_ref::<[u8]>(old_path.text).cast_mut()) + }); + } + self.log = bun_ast::Log::init(); + self.parse_error = None; + return; + } let vm = self.vm; - // SAFETY: vm outlives the job (BACKREF — VM owns the store). + // SAFETY: `is_alive()` is best-effort (lock released before this deref); + // `transpiler_store` is a VM field, so this and the push race the narrow + // is_alive→dealloc window, same as every `(*vm).*` read in `run()`. The + // authoritative fence is `post_concurrent_task` below. let transpiler_store: *mut RuntimeTranspilerStore = unsafe { ptr::addr_of_mut!((*vm).transpiler_store) }; let job = NonNull::from(&mut *self); // SAFETY: queue is concurrent-safe (UnboundedQueue uses atomics). unsafe { (*transpiler_store).queue.push(job) }; // Another thread may free `self` at any time after .push, so we cannot use it any more. - // SAFETY: vm outlives the job; event_loop() returns the live self-pointer. - unsafe { &*(*vm).event_loop() } - .enqueue_task_concurrent(ConcurrentTask::create_from(transpiler_store)); + let task = ConcurrentTask::create_from(transpiler_store); + if !context_id.post_concurrent_task(task) { + // Raced with teardown; job already in the VM-owned queue, reclaim only the node. + // SAFETY: ownership not transferred; `task` is the `create_from` heap node above. + drop(unsafe { bun_core::heap::take(task.as_ptr()) }); + } } pub(crate) fn run_from_js_thread(&mut self) -> JsResult<()> { diff --git a/src/jsc/WorkTask.rs b/src/jsc/WorkTask.rs index aa90ad9373a0..fd90cc09878c 100644 --- a/src/jsc/WorkTask.rs +++ b/src/jsc/WorkTask.rs @@ -4,7 +4,7 @@ use bun_threading::{IntrusiveWorkTask as _, WorkPoolTask, work_pool::WorkPool}; use crate::JSGlobalObject; use crate::debugger::AsyncTaskTracker; -use crate::event_loop::EventLoop; +use crate::js_global_object::ScriptExecutionContextIdentifier; use bun_ptr::BackRef; /// A generic task that runs work on a thread pool and executes a callback on the main JavaScript thread. @@ -34,9 +34,8 @@ pub trait WorkTaskContext: Sized { pub struct WorkTask { pub ctx: *mut Context, pub task: WorkPoolTask, - /// BACKREF — captured from the JS-thread VM at create time; the VM (and its - /// `EventLoop`) outlives every task scheduled on it. - pub event_loop: BackRef, + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + pub context_id: ScriptExecutionContextIdentifier, // allocator field dropped — global mimalloc (see PORTING.md §Allocators) pub global_this: BackRef, pub concurrent_task: ConcurrentTask, @@ -61,9 +60,8 @@ impl Taskable for WorkTask { impl WorkTask { pub fn create_on_js_thread(global_this: &JSGlobalObject, value: *mut Context) -> *mut Self { let vm = global_this.bun_vm().as_mut(); - let event_loop = BackRef::new(vm.event_loop_shared()); let mut this = Box::new(Self { - event_loop, + context_id: global_this.script_execution_context_identifier(), ctx: value, global_this: BackRef::new(global_this), task: WorkPoolTask { @@ -129,14 +127,13 @@ impl WorkTask { // re-initializes it in place and returns the same address. Passing // `this_ptr` while holding `&mut *this` is sound because `from` only // stores the pointer (does not dereference it). - let event_loop = this.event_loop; + let context_id = this.context_id; let this_ptr: *mut Self = this; let task = core::ptr::NonNull::from( this.concurrent_task .from(this_ptr, AutoDeinit::ManualDeinit), ); - // `task` is the inline `concurrent_task` field of the live - // heap-allocated `*this`; `event_loop` is the JS-thread loop stored at init. - event_loop.enqueue_task_concurrent(task); + // Abandon: JSC handles cannot drop off-thread, leak the box (task is intrusive). + let _ = context_id.post_concurrent_task(task); } } diff --git a/src/jsc/any_task_job.rs b/src/jsc/any_task_job.rs index 870232d0fae8..717e5dc005f5 100644 --- a/src/jsc/any_task_job.rs +++ b/src/jsc/any_task_job.rs @@ -15,6 +15,7 @@ use bun_io::KeepAlive; use bun_threading::work_pool::{IntrusiveWorkTask as _, Task as WorkPoolTask, WorkPool}; use crate::event_loop::ConcurrentTask; +use crate::js_global_object::ScriptExecutionContextIdentifier; use crate::{JSGlobalObject, JsResult, VirtualMachineRef as VirtualMachine}; /// Per-job payload trait. Implementors own the off-thread work body and the @@ -49,6 +50,10 @@ pub trait AnyTaskJobCtx: Sized { /// e.g. a `JSPromiseStrong` field after scheduling. pub struct AnyTaskJob { vm: bun_ptr::BackRef, + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + context_id: ScriptExecutionContextIdentifier, + /// Captured on the JS thread so the pool thread never derefs `vm`. + global: *mut JSGlobalObject, task: WorkPoolTask, any_task: AnyTask, poll: KeepAlive, @@ -75,6 +80,8 @@ impl AnyTaskJob { let vm = bun_ptr::BackRef::new(global.bun_vm()); let job = bun_core::heap::into_raw(Box::new(Self { vm, + context_id: global.script_execution_context_identifier(), + global: core::ptr::from_ref(global).cast_mut(), task: WorkPoolTask { node: Default::default(), callback: Self::run_task, @@ -137,14 +144,21 @@ impl AnyTaskJob { fn run_task(task: *mut WorkPoolTask) { // SAFETY: only reachable via the `WorkPoolTask::callback` slot wired // in `create`; `task` points to `Self.task` and the job is live until - // `run_from_js` reclaims it. + // `run_from_js` reclaims it (or is leaked on abandon below). let job = unsafe { &mut *Self::from_task_ptr(task) }; - let vm = job.vm; - job.ctx.run(vm.global); - // `ConcurrentTask::create` heap-allocates a fresh task; the queue takes - // ownership of it. - vm.event_loop_shared() - .enqueue_task_concurrent(ConcurrentTask::create(job.any_task.task())); + let context_id = job.context_id; + // Skip the work body on shutdown: some ctxs write into JSC-heap-backed + // `ArrayBuffer`s. Fast-path only; `post_concurrent_task` below is the gate. + if context_id.is_alive() { + job.ctx.run(job.global); + } + let task = ConcurrentTask::create(job.any_task.task()); + if context_id.post_concurrent_task(task) { + return; + } + // Abandon: JSC handles cannot drop off-thread, leak the job box. + // SAFETY: ownership not transferred; `task` was `ConcurrentTask::create`-allocated above. + drop(unsafe { bun_core::heap::take(task.as_ptr()) }); } /// `AnyTask` callback — runs ON the JS thread. Reclaims the heap diff --git a/src/jsc/bindings/EventLoopTaskNoContext.cpp b/src/jsc/bindings/EventLoopTaskNoContext.cpp index 6f1f3d367ad9..028765551391 100644 --- a/src/jsc/bindings/EventLoopTaskNoContext.cpp +++ b/src/jsc/bindings/EventLoopTaskNoContext.cpp @@ -12,4 +12,9 @@ extern "C" void* Bun__EventLoopTaskNoContext__createdInBunVm(const EventLoopTask return task->createdInBunVm(); } +extern "C" uint32_t Bun__EventLoopTaskNoContext__contextId(const EventLoopTaskNoContext* task) +{ + return task->contextId(); +} + } // namespace Bun diff --git a/src/jsc/bindings/EventLoopTaskNoContext.h b/src/jsc/bindings/EventLoopTaskNoContext.h index fede33f2603c..c7df6a092ed5 100644 --- a/src/jsc/bindings/EventLoopTaskNoContext.h +++ b/src/jsc/bindings/EventLoopTaskNoContext.h @@ -1,6 +1,7 @@ #pragma once #include "ZigGlobalObject.h" +#include "ScriptExecutionContext.h" #include "root.h" namespace Bun { @@ -12,6 +13,7 @@ class EventLoopTaskNoContext { public: EventLoopTaskNoContext(JSC::JSGlobalObject* globalObject, Function&& task) : m_createdInBunVm(defaultGlobalObject(globalObject)->bunVM()) + , m_contextId(defaultGlobalObject(globalObject)->scriptExecutionContext()->identifier()) , m_task(WTF::move(task)) { } @@ -23,13 +25,16 @@ class EventLoopTaskNoContext { } void* createdInBunVm() const { return m_createdInBunVm; } + WebCore::ScriptExecutionContextIdentifier contextId() const { return m_contextId; } private: void* m_createdInBunVm; + WebCore::ScriptExecutionContextIdentifier m_contextId; Function m_task; }; extern "C" void Bun__EventLoopTaskNoContext__performTask(EventLoopTaskNoContext* task); extern "C" void* Bun__EventLoopTaskNoContext__createdInBunVm(const EventLoopTaskNoContext* task); +extern "C" uint32_t Bun__EventLoopTaskNoContext__contextId(const EventLoopTaskNoContext* task); } // namespace Bun diff --git a/src/jsc/bindings/ScriptExecutionContext.cpp b/src/jsc/bindings/ScriptExecutionContext.cpp index 60f4ca58b9ef..892d23c0ac68 100644 --- a/src/jsc/bindings/ScriptExecutionContext.cpp +++ b/src/jsc/bindings/ScriptExecutionContext.cpp @@ -302,6 +302,50 @@ extern "C" JSC::JSGlobalObject* ScriptExecutionContextIdentifier__getGlobalObjec return context->globalObject(); } +// True while the context exists and has not been marked terminating. +extern "C" bool ScriptExecutionContext__isAlive(ScriptExecutionContextIdentifier id) +{ + if (!id) return false; + Locker locker { allScriptExecutionContextsMapLock }; + auto* context = allScriptExecutionContextsMap().get(id); + return context && !context->isTerminating(); +} + +extern "C" void Bun__EventLoop__enqueueConcurrentTask(JSC::JSGlobalObject*, void* task); + +// postTaskTo() for a Rust-allocated ConcurrentTask. Returns true if the task +// was enqueued (ownership transferred to the target context's concurrent +// queue); false if the context is gone or terminating, in which case the +// caller retains ownership. The map lock is held across the enqueue, which is +// what serializes this with markTerminating(): either the poster's whole +// critical section runs before worker shutdown's markTerminating() (task +// enqueued, and shutdown's subsequent release_queued_tasks_for_shutdown() +// drain observes it), or after (isTerminating() true; poster drops the task +// without touching the freed VM). See markTerminating()'s comment for the +// ordering argument. +extern "C" bool ScriptExecutionContext__postConcurrentTask(ScriptExecutionContextIdentifier id, void* task) +{ + if (!id) return false; + Locker locker { allScriptExecutionContextsMapLock }; + auto* context = allScriptExecutionContextsMap().get(id); + if (!context || context->isTerminating()) + return false; + Bun__EventLoop__enqueueConcurrentTask(context->globalObject(), task); + return true; +} + +// Decrement the target context's event-loop concurrent refcount under the +// contexts-map lock; no-op when the context is gone or terminating. +extern "C" void ScriptExecutionContext__unrefEventLoopConcurrently(ScriptExecutionContextIdentifier id) +{ + if (!id) return; + Locker locker { allScriptExecutionContextsMapLock }; + auto* context = allScriptExecutionContextsMap().get(id); + if (!context || context->isTerminating()) + return; + Bun__eventLoop__incrementRefConcurrently(WebCore::clientData(context->vm())->bunVM, -1); +} + extern "C" void ScriptExecutionContext__markTerminating(JSC::JSGlobalObject* globalObject) { if (auto* context = defaultGlobalObject(globalObject)->scriptExecutionContext()) diff --git a/src/jsc/virtual_machine_exports.rs b/src/jsc/virtual_machine_exports.rs index 749067bd49ce..1e84cee2346a 100644 --- a/src/jsc/virtual_machine_exports.rs +++ b/src/jsc/virtual_machine_exports.rs @@ -141,6 +141,23 @@ pub fn queue_task_concurrently(global: &JSGlobalObject, task: *mut crate::cpp_ta } } +/// Called from C++ `ScriptExecutionContext__postConcurrentTask` with the +/// contexts-map lock held, so `global`'s VM/`EventLoop` are live for this call. +// HOST_EXPORT(Bun__EventLoop__enqueueConcurrentTask, c) +#[allow(clippy::not_unsafe_ptr_arg_deref)] +pub fn event_loop_enqueue_concurrent_task( + global: &JSGlobalObject, + task: *mut crate::event_loop::ConcurrentTaskItem, +) { + crate::mark_binding!(); + // SAFETY: see fn doc — VM/EventLoop live under the held contexts-map lock; + // `task` is the non-null heap allocation forwarded from Rust. + unsafe { + (*(*global.bun_vm_concurrently()).event_loop()) + .enqueue_task_concurrent(core::ptr::NonNull::new_unchecked(task)); + } +} + // HOST_EXPORT(Bun__handleRejectedPromise, c) pub fn handle_rejected_promise(global: &JSGlobalObject, promise: &mut JSPromise) { crate::mark_binding!(); diff --git a/src/runtime/api/Archive.rs b/src/runtime/api/Archive.rs index 39cbbdffc4d6..3bbff9c9fae2 100644 --- a/src/runtime/api/Archive.rs +++ b/src/runtime/api/Archive.rs @@ -11,6 +11,7 @@ use bun_event_loop::{TaskTag, Taskable, task_tag}; use bun_glob as glob; use bun_io::KeepAlive; use bun_jsc::ConcurrentTask::{AutoDeinit, ConcurrentTask}; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{ self as jsc, CallFrame, JSGlobalObject, JSMap, JSPromise, JSPromiseStrong, JSValue, JsResult, @@ -682,7 +683,7 @@ pub trait TaskContext: Send { pub struct AsyncTask { ctx: C, promise: JSPromiseStrong, - vm: *mut VirtualMachine, + context_id: ScriptExecutionContextIdentifier, task: WorkPoolTask, concurrent_task: ConcurrentTask, keep_alive: KeepAlive, @@ -694,15 +695,10 @@ impl Taskable for AsyncTask { impl AsyncTask { fn create(global: &JSGlobalObject, ctx: C) -> Result<*mut Self, bun_alloc::AllocError> { - // `bun_vm_ptr()` returns `*mut VirtualMachine` with write provenance; valid for - // process lifetime. Do NOT launder `bun_vm()` (a `&VirtualMachine`) through - // `*const _ as *mut _` — that derives a writeable pointer from a shared - // reference and is UB under Stacked Borrows. - let vm: *mut VirtualMachine = global.bun_vm_ptr(); let this = Box::new(AsyncTask { ctx, promise: JSPromiseStrong::init(global), - vm, + context_id: global.script_execution_context_identifier(), task: WorkPoolTask { callback: Self::run_callback, node: Default::default(), @@ -746,12 +742,14 @@ impl AsyncTask { let this: *mut Self = unsafe { bun_core::from_field_ptr!(Self, task, work_task) }; // SAFETY: thread-pool has exclusive access to ctx until it enqueues the concurrent task. unsafe { (*this).ctx.run() }; - // SAFETY: vm points to the live owning VM; concurrent_task is intrusive on the same allocation. + // SAFETY: concurrent_task is intrusive on the same allocation. unsafe { + let context_id = (*this).context_id; let ct = core::ptr::NonNull::from( (*this).concurrent_task.from(this, AutoDeinit::ManualDeinit), ); - (*(*this).vm).enqueue_task_concurrent(ct); + // Abandon: JSC handles cannot drop off-thread, leak the box (ct is intrusive). + let _ = context_id.post_concurrent_task(ct); } } diff --git a/src/runtime/api/JSBundler.rs b/src/runtime/api/JSBundler.rs index cfbd22382da3..c9d810473d72 100644 --- a/src/runtime/api/JSBundler.rs +++ b/src/runtime/api/JSBundler.rs @@ -1339,8 +1339,6 @@ pub mod js_bundler { let mut plugins: Option<*mut Plugin> = None; let config = Config::from_js(global_this, arguments[0], &mut plugins)?; - let event_loop = vm.event_loop(); - // `BundleV2.generateFromJavaScript` — the completion-task struct lives in // `crate::api::js_bundle_completion_task` (bun_runtime owns it because its // fields name `Config`/`Plugin`/`HTMLBundle::Route`; lower-tier crates @@ -1350,7 +1348,6 @@ pub mod js_bundler { config, plugins.and_then(core::ptr::NonNull::new), global_this, - event_loop, ) .map_err(|_| JsError::OutOfMemory)?; // SAFETY: `completion` is the freshly-boxed allocation returned above; diff --git a/src/runtime/api/js_bundle_completion_task.rs b/src/runtime/api/js_bundle_completion_task.rs index 4932a07891c6..3c1b7e5b6422 100644 --- a/src/runtime/api/js_bundle_completion_task.rs +++ b/src/runtime/api/js_bundle_completion_task.rs @@ -25,7 +25,8 @@ use bun_core::env::OperatingSystem; use bun_io::KeepAlive; use bun_jsc::AnyTask::AnyTask; use bun_jsc::WorkPool; -use bun_jsc::event_loop::EventLoop; + +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::{self as jsc, JSGlobalObject, JSPromise, JSValue}; use bun_options_types::WindowsOptions; use bun_options_types::schema::api; @@ -58,9 +59,7 @@ pub struct JSBundleCompletionTask { // `unsafe impl Send` below for the thread-affinity constraint this imposes. pub ref_count: RefCount, pub config: JSBundlerConfig, - // BACKREF — the JS-thread `EventLoop` outlives every completion task; safe - // `Deref` so call sites read `self.jsc_event_loop.enqueue_task_concurrent(..)`. - pub jsc_event_loop: BackRef, + pub context_id: ScriptExecutionContextIdentifier, pub task: AnyTask, pub global_this: BackRef, pub promise: jsc::JSPromiseStrong, @@ -115,16 +114,13 @@ pub(crate) fn create_and_schedule_completion_task( config: JSBundlerConfig, plugins: Option>, global_this: &JSGlobalObject, - event_loop: *mut EventLoop, ) -> crate::Result<*mut JSBundleCompletionTask> { let vm = global_this.bun_vm_ptr(); let env = global_this.bun_vm().transpiler.env; let completion = bun_core::heap::into_raw(Box::new(JSBundleCompletionTask { ref_count: RefCount::init(), config, - // `event_loop` is the live JS-thread loop (caller derives it from - // `vm.event_loop()`); never null once `Bun.build` is reachable. - jsc_event_loop: BackRef::from(core::ptr::NonNull::new(event_loop).expect("event_loop")), + context_id: global_this.script_execution_context_identifier(), task: AnyTask::default(), global_this: BackRef::new(global_this), promise: jsc::JSPromiseStrong::default(), @@ -761,13 +757,14 @@ fn from_completion_handle<'a>(c: NonNull) -> &'a JSBundleCo static COMPLETION_VTABLE: dispatch::CompletionDispatch = dispatch::CompletionDispatch { result_is_err: |c| matches!(from_completion_handle(c).result, BundleV2Result::Err(_)), enqueue_task_concurrent: |c, task| { - // `jsc_event_loop` is a `BackRef` — safe Deref. - // SAFETY: `task` is a fresh heap-allocated non-null `ConcurrentTaskItem` - // passed through from the bundler vtable; the queue takes ownership. - unsafe { - from_completion_handle(c) - .jsc_event_loop - .enqueue_task_concurrent(core::ptr::NonNull::new_unchecked(task)) + // SAFETY: `task` is a fresh heap-allocated non-null `ConcurrentTaskItem`. + let node = unsafe { core::ptr::NonNull::new_unchecked(task) }; + if !from_completion_handle(c) + .context_id + .post_concurrent_task(node) + { + // SAFETY: ownership not transferred; reclaim the heap node. + drop(unsafe { bun_core::heap::take(node.as_ptr()) }); } }, }; @@ -979,11 +976,12 @@ impl CompletionStruct for JSBundleCompletionTask { } fn complete_on_bundle_thread(&mut self) { - // `jsc_event_loop` is a `BackRef` — safe Deref. - // `ConcurrentTask::create` heap-allocates a fresh task; the - // queue takes ownership of it. - self.jsc_event_loop - .enqueue_task_concurrent(jsc::ConcurrentTask::create(self.task.task())); + let task = jsc::ConcurrentTask::create(self.task.task()); + if !self.context_id.post_concurrent_task(task) { + // Abandon: JSC handles cannot drop off-thread, leak the box. + // SAFETY: ownership not transferred; `task` is the fresh `create` heap node above. + drop(unsafe { bun_core::heap::take(task.as_ptr()) }); + } } fn set_result(&mut self, result: BundleV2Result) { self.result = result; diff --git a/src/runtime/crypto/PasswordObject.rs b/src/runtime/crypto/PasswordObject.rs index 9fe555e60420..eee4d40ec0fb 100644 --- a/src/runtime/crypto/PasswordObject.rs +++ b/src/runtime/crypto/PasswordObject.rs @@ -9,12 +9,12 @@ use bun_jsc::{ }; // `bun_jsc::{AnyTask, ConcurrentTask, EventLoop}` are *modules* (re-exported from // `bun_event_loop`); pull the concrete types out by name. -use bun_jsc::event_loop::EventLoop; // JSC-side ZigString carries `to_js` (the `bun_core::ZigString` repr-twin // lives in `bun_jsc::zig_string`); used for ASCII→JS conversions only. use bun_jsc::AnyTask::{AnyTask, JsResult as AnyTaskJsResult}; use bun_jsc::ConcurrentTask::ConcurrentTask; use bun_jsc::ZigStringJsc as _; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::zig_string::ZigString as JscZigString; use bun_jsc::{JSPromise, JSPromiseStrong}; use bun_threading::work_pool::WorkPool; @@ -549,7 +549,7 @@ struct PasswordJob { op: Op, password: Box<[u8]>, promise: JSPromiseStrong, - event_loop: *mut EventLoop, + context_id: ScriptExecutionContextIdentifier, global: *const JSGlobalObject, r#ref: KeepAlive, task: WorkPoolTask, @@ -585,13 +585,13 @@ impl PasswordJob { unsafe { (*result).task = AnyTask::from_typed(result, PasswordResult::::run_from_js_erased); } - // SAFETY: `event_loop` was stored from the JS-thread VM and outlives the - // job; ownership of `result` transfers to the event loop here. `task` is - // an intrusive field at a stable address. - unsafe { - (*self.event_loop).enqueue_task_concurrent(ConcurrentTask::create_from( - core::ptr::addr_of_mut!((*result).task), - )); + // SAFETY: `result.task` is an intrusive field at a stable heap address. + let node = ConcurrentTask::create_from(unsafe { core::ptr::addr_of_mut!((*result).task) }); + if !self.context_id.post_concurrent_task(node) { + // Abandon: free node + result box; leak the `JSPromiseStrong` (JSC heap is dead). + drop(unsafe { bun_core::heap::take(node.as_ptr()) }); + let mut r = unsafe { bun_core::heap::take(result) }; + core::mem::forget(core::mem::take(&mut r.promise)); } // `self: Box` drops here; Drop runs secure_zero on password (+op). } @@ -670,8 +670,7 @@ impl JSPasswordObject { op, password, promise, - // SAFETY: bun_vm() is non-null for a Bun-owned global; VM outlives the job. - event_loop: global_object.bun_vm().event_loop(), + context_id: global_object.script_execution_context_identifier(), global: std::ptr::from_ref(global_object), r#ref: KeepAlive::default(), task: WorkPoolTask::default(), diff --git a/src/runtime/node/node_fs.rs b/src/runtime/node/node_fs.rs index 3ee2ace9e788..17c6a833f0a4 100644 --- a/src/runtime/node/node_fs.rs +++ b/src/runtime/node/node_fs.rs @@ -17,6 +17,7 @@ use bun_io::KeepAlive; use bun_jsc::AbortSignal; use bun_jsc::EventLoopTaskPtr; use bun_jsc::debugger::AsyncTaskTracker; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{EventLoopHandle, JSGlobalObject, JSValue, JsResult, Task, ThreadSafe, Unprotect}; use bun_paths::{self as paths, OSPathBuffer, OSPathChar, OSPathSliceZ, PathBuffer}; @@ -1238,6 +1239,7 @@ mod _async_tasks { /// Wrapped in [`ThreadSafe`] so the paired `unprotect()` runs on drop. pub args: ThreadSafe, pub global_object: bun_ptr::BackRef, + pub context_id: ScriptExecutionContextIdentifier, pub task: WorkPoolTask, pub result: Maybe, pub r#ref: KeepAlive, @@ -1281,6 +1283,7 @@ mod _async_tasks { // niche-optimised; never construct an all-zero `Result` value. result: Err(sys::Error::default()), global_object: bun_ptr::BackRef::new(global_object), + context_id: global_object.script_execution_context_identifier(), task: work_pool_task(Self::work_pool_callback), r#ref: KeepAlive::default(), tracker: AsyncTaskTracker::init(vm), @@ -1304,16 +1307,12 @@ mod _async_tasks { // `sys::Error::path` is `Box<[u8]>` boxed at the // `errno_sys_p` construction site, so no clone is needed — `node_fs` may drop. - // `bun_vm_concurrently()` skips the JS-thread debug assert and is the - // documented accessor for off-thread (work-pool) callers; the - // event-loop's concurrent queue is MPSC-safe. - let vm = this.global_object().bun_vm_concurrently(); - // SAFETY: VirtualMachine and its event loop are process-static - // (LIFETIMES.tsv); the concurrent queue is MPSC-safe. - unsafe { - (*(*vm).event_loop()).enqueue_task_concurrent(ConcurrentTask::create_from( - std::ptr::from_mut::(this), - )); + let context_id = this.context_id; + let node = ConcurrentTask::create_from(std::ptr::from_mut::(this)); + if !context_id.post_concurrent_task(node) { + // Abandon: JSC handles cannot drop off-thread, leak the box. + // SAFETY: ownership not transferred; `node` was `create_from`-allocated above. + drop(unsafe { bun_core::heap::take(node.as_ptr()) }); } } @@ -1395,6 +1394,7 @@ mod _async_tasks { /// Wrapped in [`ThreadSafe`] so the paired `unprotect()` runs on drop. pub args: ThreadSafe, pub evtloop: EventLoopHandle, + pub context_id: ScriptExecutionContextIdentifier, pub task: WorkPoolTask, /// Written from any workpool thread (first `finish_concurrently` caller wins via /// `has_result` CAS); read on the JS thread in `run_from_js_thread`. Wrapped in @@ -1582,6 +1582,7 @@ mod _async_tasks { result: core::cell::Cell::new(Ok(())), // `vm.event_loop` is the live per-thread `jsc::EventLoop` field. evtloop: EventLoopHandle::init(vm.event_loop.cast()), + context_id: global_object.script_execution_context_identifier(), task: work_pool_task(Self::work_pool_callback), r#ref: KeepAlive::default(), tracker: AsyncTaskTracker::init(vm), @@ -1615,6 +1616,7 @@ mod _async_tasks { // `has_result` CAS) before any read on the JS thread. result: core::cell::Cell::new(Ok(())), evtloop: EventLoopHandle::init_mini(mini), + context_id: ScriptExecutionContextIdentifier(0), task: work_pool_task(Self::work_pool_callback), r#ref: KeepAlive::default(), tracker: AsyncTaskTracker { id: 0 }, @@ -1690,14 +1692,17 @@ mod _async_tasks { // provenance from `Box::leak`, so the enqueued callback may safely // form `&mut *this` on the JS thread. if matches!(this_ref.evtloop, EventLoopHandle::Js { .. }) { - this_ref.evtloop.enqueue_task_concurrent(EventLoopTaskPtr { - js: ConcurrentTask::from_callback(this, |p| { - // SAFETY: `p` is the `Box::leak`'d task; subtask count hit zero so this - // JS-thread callback holds the only live reference (exclusive `&mut`). - unsafe { (&mut *p).run_from_js_thread().map_err(Into::into) } - }) - .as_ptr(), + let node = ConcurrentTask::from_callback(this, |p| { + // SAFETY: `p` is the `Box::leak`'d task; subtask count hit zero so this + // JS-thread callback holds the only live reference (exclusive `&mut`). + unsafe { (&mut *p).run_from_js_thread().map_err(Into::into) } }); + if !this_ref.context_id.post_concurrent_task(node) { + // SAFETY: ownership not transferred; `node` is a `from_callback` allocation. + unsafe { + bun_event_loop::ConcurrentTask::ConcurrentTask::destroy_from_callback(node) + }; + } } else { this_ref.evtloop.enqueue_task_concurrent(EventLoopTaskPtr { mini: AnyTaskWithExtraContext::from_callback_auto_deinit( @@ -2163,6 +2168,7 @@ mod _async_tasks { /// Wrapped in [`ThreadSafe`] so the paired `unprotect()` runs on drop. pub args: ThreadSafe, pub global_object: bun_ptr::BackRef, + pub context_id: ScriptExecutionContextIdentifier, pub task: WorkPoolTask, pub r#ref: KeepAlive, pub tracker: AsyncTaskTracker, @@ -2358,6 +2364,7 @@ mod _async_tasks { args: FsArgument::into_thread_safe(args), has_result: AtomicBool::new(false), global_object: bun_ptr::BackRef::new(global_object), + context_id: global_object.script_execution_context_identifier(), task: work_pool_task(Self::work_pool_callback), r#ref: KeepAlive::default(), tracker: AsyncTaskTracker::init(vm), @@ -2533,16 +2540,13 @@ mod _async_tasks { } } - // `bun_vm_concurrently()` skips the JS-thread debug assert and is the - // documented accessor for off-thread (work-pool) callers. - // SAFETY: `bun_vm_concurrently()` returns the process-singleton VM; - // sole `&mut` borrow at this point on the work-pool thread. - let vm = unsafe { &mut *self.global_object().bun_vm_concurrently() }; - // `ConcurrentTask::create` heap-allocates a fresh task; the - // queue takes ownership of it. - vm.enqueue_task_concurrent(ConcurrentTask::create(Task::init(std::ptr::from_mut::< - Self, - >(self)))); + let context_id = self.context_id; + let node = ConcurrentTask::create(Task::init(std::ptr::from_mut::(self))); + if !context_id.post_concurrent_task(node) { + // Abandon: JSC handles cannot drop off-thread, leak the box. + // SAFETY: ownership not transferred; `node` was `create`-allocated above. + drop(unsafe { bun_core::heap::take(node.as_ptr()) }); + } } fn clear_result_list(&mut self) { diff --git a/src/runtime/node/node_zlib_binding.rs b/src/runtime/node/node_zlib_binding.rs index 773615da2ecb..bc71e96c62e4 100644 --- a/src/runtime/node/node_zlib_binding.rs +++ b/src/runtime/node/node_zlib_binding.rs @@ -9,6 +9,7 @@ use bun_core::{String as BunString, ZigStringSlice}; use bun_event_loop::Taskable; use bun_io::KeepAlive; use bun_jsc::ConcurrentTask::{ConcurrentTask, Task}; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{ self as jsc, CallFrame, ErrorCode, JSGlobalObject, JSValue, JsCell, JsResult, StringJsc as _, @@ -210,6 +211,8 @@ pub(crate) trait CompressionStreamImpl: Sized + Taskable + 'static { /// Implementations store a `BackRef`; the single unsafe /// deref lives in `BackRef::get`, so callers and impls are safe. fn global_this(&self) -> &JSGlobalObject; + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + fn context_id(&self) -> ScriptExecutionContextIdentifier; fn stream(&self) -> &JsCell; /// Write `(avail_out, avail_in)` into the JS-owned 2-element `Uint32Array` @@ -472,24 +475,19 @@ impl CompressionStream { // `ref_()` in `write()`); bodies use the `&self` accessor surface // (R-2). `ParentRef` Deref collapses the per-site raw deref. let this_ref = ParentRef::from(NonNull::new(this).expect("async_job_run: this")); - let global_this: &JSGlobalObject = this_ref.global_this(); - // `bun_vm_concurrently()` is the thread-safe accessor (skips the - // JS-thread debug assert; same backing pointer as `bun_vm()`). - // BACKREF — `bun_vm_concurrently()` never returns null for a Bun-owned - // global; wrap once so the `event_loop()` read below is safe Deref. - let vm = ParentRef::from( - NonNull::new(global_this.bun_vm_concurrently()).expect("bun_vm_concurrently"), - ); + let context_id = this_ref.context_id(); - this_ref.stream().with_mut(|s| s.do_work()); + // `do_work()` writes into the pinned JS ArrayBuffer; best-effort skip on shutdown. + if context_id.is_alive() { + this_ref.stream().with_mut(|s| s.do_work()); + } - // SAFETY: `event_loop()` is a self-pointer into a live VM; the - // `enqueue_task_concurrent` body only touches the lock-free - // `concurrent_tasks` queue (thread-safe). `this` is the heap-allocated - // `m_ctx` payload — the matching `ref()` in `write()` keeps it alive - // until `run_from_js_thread` runs and calls `deref()`. - unsafe { - (*vm.event_loop()).enqueue_task_concurrent(ConcurrentTask::create(Task::init(this))); + // `this` is the heap `m_ctx` payload; the `ref()` in `write()` keeps it alive. + let node = ConcurrentTask::create(Task::init(this)); + if !context_id.post_concurrent_task(node) { + // Abandon: `deref()` is not safe off-thread, leak `this`. + // SAFETY: ownership not transferred; `node` was `create`-allocated above. + drop(unsafe { bun_core::heap::take(node.as_ptr()) }); } } @@ -998,6 +996,7 @@ macro_rules! __impl_compression_stream { type Stream = $ctx; #[inline] fn global_this(&self) -> &::bun_jsc::JSGlobalObject { self.global_this.get() } + #[inline] fn context_id(&self) -> ::bun_jsc::js_global_object::ScriptExecutionContextIdentifier { self.context_id } #[inline] fn stream(&self) -> &::bun_jsc::JsCell { &self.stream } #[inline] fn poll_ref(&self) -> &::bun_jsc::JsCell<$crate::node::node_zlib_binding::CountedKeepAlive> { &self.poll_ref } #[inline] fn this_value(&self) -> &::bun_jsc::JsCell<::bun_jsc::StrongOptional> { &self.this_value } diff --git a/src/runtime/node/zlib/NativeBrotli.rs b/src/runtime/node/zlib/NativeBrotli.rs index 7a5b889ba935..b277253bed29 100644 --- a/src/runtime/node/zlib/NativeBrotli.rs +++ b/src/runtime/node/zlib/NativeBrotli.rs @@ -85,6 +85,7 @@ mod _impl { // JSC_BORROW backref; global outlives this m_ctx payload. `BackRef` // centralises the single unsafe deref so the trait impl is safe. pub global_this: bun_ptr::BackRef, + pub context_id: bun_jsc::js_global_object::ScriptExecutionContextIdentifier, pub stream: JsCell, pub poll_ref: JsCell, // TODO: Strong self-ref on the wrapper → JsRef per PORTING.md §JSC (Strong back-ref to own wrapper leaks) @@ -148,6 +149,7 @@ mod _impl { ref_count: Cell::new(1), // JSC_BORROW backref — the global outlives this m_ctx payload. global_this: bun_ptr::BackRef::new(global_this), + context_id: global_this.script_execution_context_identifier(), stream: JsCell::new(stream), poll_ref: JsCell::new(CountedKeepAlive::default()), this_value: JsCell::new(StrongOptional::empty()), diff --git a/src/runtime/node/zlib/NativeZlib.rs b/src/runtime/node/zlib/NativeZlib.rs index 2abe06797c9c..f1ac1afd98d1 100644 --- a/src/runtime/node/zlib/NativeZlib.rs +++ b/src/runtime/node/zlib/NativeZlib.rs @@ -42,6 +42,7 @@ mod _impl { // JSC_BORROW backref; global outlives this m_ctx payload. `BackRef` // centralises the single unsafe deref so the trait impl is safe. pub global_this: bun_ptr::BackRef, + pub context_id: bun_jsc::js_global_object::ScriptExecutionContextIdentifier, pub stream: JsCell, pub poll_ref: JsCell, pub this_value: JsCell, // jsc.Strong.Optional @@ -93,6 +94,7 @@ mod _impl { ref_count: Cell::new(1), // JSC_BORROW backref — the global outlives this m_ctx payload. global_this: bun_ptr::BackRef::new(global), + context_id: global.script_execution_context_identifier(), stream: JsCell::new(stream), poll_ref: JsCell::new(CountedKeepAlive::default()), this_value: JsCell::new(StrongOptional::empty()), diff --git a/src/runtime/node/zlib/NativeZstd.rs b/src/runtime/node/zlib/NativeZstd.rs index 8ca905cf12f4..6b2b2320963b 100644 --- a/src/runtime/node/zlib/NativeZstd.rs +++ b/src/runtime/node/zlib/NativeZstd.rs @@ -41,6 +41,7 @@ mod _impl { // LIFETIMES.tsv: JSC_BORROW. The global outlives this m_ctx payload; // `BackRef` centralises the single unsafe deref so the trait impl is safe. pub global_this: bun_ptr::BackRef, + pub context_id: bun_jsc::js_global_object::ScriptExecutionContextIdentifier, pub stream: JsCell, pub poll_ref: JsCell, pub this_value: JsCell, // jsc.Strong.Optional @@ -105,6 +106,7 @@ mod _impl { // JSC_BORROW — the JSGlobalObject outlives this payload (the C++ // wrapper is owned by that global's heap). global_this: bun_ptr::BackRef::new(global), + context_id: global.script_execution_context_identifier(), stream: JsCell::new(stream), poll_ref: JsCell::new(CountedKeepAlive::default()), this_value: JsCell::new(StrongOptional::empty()), diff --git a/src/runtime/server/HTMLBundle.rs b/src/runtime/server/HTMLBundle.rs index 937785d53a27..28de23af4537 100644 --- a/src/runtime/server/HTMLBundle.rs +++ b/src/runtime/server/HTMLBundle.rs @@ -498,8 +498,7 @@ impl Route { } config.source_map = bundler_options::SourceMapOption::Linked; - let completion_task = - create_and_schedule_completion_task(config, plugins, global, vm.event_loop())?; + let completion_task = create_and_schedule_completion_task(config, plugins, global)?; // SAFETY: `completion_task` is the freshly-boxed allocation (refcount==1); sole owner. unsafe { (*completion_task).started_at_ns = diff --git a/src/runtime/webcore/fetch/FetchTasklet.rs b/src/runtime/webcore/fetch/FetchTasklet.rs index 9bff5a530a83..67422cb6c9d0 100644 --- a/src/runtime/webcore/fetch/FetchTasklet.rs +++ b/src/runtime/webcore/fetch/FetchTasklet.rs @@ -17,6 +17,7 @@ use bun_http::{ }; use bun_io::KeepAlive; use bun_jsc::debugger::AsyncTaskTracker; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::virtual_machine::VirtualMachine; use bun_jsc::{ self as jsc, GlobalRef, JSGlobalObject, JSValue, JsResult, StringJsc, StrongOptional, @@ -69,6 +70,7 @@ pub struct FetchTasklet { pub metadata: Option, pub javascript_vm: &'static VirtualMachine, pub global_this: GlobalRef, + pub context_id: ScriptExecutionContextIdentifier, pub request_body: HTTPRequestBody, // ThreadSafeStreamBuffer is intrusively refcounted (`ref_count: AtomicU32`, // starts at 2) and shared with the HTTP thread via raw ptr; `Arc` can't be mutably @@ -295,17 +297,13 @@ impl FetchTasklet { unsafe { &*this } } - /// Enqueue a concurrent task on the JS-thread event loop. - /// - /// Centralises the `(*vm.event_loop()).enqueue_task_concurrent(..)` raw - /// deref. `event_loop()` returns a self-ptr into the VirtualMachine that - /// is valid for the VM's lifetime; `enqueue_task_concurrent` takes `&self` - /// and is thread-safe (lock-free MPSC push). `task` is a live - /// `ConcurrentTaskItem` that the queue takes ownership of via its - /// intrusive `next` link. + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. #[inline] - fn enqueue_concurrent(vm: &VirtualMachine, task: core::ptr::NonNull) { - vm.event_loop_shared().enqueue_task_concurrent(task); + fn enqueue_concurrent( + context_id: ScriptExecutionContextIdentifier, + task: core::ptr::NonNull, + ) -> bool { + context_id.post_concurrent_task(task) } /// Wrap a borrowed body chunk in a `StreamResult::Temporary*` for @@ -393,7 +391,7 @@ impl FetchTasklet { return; } let self_ = Self::from_raw_ref(this); - if self_.javascript_vm.is_shutting_down() { + if !self_.context_id.is_alive() { // SAFETY: last ref; exclusive access. `deinit()` would run // `clear_data()` + `Drop` for the JSC `Strong`/`Weak` fields, which // reach into the VM's HandleSet from this (HTTP) thread — not @@ -406,10 +404,13 @@ impl FetchTasklet { // lets make sure that we always call deinit from main thread // `from_callback` heap-allocates a fresh `ConcurrentTaskItem`; the queue // takes ownership of it. - Self::enqueue_concurrent( - self_.javascript_vm, - ConcurrentTask::from_callback(this, FetchTasklet::deinit_callback), - ); + let node = ConcurrentTask::from_callback(this, FetchTasklet::deinit_callback); + if !Self::enqueue_concurrent(self_.context_id, node) { + // SAFETY: ownership not transferred; `node` is a `from_callback` allocation. + unsafe { ConcurrentTask::destroy_from_callback(node) }; + // SAFETY: last ref; see the `!is_alive()` branch above. + unsafe { FetchTasklet::dealloc_for_shutdown(this) }; + } } // ConcurrentTask::from_callback takes `fn(*mut T) -> bun_event_loop::JsResult<()>` @@ -1879,6 +1880,7 @@ impl FetchTasklet { metadata: None, javascript_vm: jsc_vm, global_this: GlobalRef::from(global_this), + context_id: global_this.script_execution_context_identifier(), request_body: fetch_options.body, request_body_streaming_buffer: None, response_buffer: MutableString::default(), @@ -2131,17 +2133,19 @@ impl FetchTasklet { /// This is ALWAYS called from the http thread and we cannot touch the buffer here because is locked pub(crate) fn on_write_request_data_drain(this: *mut FetchTasklet) { let this_ref = Self::from_raw_ref(this); - if this_ref.javascript_vm.is_shutting_down() { + if !this_ref.context_id.is_alive() { return; } // ref until the main thread callback is called this_ref.ref_(); // `from_callback` heap-allocates a fresh `ConcurrentTaskItem`; the queue // takes ownership of it. - Self::enqueue_concurrent( - this_ref.javascript_vm, - ConcurrentTask::from_callback(this, FetchTasklet::resume_request_data_stream), - ); + let node = ConcurrentTask::from_callback(this, FetchTasklet::resume_request_data_stream); + if !Self::enqueue_concurrent(this_ref.context_id, node) { + // SAFETY: ownership not transferred; `node` is a `from_callback` allocation. + unsafe { ConcurrentTask::destroy_from_callback(node) }; + FetchTasklet::deref_from_thread(this); + } } /// This is ALWAYS called from the main thread @@ -2451,36 +2455,28 @@ impl FetchTasklet { } } // will deinit when done with the http client (when is_done = true) - if task_ref.javascript_vm.is_shutting_down() { - // VM teardown: the JS-thread side will never drain this buffer (its - // on_progress_update bails the same way), so free the body bytes now. + // Abandon path: Rust-side cleanup the JS-thread `on_progress_update` would + // never reach (body buffer, parked socket, CAS flag, unlock, final derefs). + let abandon = |task_ref: &mut FetchTasklet| { task_ref.scheduled_response_buffer = MutableString::default(); - // The certificate will never be checked; release the parked - // socket instead of leaving it occupying an active request slot - // until the idle timeout. if task_ref.result.certificate_info.take().is_some() { if let Some(http_) = task_ref.http.as_mut() { http::http_thread().schedule_shutdown(http_); } } - // We won the `has_schedule_callback` CAS above but are not - // enqueueing the on_progress_update task; undo the flag so a later - // (final) callback can re-enter this branch instead of taking the - // already-scheduled early return. task_ref .has_schedule_callback .store(false, Ordering::Release); task_ref.mutex.unlock(); if is_done { - // No on_progress_update will ever run for this final result, so - // release the JS-side ref it would have dropped, then the - // HTTP-side ref. The 1→0 transition runs `dealloc_for_shutdown` - // (Rust boxes only — JSC handles are leaked to destructOnExit). // SAFETY: `task` is the live heap tasklet; both refs held. FetchTasklet::deref_from_thread(task); // SAFETY: second ref still held until this 1→0 transition. FetchTasklet::deref_from_thread(task); } + }; + if !task_ref.context_id.is_alive() { + abandon(task_ref); return; } let ct = core::ptr::NonNull::from( @@ -2490,7 +2486,10 @@ impl FetchTasklet { ); // `ct` is the inline `concurrent_task` field of the heap tasklet; the // queue takes ownership of its `next` link. - Self::enqueue_concurrent(task_ref.javascript_vm, ct); + if !Self::enqueue_concurrent(task_ref.context_id, ct) { + abandon(task_ref); + return; + } task_ref.mutex.unlock(); // we are done with the http client so we can deref our side diff --git a/src/runtime/webcore/s3/client.rs b/src/runtime/webcore/s3/client.rs index 80271e2b89eb..5be541d367d9 100644 --- a/src/runtime/webcore/s3/client.rs +++ b/src/runtime/webcore/s3/client.rs @@ -304,6 +304,9 @@ pub(crate) fn list_objects( callback: s3_simple_request::Callback::ListObjects(callback), headers, vm: Some(bun_ptr::BackRef::new(VirtualMachine::get())), + context_id: VirtualMachine::get() + .global() + .script_execution_context_identifier(), response_buffer: MutableString::default(), result: bun_http::HTTPClientResult::default(), concurrent_task: Default::default(), @@ -1009,8 +1012,9 @@ pub(crate) fn download_stream( callback, range: range.map(Vec::into_boxed_slice), headers, - // `VirtualMachine::get()` returns the live per-thread VM singleton. - vm: Some(bun_ptr::BackRef::new(VirtualMachine::get())), + context_id: VirtualMachine::get() + .global() + .script_execution_context_identifier(), has_schedule_callback: core::sync::atomic::AtomicBool::new(false), signal_store: Default::default(), signals: Default::default(), diff --git a/src/runtime/webcore/s3/download_stream.rs b/src/runtime/webcore/s3/download_stream.rs index 6863b5734e73..4f9c51710f18 100644 --- a/src/runtime/webcore/s3/download_stream.rs +++ b/src/runtime/webcore/s3/download_stream.rs @@ -7,7 +7,7 @@ use bun_event_loop::ConcurrentTask::{AutoDeinit, ConcurrentTask}; use bun_event_loop::{TaskTag, Taskable, task_tag}; use bun_http::{AsyncHTTP, HTTPClientResult, Headers, Signals}; use bun_io::KeepAlive; -use bun_jsc::virtual_machine::VirtualMachine; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_s3_signing::credentials::SignResult; use bun_s3_signing::error::S3Error; use bun_threading::Mutex; @@ -18,9 +18,8 @@ pub struct S3HttpDownloadStreamingTask { // `MaybeUninit` because `AsyncHTTP` contains non-null references, so // `mem::zeroed()` can't be used here (mirrors `S3HttpSimpleTask`). pub http: core::mem::MaybeUninit>, - /// JSC_BORROW: per-thread VM singleton, outlives every task. `None` only in - /// the inert `Default` placeholder (overwritten before the task escapes). - pub vm: Option>, + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + pub context_id: ScriptExecutionContextIdentifier, pub sign_result: SignResult, pub headers: Headers, pub callback_context: NonNull<()>, @@ -61,7 +60,7 @@ impl Default for S3HttpDownloadStreamingTask { Self { // never read — fully overwritten by `AsyncHTTP::init` before first use. http: core::mem::MaybeUninit::uninit(), - vm: None, + context_id: ScriptExecutionContextIdentifier(0), sign_result: SignResult::default(), headers: Headers::default(), callback_context: NonNull::dangling(), @@ -344,15 +343,8 @@ impl S3HttpDownloadStreamingTask { let task = core::ptr::NonNull::from( self_.concurrent_task.from(this, AutoDeinit::ManualDeinit), ); - // `vm` is the live per-thread VM BackRef captured at task creation; event_loop - // is initialized for the request's lifetime and enqueue is thread-safe (`&self`). - // `task` is the inline `concurrent_task` field of this heap request; - // the queue takes ownership of its `next` link. - self_ - .vm - .expect("vm set at task creation") - .event_loop_shared() - .enqueue_task_concurrent(task); + // Abandon: `Drop` touches the VM event loop, leak the box (task is intrusive). + let _ = self_.context_id.post_concurrent_task(task); } } } diff --git a/src/runtime/webcore/s3/simple_request.rs b/src/runtime/webcore/s3/simple_request.rs index b4114d0536f4..353b8fcfc6a3 100644 --- a/src/runtime/webcore/s3/simple_request.rs +++ b/src/runtime/webcore/s3/simple_request.rs @@ -10,6 +10,7 @@ use bun_http::{ Method, }; use bun_io::KeepAlive; +use bun_jsc::js_global_object::ScriptExecutionContextIdentifier; use bun_jsc::virtual_machine::VirtualMachine; use bun_picohttp as picohttp; use bun_s3_signing::acl::ACL; @@ -120,6 +121,8 @@ pub struct S3HttpSimpleTask { /// JSC_BORROW: per-thread VM singleton, outlives every task. `None` only in /// the inert `Default` placeholder (overwritten before the task escapes). pub vm: Option>, + /// See [`ScriptExecutionContextIdentifier::post_concurrent_task`]. + pub context_id: ScriptExecutionContextIdentifier, pub sign_result: SignResult, pub headers: Headers, pub callback_context: *mut c_void, @@ -156,6 +159,7 @@ impl Default for S3HttpSimpleTask { Self { http: core::mem::MaybeUninit::uninit(), vm: None, + context_id: ScriptExecutionContextIdentifier(0), sign_result: SignResult::default(), headers: Headers::default(), callback_context: core::ptr::null_mut(), @@ -487,14 +491,8 @@ impl S3HttpSimpleTask { this.concurrent_task .from(this_ptr, AutoDeinit::ManualDeinit), ); - // `vm` is the live per-thread VM BackRef captured at task creation; event_loop - // is set during VM init and outlives this task. `enqueue_task_concurrent` is `&self`. - // `task` is the inline `concurrent_task` field of this heap request; - // the queue takes ownership of its `next` link. - this.vm - .expect("vm set at task creation") - .event_loop_shared() - .enqueue_task_concurrent(task); + // Abandon: `Drop` touches the VM event loop, leak the box (task is intrusive). + let _ = this.context_id.post_concurrent_task(task); } } } @@ -635,6 +633,9 @@ pub(crate) fn execute_simple_s3_request( range: options.range, headers, vm: Some(bun_ptr::BackRef::new(VirtualMachine::get())), + context_id: VirtualMachine::get() + .global() + .script_execution_context_identifier(), response_buffer: MutableString::default(), result: HTTPClientResult::default(), concurrent_task: ConcurrentTask::default(), diff --git a/test/js/web/workers/worker-terminate-lifetime.test.ts b/test/js/web/workers/worker-terminate-lifetime.test.ts index 9d476d1f3d53..a368839d5c8f 100644 --- a/test/js/web/workers/worker-terminate-lifetime.test.ts +++ b/test/js/web/workers/worker-terminate-lifetime.test.ts @@ -1,5 +1,5 @@ import { expect, test } from "bun:test"; -import { bunEnv, bunExe, isASAN, isDebug } from "harness"; +import { bunEnv, bunExe, isASAN, isDebug, tempDir } from "harness"; import { join } from "path"; // Worker VM startup/teardown is much slower under debug and/or ASAN; these @@ -176,3 +176,91 @@ test.skipIf(!isASAN)( }, timeout, ); + +// Regression: off-thread completions (WorkPool / HTTP thread / bundle thread) +// posted back to a worker's EventLoop via a raw BackRef/&VirtualMachine +// captured at schedule time. worker.terminate() freed the VM box before the +// pool thread ran, so the enqueue dereferenced freed memory. Every cross-thread +// source is now posted by ScriptExecutionContextIdentifier under the contexts- +// map lock (serializing with markTerminating). ASAN-only: release builds read +// freed memory without crashing; under ASAN the heap-use-after-free is +// deterministic on the first teardown. +test.skipIf(!isASAN)( + "terminate() while cross-thread WorkPool completions are in flight does not UAF on enqueue", + async () => { + // Worker body: arm one in-flight op per cross-thread source, signal the + // parent, then sit. All loopback-only; results ignored. + const body = ` + import { parentPort, threadId } from "node:worker_threads"; + import fs from "node:fs"; import fsp from "node:fs/promises"; + import zlib from "node:zlib"; import crypto from "node:crypto"; + import path from "node:path"; + const sink = () => {}; const swallow = p => Promise.resolve(p).then(sink, sink); + // Scratch rooted under body.mjs's dir (the harness tempDir) so cleanup is automatic. + const tmp = path.join(path.dirname(new URL(import.meta.url).pathname), "scratch-" + threadId); + fs.mkdirSync(tmp, { recursive: true }); + const tf = path.join(tmp, "a.txt"); + fs.writeFileSync(tf, Buffer.alloc(1 << 14, "x").toString()); + // WorkTask + swallow(Bun.write(path.join(tmp, "b.txt"), Buffer.alloc(1 << 14, "y").toString())); + swallow(Bun.file(tf).text()); + // AsyncFSTask + swallow(fsp.readFile(tf)); swallow(fsp.stat(tf)); + swallow(fsp.readdir(tmp, { recursive: true })); + // AnyTaskJob (pbkdf2/scrypt/keygen/zstd) + crypto.pbkdf2("p", "s", 50000, 64, "sha512", sink); + crypto.scrypt("p", "saltsalt", 64, sink); + crypto.generateKeyPair("rsa", { modulusLength: 1024 }, sink); + swallow(Bun.zstdCompress(Buffer.alloc(1 << 14))); + // PasswordJob + swallow(Bun.password.hash("hunter2", { algorithm: "bcrypt", cost: 4 })); + // CompressionStream async_job_run + zlib.deflate(Buffer.alloc(1 << 16), sink); zlib.gzip(Buffer.alloc(1 << 16), sink); + // ConcurrentPromiseTask + swallow(new Bun.Transpiler({ loader: "ts" }).transform("const x: number = 1;")); + swallow(Array.fromAsync(new Bun.Glob("**/*").scan(tmp))); + // ConcurrentCppTask (WebCrypto) + swallow(crypto.subtle.digest("SHA-256", new Uint8Array(1 << 14))); + // JSBundleCompletionTask (+ COMPLETION_VTABLE.enqueue_task_concurrent via plugins) + const e = path.join(tmp, "e.ts"); fs.writeFileSync(e, "export const x=1;"); + swallow(Bun.build({ + entrypoints: [e], target: "bun", + plugins: [{ name: "p", setup(b) { b.onLoad({ filter: /\\.ts$/ }, () => undefined); } }], + })); + // FetchTasklet (HTTP thread) + const srv = Bun.serve({ port: 0, fetch: () => new Response(Buffer.alloc(1 << 12, "x")) }); + swallow(fetch("http://127.0.0.1:" + srv.port + "/").then(r => r.text())); + parentPort.postMessage("armed"); + setInterval(sink, 1 << 30); + `; + using dir = tempDir("worker-terminate-crossthread", { + "body.mjs": body, + "main.mjs": ` + import { Worker } from "node:worker_threads"; + const n = ${slow ? 4 : 12}; + for (let i = 0; i < n; i++) { + const w = new Worker(new URL("./body.mjs", import.meta.url)); + await new Promise((res, rej) => { + w.once("message", res); w.once("error", rej); + }); + await w.terminate(); + } + console.log("ok"); + `, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "main.mjs"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + // No UAF, no panic, no assert: stderr must be clean and the driver must + // have finished every round. + expect(stderr).toBe(""); + expect(stdout).toBe("ok\n"); + expect(exitCode).toBe(0); + }, + timeout * 2, +);