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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions src/jsc/VirtualMachine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,7 @@ unsafe extern "C" {
safe fn Bun__closeAllSQLiteDatabasesForTermination();
safe fn Bun__WebView__closeAllForTermination();
safe fn Zig__GlobalObject__destructOnExit(global: &JSGlobalObject);
safe fn Bun__JSCTaskScheduler__markShuttingDown(global: &JSGlobalObject);
}

pub const HOT_RELOAD_HOT: u8 = 1;
Expand Down Expand Up @@ -1561,6 +1562,13 @@ impl VirtualMachine {
(hooks.terminate_all_workers_and_wait)(10_000);
}

// Mirror web_worker.rs::shutdown(): fence DeferredWorkTimer
// producers before the drain so a cross-thread scheduleWorkSoon
// that raced the shutdown either enqueued (and is caught by the
// drain below) or observes the flag under m_lock and drops.
// destructOnExit sets it again (idempotently).
Bun__JSCTaskScheduler__markShuttingDown(self.global());

// Every worker has now posted its close task to our concurrent
// queue (OUTSTANDING is decremented after dispatchExit). Drop
// those queued lambdas — without running them — so the captured
Expand Down
77 changes: 65 additions & 12 deletions src/jsc/bindings/JSCTaskScheduler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,28 @@ static JSC::VM& getVM(Ticket& ticket)
return ticket->scriptExecutionOwner()->vm();
}

// Drop `ticket` from whichever pending set holds it. Caller holds m_lock; the
// event-loop ref is balanced after the caller releases the lock.
static bool dropPendingTicketLocked(Bun::JSCTaskScheduler& scheduler, Ticket ticket) WTF_REQUIRES_LOCK(scheduler.m_lock)
{
bool isKeepingEventLoopAlive = scheduler.m_pendingTicketsKeepingEventLoopAlive.removeIf([ticket](auto pendingTicket) {
return pendingTicket.ptr() == ticket;
});
// -- At this point, ticket may be an invalid pointer.
if (!isKeepingEventLoopAlive) {
scheduler.m_pendingTicketsOther.removeIf([ticket](auto pendingTicket) {
return pendingTicket.ptr() == ticket;
});
}
return isKeepingEventLoopAlive;
}

