Skip to content
Closed
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
172 changes: 96 additions & 76 deletions src/runtime/napi/napi_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2438,9 +2438,9 @@ extern "C" fn napi_internal_enqueue_finalizer(
// ThreadSafeFunction
// ──────────────────────────────────────────────────────────────────────────

/// Ownership: the JS thread owns this allocation while the env lives and frees
/// it in `destroy`; from `env_teardown_done` on it belongs to the remaining
/// `thread_count` references, and whoever drops the last one frees it.
/// Ownership: the JS thread owns this allocation until `finalize` (or
/// `env_teardown`) sets `resources_released`; then the remaining
/// `thread_count` references own it and the last one dropped frees it.
Comment thread
robobun marked this conversation as resolved.
// TODO: generate a compile-time version of this instead of runtime checking
pub(crate) struct ThreadSafeFunction {
/// thread-safe functions can be "referenced" and "unreferenced". A
Expand Down Expand Up @@ -2492,11 +2492,11 @@ pub(crate) struct ThreadSafeFunction {
/// that would reach `event_loop` from another thread reads it under the
/// same lock, so teardown cannot land between the check and the enqueue.
pub(crate) env_dead: AtomicBool,
/// Also written under `lock`, once `env_teardown` has released every
/// JS-thread-owned resource. Until then teardown still owns this object,
/// so a thread that drops the last `thread_count` reference must not free
/// it (Node's `kClosed`).
pub(crate) env_teardown_done: AtomicBool,
/// Also written under `lock`, once `finalize` or `env_teardown` has
/// released every JS-thread-owned resource. Until then the thread that
/// drops the last `thread_count` reference must not free this (Node's
/// `kClosed`).
Comment thread
robobun marked this conversation as resolved.
pub(crate) resources_released: AtomicBool,
}

pub(crate) enum TsfnCallback {
Expand Down Expand Up @@ -2593,9 +2593,9 @@ impl ThreadSafeFunction {
}
// SAFETY: as above.
if unsafe { (*this).closing.load(Ordering::SeqCst) } == ClosingState::Closed as u8 {
// Finalize the ThreadSafeFunction.
// SAFETY: `this` is the live heap allocation we own; closed state guarantees no other thread will touch it.
unsafe { ThreadSafeFunction::destroy(this) };
// SAFETY: `this` is the live heap allocation `maybe_queue_finalizer`
// queued this task for.
unsafe { ThreadSafeFunction::finalize(this) };
return;
}

Expand Down Expand Up @@ -2714,15 +2714,14 @@ impl ThreadSafeFunction {
// Closing (napi_tsfn_abort, or the last call already ran):
// nothing still queued runs any more, as in Node's DispatchOne.
// An abort's leftovers go back to the addon, with no lock held
// since that re-enters it; the function finalizes once the last
// thread reference is gone.
// since that re-enters it. Then the function finalizes whatever
// thread_count is: after an abort the other threads may never
// release (Node's CloseHandlesAndMaybeDelete).
Comment thread
robobun marked this conversation as resolved.
let leftovers = self.take_queue();
drop(_g);
self.hand_back(leftovers);
let _g = self.lock.lock_guard();
if self.thread_count.load(Ordering::SeqCst) == 0 {
self.maybe_queue_finalizer();
}
self.maybe_queue_finalizer();
return Ok(false);
}
let was_blocked = self.queue.is_blocked();
Expand Down Expand Up @@ -2865,8 +2864,8 @@ impl ThreadSafeFunction {
let (status, orphaned) = unsafe { (*this).enqueue(ctx, block) };

if orphaned {
// SAFETY: the lock is dropped, we dropped the last thread reference
// and `env_teardown` already released everything it owned.
// SAFETY: the lock is dropped, this was the last thread reference,
// and the JS thread already released what only it may release.
unsafe { ThreadSafeFunction::free_orphaned(this) };
}
status
Expand Down Expand Up @@ -2943,40 +2942,64 @@ impl ThreadSafeFunction {
}
}

/// Consumes and frees a heap-allocated ThreadSafeFunction (allocated by `new`).
/// SAFETY: `this` must be a live `*mut ThreadSafeFunction` returned from `heap::alloc`
/// and not aliased; caller transfers ownership.
pub(crate) unsafe fn destroy(this: *mut ThreadSafeFunction) {
// SAFETY: caller contract — `this` is a live heap allocation and we are
// the sole owner; reclaim the Box up front so the body works on owned
// state and the drop at scope end frees it.
let mut self_ = unsafe { bun_core::heap::take(this) };
self_.unref();

if let Some(env) = self_.env.as_ref() {
// SAFETY: env is live (we hold a ref); drops our registry entry so
// teardown cannot hand this pointer out after we free it. `this` is
// passed as an opaque registry key only, never dereferenced.
unsafe { NapiEnv__unregisterThreadSafeFunction(env.get(), this.cast()) };
/// Runs on the JS thread once the function has closed. Runs the addon's
/// finalizer, releases what only this thread may release, and frees the
/// allocation unless another thread still holds a reference; then the last
/// release frees it (Node's Finalize + MaybeDelete).
///
/// SAFETY: `this` is a live allocation from `new` that no other event-loop
/// task will reach again.
unsafe fn finalize(this: *mut ThreadSafeFunction) {
// SAFETY: caller contract. The borrow ends before the addon's finalizer
// runs; it may still use the handle (napi_get_threadsafe_function_context).
let finalizer = unsafe {
let self_ = &mut *this;
if let Some(env) = self_.env.as_ref() {
// SAFETY: env is live (we hold a ref). `this` is an opaque
// registry key here, never dereferenced.
NapiEnv__unregisterThreadSafeFunction(env.get(), this.cast());
}
self_
.finalizer_fun
.take()
.zip(self_.env.as_ref())
.map(|(fun, env)| Finalizer {
env: env.clone(),
fun,
data: self_.finalizer_data,
hint: self_.ctx,
})
};

// Before anything is released, as in Node and `env_teardown`.
if let Some(mut finalizer) = finalizer {
crate::dispatch::fold(finalizer.run());
}

if let (Some(fun), Some(env)) = (self_.finalizer_fun, self_.env.as_ref()) {
// Note: ownership transfer of `env` into the Finalizer. We clone (bumps the
// external refcount) and let the original drop with the Box below — net refcount
// delta is zero.
let finalizer = Finalizer {
env: env.clone(),
fun,
data: self_.finalizer_data,
hint: self_.ctx,
};
finalizer.enqueue();
// SAFETY: caller contract; the finalizer has returned and the borrow
// ends before the free below.
let free = unsafe {
let self_ = &mut *this;
// The same critical section reads thread_count, so a thread that
// drops the last reference frees only if it sees this store.
Comment thread
robobun marked this conversation as resolved.
let _g = self_.lock.lock_guard();
self_.event_loop = None;
drop(self_.env.take());
self_.resources_released.store(true, Ordering::SeqCst);
self_.thread_count.load(Ordering::SeqCst) <= 0
};

if free {
// SAFETY: no thread reference is left, the lock is dropped, and
// nothing else can reach this allocation.
unsafe { ThreadSafeFunction::free_orphaned(this) };
}
}

/// Frees the allocation and nothing else: no finalizer, no registry entry,
/// no event loop. Every JS-thread-owned resource must already be released
/// (`env_teardown`) or be safe to drop here (a creation that failed).
/// (`finalize`, `env_teardown`) or be safe to drop here (a creation that
/// failed).
Comment thread
robobun marked this conversation as resolved.
///
/// SAFETY: `this` is a live allocation from `new`, the caller holds no
/// lock on it, and no other thread holds a reference.
Expand Down Expand Up @@ -3056,15 +3079,14 @@ impl ThreadSafeFunction {
}

// Phase 3: release what only the JS thread may release, then hand the
// allocation over: `env_teardown_done` is what lets another thread free
// it, so it is published in the same critical section that reads
// thread_count (Node's ReleaseResources + MaybeDelete).
// allocation over. `resources_released` is published in the critical
// section that reads thread_count (Node's ReleaseResources + MaybeDelete).
Comment thread
robobun marked this conversation as resolved.
let _g = self.lock.lock_guard();
self.callback = TsfnCallback::Js(StrongOptional::empty());
self.poll_ref.disable();
self.event_loop = None;
drop(self.env.take());
self.env_teardown_done.store(true, Ordering::SeqCst);
self.resources_released.store(true, Ordering::SeqCst);
// Cleanup hooks are the loop's last tick: a task still queued for this
// TSFN will never run (and its `release_unrun` does not dereference it).
// With no thread_count reference left, nobody else can reach this, so
Expand Down Expand Up @@ -3111,8 +3133,8 @@ impl ThreadSafeFunction {
};

if orphaned {
// SAFETY: the lock is dropped, we dropped the last thread reference
// and `env_teardown` already released everything it owned.
// SAFETY: the lock is dropped, this was the last thread reference,
// and the JS thread already released what only it may release.
unsafe { ThreadSafeFunction::free_orphaned(this) };
}
status
Expand All @@ -3130,34 +3152,32 @@ impl ThreadSafeFunction {

let prev_remaining = self.thread_count.fetch_sub(1, Ordering::SeqCst);

if self.env_dead.load(Ordering::SeqCst) {
// The event loop we were created on is gone (`env_teardown` set
// this under the lock we hold). Never schedule onto it. Whoever
// drops the last reference frees us -- but only once teardown has
// released the JS-thread-owned resources; until then it owns us
// and will free us itself if we are the last to let go.
let orphaned = prev_remaining == 1 && self.env_teardown_done.load(Ordering::SeqCst);
if self.env_dead.load(Ordering::SeqCst)
|| self.closing.load(Ordering::SeqCst) == ClosingState::Closed as u8
{
// The JS thread owns the finalization (`finalize` is queued or
// done, or `env_teardown` ran): never schedule onto the loop again.
// The last reference frees us once the JS thread has released what
// only it may release; before that, it frees us itself.
Comment thread
robobun marked this conversation as resolved.
let orphaned = prev_remaining == 1 && self.resources_released.load(Ordering::SeqCst);
return (NapiStatus::ok as napi_status, orphaned);
}

if mode == napi_threadsafe_function_release_mode::abort || prev_remaining == 1 {
if !self.is_closing() {
if mode == napi_threadsafe_function_release_mode::abort {
self.closing
.store(ClosingState::Closing as u8, Ordering::SeqCst);
if self.queue.max_queue_size > 0 {
// Wake all producers blocked in enqueue()'s bounded
// queue wait so they observe is_closing and release.
self.blocking_condvar.broadcast();
}
// Already closing: the abort's dispatch is pending or running and
// finalizes whatever thread_count is by then.
Comment thread
robobun marked this conversation as resolved.
if (mode == napi_threadsafe_function_release_mode::abort || prev_remaining == 1)
&& !self.is_closing()
{
if mode == napi_threadsafe_function_release_mode::abort {
self.closing
.store(ClosingState::Closing as u8, Ordering::SeqCst);
if self.queue.max_queue_size > 0 {
// Wake all producers blocked in enqueue()'s bounded
// queue wait so they observe is_closing and release.
Comment thread
robobun marked this conversation as resolved.
self.blocking_condvar.broadcast();
}
self.schedule_dispatch();
} else if prev_remaining == 1 {
// Already closing from an earlier abort. The last release must
// still reach dispatch_one's thread_count==0 path so the
// finalizer runs and the event-loop keepalive is dropped.
self.schedule_dispatch();
}
self.schedule_dispatch();
}

(NapiStatus::ok as napi_status, false)
Expand All @@ -3169,7 +3189,7 @@ impl ThreadSafeFunction {
#[unsafe(no_mangle)]
extern "C" fn napi_internal_threadsafe_function_env_teardown(tsfn: *mut c_void) {
let this = tsfn.cast::<ThreadSafeFunction>();
// SAFETY: the registry only holds live TSFN pointers — `destroy` and
// SAFETY: the registry only holds live TSFN pointers — `finalize` and
// `env_teardown` both remove the entry before freeing. Exclusive borrow
// scoped to this call.
if unsafe { (*this).env_teardown() } {
Expand Down Expand Up @@ -3245,7 +3265,7 @@ extern "C" fn napi_create_threadsafe_function(
blocking_condvar: Condvar::default(),
closing: AtomicU8::new(ClosingState::NotClosing as u8),
env_dead: AtomicBool::new(false),
env_teardown_done: AtomicBool::new(false),
resources_released: AtomicBool::new(false),
});

// Register with the env so that VM/worker teardown neutralizes this TSFN
Expand Down
29 changes: 29 additions & 0 deletions test/napi/napi-app/module.js
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,35 @@ nativeTests.test_threadsafe_function_abort_blocked_producers = async () => {
console.log("finalized:", nativeTests.test_napi_threadsafe_function_abort_blocked_producers_finalized());
};

nativeTests.test_threadsafe_function_abort_with_outstanding_ref = async () => {
// create (thread_count=2) with the second reference held by a thread that
// makes no calls, then abort from this thread
nativeTests.test_napi_threadsafe_function_abort_with_outstanding_ref();
// the finalizer runs from the abort's dispatch on this thread, so a bounded
// number of turns is enough; it must not wait for the other thread
for (let i = 0; i < 1000; i++) {
if (nativeTests.test_napi_threadsafe_function_abort_with_outstanding_ref_finalized()) break;
await new Promise(resolve => setImmediate(resolve));
}
console.log("finalized:", nativeTests.test_napi_threadsafe_function_abort_with_outstanding_ref_finalized());
// the other thread releases once it has seen the finalizer run (or gives up
// after 2s and says so); that is the one step here that depends on OS
// scheduling, so wait for it by deadline, under the test runner's 5s
const deadline = Date.now() + 4_000;
while (
nativeTests.test_napi_threadsafe_function_abort_with_outstanding_ref_release_status() === -1 &&
Date.now() < deadline
) {
await new Promise(resolve => setTimeout(resolve, 1));
}
console.log(
"released after finalize:",
nativeTests.test_napi_threadsafe_function_abort_with_outstanding_ref_released_after_finalize(),
"status:",
nativeTests.test_napi_threadsafe_function_abort_with_outstanding_ref_release_status(),
);
};

nativeTests.test_get_exception = (_, value) => {
function thrower() {
throw value;
Expand Down
Loading
Loading