void JSCTaskScheduler::onAddPendingWork(WebCore::JSVMClientData* clientData, Ref<TicketData>&& ticket, JSC::DeferredWorkTimer::WorkType kind)
{
auto& scheduler = clientData->deferredWorkTimer;
Locker<Lock> holder { scheduler.m_lock };
if (scheduler.m_isShuttingDown) [[unlikely]]
return;
if (kind == DeferredWorkTimer::WorkType::ImminentlyScheduled) {
Bun__eventLoop__incrementRefConcurrently(clientData->bunVM, 1);
scheduler.m_pendingTicketsKeepingEventLoopAlive.add(WTF::move(ticket));
Expand All @@ -50,6 +68,24 @@ void JSCTaskScheduler::onAddPendingWork(WebCore::JSVMClientData* clientData, Ref
}
void JSCTaskScheduler::onScheduleWorkSoon(WebCore::JSVMClientData* clientData, Ticket ticket, Task&& task)
{
auto& scheduler = clientData->deferredWorkTimer;
Locker<Lock> holder { scheduler.m_lock };
// The event loop is past its last tick; a JSCDeferredWorkTask enqueued now
// would never run and its ConcurrentTask wrapper would leak once the Bun
// VirtualMachine box is dealloc'd. Reached from ~VM -> WaiterListManager::
// unregister -> Waiter::cancelAndClear for every outstanding
// Atomics.waitAsync on a terminating worker, and from collectNow ->
// JSFinalizationRegistry::finalizeUnconditionally. Balance onAddPendingWork
// so the ticket-set entry and event-loop ref are released. The lock is held
// across the check and the enqueue so the transition in markShuttingDown
// cannot race a cross-thread Atomics.notify.
if (scheduler.m_isShuttingDown) [[unlikely]] {
bool wasKeepingAlive = dropPendingTicketLocked(scheduler, ticket);
holder.unlockEarly();
if (wasKeepingAlive)
Bun__eventLoop__incrementRefConcurrently(clientData->bunVM, -1);
return;
}
auto* job = new JSCDeferredWorkTask(*ticket, WTF::move(task));
Bun__queueJSCDeferredWorkTaskConcurrently(clientData->bunVM, job);
}
Expand All @@ -60,19 +96,10 @@ void JSCTaskScheduler::onCancelPendingWork(WebCore::JSVMClientData* clientData,
auto& scheduler = clientData->deferredWorkTimer;

Locker<Lock> holder { scheduler.m_lock };
bool isKeepingEventLoopAlive = scheduler.m_pendingTicketsKeepingEventLoopAlive.removeIf([ticket](auto pendingTicket) {
return pendingTicket.ptr() == ticket;
});
// -- At this point, ticket may be an invalid pointer.

if (isKeepingEventLoopAlive) {
holder.unlockEarly();
bool wasKeepingAlive = dropPendingTicketLocked(scheduler, ticket);
holder.unlockEarly();
if (wasKeepingAlive)
Bun__eventLoop__incrementRefConcurrently(bunVM, -1);
} else {
scheduler.m_pendingTicketsOther.removeIf([ticket](auto pendingTicket) {
return pendingTicket.ptr() == ticket;
});
}
}

static void runPendingWork(void* bunVM, Bun::JSCTaskScheduler& scheduler, JSCDeferredWorkTask* job)
Expand Down Expand Up @@ -101,4 +128,30 @@ extern "C" void Bun__runDeferredWork(Bun::JSCDeferredWorkTask* job)
runPendingWork(clientData->bunVM, clientData->deferredWorkTimer, job);
}

// Flip m_isShuttingDown from the owning JS thread before the final concurrent-
// task drain. Any onScheduleWorkSoon that serializes before this under m_lock
// has its enqueue visible to the drain; any that serializes after drops.
extern "C" void Bun__JSCTaskScheduler__markShuttingDown(JSC::JSGlobalObject* globalObject)
{
if (auto* clientData = WebCore::clientData(JSC::getVM(globalObject)))
clientData->deferredWorkTimer.markShuttingDown();
}

// Reclaim a queued-but-never-dispatched job during shutdown. Called while the
// JSC VM is still alive, so ~Ref<TicketData> and the captured Task lambda may
// safely touch TZone-allocated / JSC-owned state. Mirrors runPendingWork's
// ticket take() so the pending set and event-loop ref stay balanced.
extern "C" void Bun__deleteDeferredWorkTask(Bun::JSCDeferredWorkTask* job)
{
if (auto* clientData = WebCore::clientData(job->vm())) {
auto& scheduler = clientData->deferredWorkTimer;
Locker<Lock> holder { scheduler.m_lock };
bool wasKeepingAlive = dropPendingTicketLocked(scheduler, job->ticket.ptr());
holder.unlockEarly();
if (wasKeepingAlive)
Bun__eventLoop__incrementRefConcurrently(clientData->bunVM, -1);
}
delete job;
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

}
13 changes: 13 additions & 0 deletions src/jsc/bindings/JSCTaskScheduler.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,21 @@ class JSCTaskScheduler {
static void onScheduleWorkSoon(WebCore::JSVMClientData* clientData, JSC::DeferredWorkTimer::Ticket ticket, JSC::DeferredWorkTimer::Task&& task);
static void onCancelPendingWork(WebCore::JSVMClientData* clientData, JSC::DeferredWorkTimer::Ticket ticket);

// Set once the owning VM's event loop has taken its last tick. After this,
// onScheduleWorkSoon drops the task instead of enqueueing a ConcurrentTask
// that can never be drained (~VM -> WaiterListManager::unregister reaches
// it for every still-pending Atomics.waitAsync ticket). Guarded by m_lock
// so the check+enqueue in onScheduleWorkSoon is atomic with respect to this
// transition (a cross-thread Atomics.notify may race a worker's shutdown).
void markShuttingDown()
{
Locker<Lock> holder { m_lock };
m_isShuttingDown = true;
}

public:
Lock m_lock;
bool m_isShuttingDown WTF_GUARDED_BY_LOCK(m_lock) { false };
UncheckedKeyHashSet<Ref<JSC::DeferredWorkTimer::TicketData>> m_pendingTicketsKeepingEventLoopAlive;
UncheckedKeyHashSet<Ref<JSC::DeferredWorkTimer::TicketData>> m_pendingTicketsOther;
};
Expand Down
2 changes: 2 additions & 0 deletions src/jsc/bindings/ZigGlobalObject.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -4024,6 +4024,8 @@ extern "C" void Zig__GlobalObject__destructOnExit(Zig::GlobalObject* globalObjec
// of enqueueing a ConcurrentTask that leaks past the last drain.
if (auto* ctx = globalObject->scriptExecutionContext())
ctx->markTerminating();
if (auto* clientData = WebCore::clientData(vm))
clientData->deferredWorkTimer.markShuttingDown();
Comment thread
robobun marked this conversation as resolved.
Bun__InspectorConnection__disconnectAllOnExit(globalObject);
// Hold a Ref so the RunLoop is guaranteed to outlive the VM teardown below.
Ref<WTF::RunLoop> runLoop = vm.runLoop();
Expand Down
5 changes: 5 additions & 0 deletions src/jsc/bindings/webcore/Worker.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -647,6 +647,11 @@ extern "C" void WebWorker__teardownJSCVM(Zig::GlobalObject* globalObject)
// can never run (e.g. notifyPeerClosed posted during the final collectNow).
if (auto* ctx = globalObject->scriptExecutionContext())
ctx->markTerminating();
// Same for DeferredWorkTimer: collectNow -> finalizers and ~VM ->
// WaiterListManager::unregister both reach scheduleWorkSoon; past this
// point those calls must not enqueue into our drained concurrent queue.
if (auto* clientData = WebCore::clientData(vm))
clientData->deferredWorkTimer.markShuttingDown();
Comment thread
robobun marked this conversation as resolved.

{
auto scope = DECLARE_THROW_SCOPE(vm);
Expand Down
8 changes: 8 additions & 0 deletions src/jsc/web_worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,9 @@ unsafe extern "C" {
// safe: opaque `&JSGlobalObject` handle (see above); takes the contexts-map
// lock and flips an atomic flag, no Rust-visible state touched.
safe fn ScriptExecutionContext__markTerminating(global: &JSGlobalObject);
// safe: same opaque-handle contract; flips JSCTaskScheduler::m_isShuttingDown
// under its own lock and returns. Idempotent.
safe fn Bun__JSCTaskScheduler__markShuttingDown(global: &JSGlobalObject);
// safe: `cpp_worker` is an opaque round-trip pointer owned by C++ (allocated
// there, stored in `WebWorker.cpp_worker`, and only ever passed back to C++
// — never dereferenced as Rust data); same contract as `JSC__VM__holdAPILock`'s
Expand Down Expand Up @@ -1291,6 +1294,11 @@ impl WebWorker {
// posted in the gap would sit in concurrent_tasks past the raw VM
// dealloc and leak under LSan.
ScriptExecutionContext__markTerminating(vm.global());
// Same for JSCTaskScheduler: a cross-thread Atomics.notify that
// races this shutdown either enqueues (and is caught by the drain)
// or observes m_isShuttingDown under m_lock and drops. Idempotent;
// teardownJSCVM sets it again.
Bun__JSCTaskScheduler__markShuttingDown(vm.global());
// Reclaim queued CppTasks (the per-worker stdio/messaging
// MessagePort drain tasks that can be in self.tasks mid-tick when
// terminate() lands, and any Worker dispatchExit close task from a
Expand Down
15 changes: 15 additions & 0 deletions src/runtime/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1198,6 +1198,21 @@ pub(crate) fn __bun_release_task_at_shutdown(task: bun_event_loop::Task) -> bool
for_each_fs_async_op!(__fs_destroy);
true
}
// A cross-thread Atomics.notify (or Wasm/FinalizationRegistry
// completion) enqueued this after the event loop's last tick. The
// dispatch arm above would have `delete`d it; mirror that here so the
// re-queue path doesn't keep it alive past worker VM dealloc. Runs
// before JSC teardown, so ~Ref<TicketData> is safe.
task_tag::JSCDeferredWorkTask => {
unsafe extern "C" {
fn Bun__deleteDeferredWorkTask(task: *mut JSCDeferredWorkTask);
}
// SAFETY: every JSCDeferredWorkTask payload is heap-allocated by
// `new JSCDeferredWorkTask` in JSCTaskScheduler::onScheduleWorkSoon;
// we own it once popped.
unsafe { Bun__deleteDeferredWorkTask(task.ptr.cast::<JSCDeferredWorkTask>()) };
true
}
// Same reclaim `drop_concurrent_cpp_tasks` performs, but for tasks
// that were already batch-moved into `self.tasks`. Must run before
// JSC teardown: a Worker `dispatchExit` lambda's `~Ref<Worker>` walks
Expand Down
32 changes: 32 additions & 0 deletions test/js/web/timers/timer-heap-atomics-teardown-fixture.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

25 changes: 23 additions & 2 deletions test/js/web/timers/timer-heap-race.test.ts
Original file line number Diff line number Diff line change
@@ -1,15 +1,16 @@
import { expect, it } from "bun:test";
import { bunEnv, bunExe, isDebug } from "harness";
import { bunEnv, bunExe, isASAN, isDebug } from "harness";
import path from "node:path";

async function runFixture(fixture: string) {
async function runFixture(fixture: string, env: Record<string, string | undefined> = {}) {
await using proc = Bun.spawn({
cmd: [bunExe(), path.join(import.meta.dir, fixture)],
env: {
...bunEnv,
// These make the debug build an order of magnitude slower; the fixtures need real wall time.
BUN_JSC_validateExceptionChecks: undefined,
BUN_JSC_dumpSimulatedThrows: undefined,
...env,
},
stdout: "pipe",
stderr: "pipe",
Expand Down Expand Up @@ -41,3 +42,23 @@ it.skipIf(!isDebug)(
},
20_000,
);

it.skipIf(!isASAN)(
"terminating a worker with pending Atomics.waitAsync tickets does not leak deferred-work tasks",
async () => {
const { stdout, stderr, signal, exitCode } = await runFixture("timer-heap-atomics-teardown-fixture.ts", {
BUN_DESTRUCT_VM_ON_EXIT: "1",
ASAN_OPTIONS: "allow_user_segv_handler=1:disable_coredump=0:detect_leaks=1:abort_on_error=1",
LSAN_OPTIONS: `malloc_context_size=30:print_suppressions=0:suppressions=${path.join(import.meta.dir, "..", "..", "..", "leaksan.supp")}`,
});
// LSan writes its leak report to stderr and SIGABRTs; stdout holds the
// fixture's own OK line either way, so assert exitCode/signal explicitly.
expect({ stdout, stderr, signal, exitCode }).toEqual({
stdout: "OK\n",
stderr: expect.not.stringContaining("LeakSanitizer"),
signal: null,
exitCode: 0,
});
Comment thread
robobun marked this conversation as resolved.
},
20_000,
);
Loading