diff --git a/changelog.d/11445-cross-agent-pump-ownership.md b/changelog.d/11445-cross-agent-pump-ownership.md new file mode 100644 index 0000000000..f14351c562 --- /dev/null +++ b/changelog.d/11445-cross-agent-pump-ownership.md @@ -0,0 +1,7 @@ +Fixed `test_gap_turnloop_p9_worker_agent_net` failing intermittently (#11433). A `node:worker_threads` Worker doing `fetch` against a server owned by the main thread printed `immediate 000undefined`, threw `Invalid response handle`, or stalled at `WORKER STOPPED AFTER n/2`, because every JS agent runs the same process-wide pumps and three queues had no owner tag, so whichever agent pumped first drained all of them: + +- perry-stdlib's promise-resolution queues (`common/async_bridge.rs`): a Worker's fetch promise could be settled on the main thread, and vice versa. +- The `worker_threads` Worker-to-parent event queue: a Worker could drain its own `postMessage` events, so the parent never received them. +- perry-ext-http's server pump: a Worker could run the main thread's `http.Server` handler, so the request was never answered. + +Each entry now records its `AgentId`. Pumps, keep-alive checks and the async-bridge GC scanner act only on the calling agent's entries. Pool and fallback threads re-assert the submitting agent with `ResolutionOwnerScope`. `perry_runtime::agent::register_retire_hook` purges a retired agent's leftovers, and `perry_ffi::agent_post::current_agent()` exposes the agent id to ext crates. Under CPU load, the test went from 61/100 failures before the fix to 0/100 after. diff --git a/crates/perry-ext-http/src/server/mod.rs b/crates/perry-ext-http/src/server/mod.rs index ee0fe5e4b4..9fbf04ab17 100644 --- a/crates/perry-ext-http/src/server/mod.rs +++ b/crates/perry-ext-http/src/server/mod.rs @@ -277,6 +277,30 @@ mod tests { s } + /// #11433 — a server belongs to the agent that created it, and no other + /// agent's pump may treat it as its own. Every agent runs this extension's + /// pump, so this is what keeps a Worker from dispatching the primary + /// agent's requests on the Worker's thread. + #[test] + fn a_server_belongs_to_the_agent_that_created_it() { + let here = HttpServer::with_handler(0); + assert!(here.owned_here()); + let worker_owner = std::thread::spawn(|| { + perry_runtime::agent::enter_worker_agent(); + let s = HttpServer::with_handler(0); + assert!(s.owned_here(), "the creating Worker owns its server"); + s.owner_agent + }) + .join() + .unwrap(); + let mut foreign = HttpServer::with_handler(0); + foreign.owner_agent = worker_owner; + assert!( + !foreign.owned_here(), + "the primary agent must not drain a Worker's server" + ); + } + /// Issue #2210 — `HttpServer::with_handler` seeds Node's /// documented timeout defaults so a fresh server reads back the /// same numbers Node returns when no options are passed. diff --git a/crates/perry-ext-http/src/server/server.rs b/crates/perry-ext-http/src/server/server.rs index fc5a7ee79d..18ecea28c7 100644 --- a/crates/perry-ext-http/src/server/server.rs +++ b/crates/perry-ext-http/src/server/server.rs @@ -134,6 +134,21 @@ pub struct HttpServer { pub bun_error_handler: i64, /// Observable `Server.development` option. pub bun_development: bool, + /// The JS agent that created this server (#11433). Its handler and + /// listeners are closures in that agent's heap, and its connections live on + /// that agent's loop, so only that agent's pump may drain it. Every agent's + /// event loop runs this extension's pump — a `node:worker_threads` Worker's + /// included — and before this tag a Worker's pump would take a request + /// addressed to the primary agent's server and run the primary's handler on + /// the Worker's thread. + pub owner_agent: u64, +} + +impl HttpServer { + /// Whether the calling agent owns this server and may run its JS. + pub(crate) fn owned_here(&self) -> bool { + self.owner_agent == perry_ffi::agent_post::current_agent() + } } impl HttpServer { @@ -170,6 +185,7 @@ impl HttpServer { is_bun_server: false, bun_error_handler: 0, bun_development: false, + owner_agent: perry_ffi::agent_post::current_agent(), } } } @@ -1025,8 +1041,9 @@ pub unsafe extern "C" fn js_node_http_server_remove_listener( #[no_mangle] pub extern "C" fn js_node_http_server_has_active() -> i32 { let mut active = 0i32; + // Only this agent's servers keep this agent alive (#11433). iter_handles_of::(|s| { - if server_is_active(s) { + if s.owned_here() && server_is_active(s) { active = 1; } }); @@ -1034,7 +1051,7 @@ pub extern "C" fn js_node_http_server_has_active() -> i32 { return 1; } iter_handles_of::(|s| { - if server_is_active(&s.base) { + if s.base.owned_here() && server_is_active(&s.base) { active = 1; } }); @@ -1042,7 +1059,7 @@ pub extern "C" fn js_node_http_server_has_active() -> i32 { return 1; } iter_handles_of::(|s| { - if server_is_active(&s.base) { + if s.base.owned_here() && server_is_active(&s.base) { active = 1; } }); @@ -1057,6 +1074,19 @@ pub extern "C" fn js_node_http_server_has_active() -> i32 { active } +/// Whether `server_handle` names an HTTP, HTTPS or HTTP/2 server the calling +/// agent owns. A handle that names none of them reads as not owned. +pub(crate) fn server_handle_owned_here(server_handle: i64) -> bool { + if let Some(s) = get_handle::(server_handle) { + return s.owned_here(); + } + if let Some(s) = get_handle::(server_handle) { + return s.base.owned_here(); + } + get_handle::(server_handle) + .is_some_and(|s| s.base.owned_here()) +} + /// Drain pending requests + upgrades from every registered server, /// dispatching to the user handler / `'upgrade'` listener on the /// main thread. Called each tick by perry-stdlib's pump (gated on @@ -1099,9 +1129,17 @@ pub extern "C" fn js_node_http_server_process_pending() -> i32 { // #4905 — fire `'connection'` listeners for connections accepted since // the last tick, before their requests are dispatched (Node fires // `'connection'` ahead of `'request'`). + // Only the events of servers this agent owns (#11433); another agent's stay + // queued for its own pump. let connection_events: Vec<(i64, i64)> = PENDING_CONNECTION_EVENTS .lock() - .map(|mut q| q.drain(..).collect()) + .map(|mut q| { + let (mine, theirs): (Vec<_>, Vec<_>) = std::mem::take(&mut *q) + .into_iter() + .partition(|(server_handle, _)| server_handle_owned_here(*server_handle)); + *q = theirs; + mine + }) .unwrap_or_default(); for (server_handle, socket_handle) in connection_events { // The handle may back an HttpServer or an HttpsServer (whose @@ -1140,8 +1178,10 @@ pub extern "C" fn js_node_http_server_process_pending() -> i32 { // Snapshot handle ids first so we can mutate handle state // (drain channels, free per-request handles) without the // DashMap iterator dangling. + // #11433: each loop below visits only the servers this agent created. let mut http_handles: Vec = Vec::new(); perry_ffi::iter_handle_ids_of::(|id| http_handles.push(id)); + http_handles.retain(|h| get_handle::(*h).is_some_and(HttpServer::owned_here)); for h in http_handles { // #4903 — fire the deferred `'listening'` emit + listen callbacks // before draining requests: the listen callback is usually what @@ -1161,6 +1201,10 @@ pub extern "C" fn js_node_http_server_process_pending() -> i32 { perry_ffi::iter_handle_ids_of::(|id| { https_handles.push(id) }); + https_handles.retain(|h| { + get_handle::(*h) + .is_some_and(|s| s.base.owned_here()) + }); for h in https_handles { count += drain_deferred_listen_for::(h, |s| { &mut s.base @@ -1186,6 +1230,10 @@ pub extern "C" fn js_node_http_server_process_pending() -> i32 { perry_ffi::iter_handle_ids_of::(|id| { h2_handles.push(id) }); + h2_handles.retain(|h| { + get_handle::(*h) + .is_some_and(|s| s.base.owned_here()) + }); for h in h2_handles { count += drain_deferred_listen_for::( h, diff --git a/crates/perry-ext-http/src/server/server/in_flight.rs b/crates/perry-ext-http/src/server/server/in_flight.rs index ded90a3c5c..51c92af661 100644 --- a/crates/perry-ext-http/src/server/server/in_flight.rs +++ b/crates/perry-ext-http/src/server/server/in_flight.rs @@ -45,6 +45,9 @@ pub(crate) struct InFlightRequest { /// Grace deadline. Past this, synthesize the default response (unless /// `skip_default_response`) and free the handles regardless. pub(crate) deadline: Instant, + /// The agent whose pump parked it — the server's owner (#11433). Only that + /// agent reaps it, fires its `'drain'` listeners, or is kept alive by it. + owner_agent: u64, } pub(crate) static IN_FLIGHT: Mutex> = Mutex::new(Vec::new()); @@ -90,7 +93,11 @@ pub(crate) fn finalize_request_handles_deferred( /// server's handle "active" so the main loop doesn't exit before the /// pending response is flushed. pub(crate) fn has_in_flight_requests() -> bool { - IN_FLIGHT.lock().map(|g| !g.is_empty()).unwrap_or(false) + let agent = perry_ffi::agent_post::current_agent(); + IN_FLIGHT + .lock() + .map(|g| g.iter().any(|e| e.owner_agent == agent)) + .unwrap_or(false) } /// Finalize parked requests whose handler has now called `res.end()` (the @@ -110,7 +117,13 @@ pub(crate) fn reap_in_flight_requests() { return; } let now = Instant::now(); + let agent = perry_ffi::agent_post::current_agent(); guard.retain(|e| { + if e.owner_agent != agent { + // Another agent's request: its handles and listeners live in + // that agent's heap (#11433). + return true; + } let ended = response_writable_ended(e.response_handle); if !ended { // Streaming backpressure cleared — fire `'drain'` (outside @@ -193,6 +206,7 @@ pub(crate) fn finalize_or_park_request(pending: &HttpPendingRequest) { response_handle: pending.response_handle, skip_default_response: pending.skip_default_response, deadline, + owner_agent: perry_ffi::agent_post::current_agent(), }); } else { // Lock poisoned — fall back to the old immediate behavior so we diff --git a/crates/perry-ext-http/src/server/turnloop_serve/conn.rs b/crates/perry-ext-http/src/server/turnloop_serve/conn.rs index 645e2f0e9d..a88a60c56b 100644 --- a/crates/perry-ext-http/src/server/turnloop_serve/conn.rs +++ b/crates/perry-ext-http/src/server/turnloop_serve/conn.rs @@ -125,8 +125,12 @@ fn pending() -> &'static Mutex>> { /// `IncomingMessage` handles whose connection died before their response /// completed. Node raises `'aborted'` on the request; the sink cannot run JS, /// so the pump drains this and fires the listeners on its own tick. -fn aborted() -> &'static Mutex> { - static ABORTED: OnceLock>> = OnceLock::new(); +/// +/// Each entry is tagged with the agent whose loop queued it — the server's +/// owner, since the completion sink runs there — and only that agent's pump +/// takes it (#11433). +fn aborted() -> &'static Mutex> { + static ABORTED: OnceLock>> = OnceLock::new(); ABORTED.get_or_init(|| Mutex::new(Vec::new())) } @@ -142,21 +146,32 @@ pub(crate) fn note_aborted_handle(handle: i64) { aborted() .lock() .unwrap_or_else(|e| e.into_inner()) - .push(handle); + .push((perry_ffi::agent_post::current_agent(), handle)); } /// Take the `IncomingMessage` handles whose connection died mid-request. pub(crate) fn take_aborted() -> Vec { let mut queue = aborted().lock().unwrap_or_else(|e| e.into_inner()); - std::mem::take(&mut *queue) + take_owned(&mut queue) +} + +/// Remove and return the calling agent's entries, leaving every other agent's +/// in place for its own pump (#11433). +fn take_owned(queue: &mut Vec<(u64, i64)>) -> Vec { + let agent = perry_ffi::agent_post::current_agent(); + let (mine, theirs): (Vec<_>, Vec<_>) = std::mem::take(queue) + .into_iter() + .partition(|(owner, _)| *owner == agent); + *queue = theirs; + mine.into_iter().map(|(_, handle)| handle).collect() } /// Connection-socket handles (`alloc_connection_socket`) whose TCP connection /// has fully closed. Same pattern as `aborted()`: the completion sink cannot /// run JS, so the pump drains this and fires the socket's `'close'` /// listeners on its own tick. -fn closed_sockets() -> &'static Mutex> { - static CLOSED: OnceLock>> = OnceLock::new(); +fn closed_sockets() -> &'static Mutex> { + static CLOSED: OnceLock>> = OnceLock::new(); CLOSED.get_or_init(|| Mutex::new(Vec::new())) } @@ -164,13 +179,13 @@ fn note_closed_socket(socket_handle: i64) { closed_sockets() .lock() .unwrap_or_else(|e| e.into_inner()) - .push(socket_handle); + .push((perry_ffi::agent_post::current_agent(), socket_handle)); } /// Take the connection-socket handles due a `'close'` emit. pub(crate) fn take_closed_sockets() -> Vec { let mut queue = closed_sockets().lock().unwrap_or_else(|e| e.into_inner()); - std::mem::take(&mut *queue) + take_owned(&mut queue) } /// Note that this connection's in-flight request (if any) will never be @@ -184,10 +199,7 @@ fn note_aborted(id: i64) { .flatten() .filter(|h| *h != 0); if let Some(handle) = handle { - aborted() - .lock() - .unwrap_or_else(|e| e.into_inner()) - .push(handle); + note_aborted_handle(handle); } } diff --git a/crates/perry-runtime/src/agent.rs b/crates/perry-runtime/src/agent.rs index 291bb0c2d6..a90f15687e 100644 --- a/crates/perry-runtime/src/agent.rs +++ b/crates/perry-runtime/src/agent.rs @@ -54,6 +54,7 @@ use std::cell::Cell; use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Mutex, PoisonError}; /// Identity of a JS heap (arena + GC) that queued cross-thread work. pub type AgentId = u64; @@ -111,6 +112,26 @@ pub fn owns(owner: AgentId) -> bool { owner == current_agent() } +/// Purge hooks for agent-tagged queues that live OUTSIDE this crate (#11433). +/// +/// `perry-stdlib`'s promise-resolution queues (`common::async_bridge`) are +/// tagged with the enqueuing agent exactly like the timer and thread-result +/// queues below, and need the same purge when that agent dies — but they cannot +/// be named from here. A hook registers once per process and runs inside +/// [`retire_agent`] with the dying agent's id. +static RETIRE_HOOKS: Mutex> = Mutex::new(Vec::new()); + +/// Register `hook` to run from [`retire_agent`]. Idempotent per function. +pub fn register_retire_hook(hook: fn(AgentId)) { + let mut hooks = RETIRE_HOOKS.lock().unwrap_or_else(PoisonError::into_inner); + if !hooks + .iter() + .any(|existing| std::ptr::fn_addr_eq(*existing, hook)) + { + hooks.push(hook); + } +} + /// Retire a worker agent at thread exit: its arena is about to be unmapped, so /// any queue entry still tagged with it points at memory that is about to go /// away. Purge those entries rather than leaving them for a drain that can @@ -141,6 +162,14 @@ pub fn retire_agent(id: AgentId) { crate::event_pump::shutdown_agent_loop(); crate::timer::purge_agent_timers(id); crate::thread::purge_agent_thread_results(id); + // Copied out so a hook may itself take locks without holding this one. + let hooks: Vec = RETIRE_HOOKS + .lock() + .unwrap_or_else(PoisonError::into_inner) + .clone(); + for hook in hooks { + hook(id); + } // Deliberately do NOT clear `CURRENT_AGENT`. Clearing it would make // `current_agent()` fall back to `PRIMARY_AGENT` for the rest of this // thread's life — i.e. a worker that has just torn down its heap would diff --git a/crates/perry-stdlib/src/common/async_bridge.rs b/crates/perry-stdlib/src/common/async_bridge.rs index 8e42f7ed12..6fce726b0a 100644 --- a/crates/perry-stdlib/src/common/async_bridge.rs +++ b/crates/perry-stdlib/src/common/async_bridge.rs @@ -15,6 +15,28 @@ //! 2. Store raw Rust data and use deferred conversion callbacks //! 3. The conversion callbacks run on the main thread during js_stdlib_process_pending //! +//! # One queue, many agents (#11433) +//! +//! "The main thread" above means *the agent that owns the promise*. Since +//! turnloop P9 a `node:worker_threads` Worker (and a `perry/thread` worker) is +//! an agent with its own heap, its own loop, and its own call to +//! `js_stdlib_process_pending` from its await loop. The two queues below are +//! process-global, so every entry carries the [`AgentId`] it belongs to and a +//! pump settles only its own agent's entries, leaving the rest for their +//! owner — the model `perry_runtime::agent` already applies to timers and +//! thread results. Before this, whichever agent pumped first took everything: +//! the primary agent settled a Worker's `fetch` promise (running the Worker's +//! continuation on the main thread, against the main heap) and the Worker +//! settled the primary agent's, which is how +//! `test_gap_turnloop_p9_worker_agent_net` printed `immediate 000undefined`, +//! lost messages, and threw `Invalid response handle` about half the time. +//! +//! The owner is the agent current at *enqueue* time. A producer that queues +//! from a thread with no agent of its own — a pool job, a fallback OS thread — +//! would read as the primary agent there, so the spawn helpers capture the +//! submitting agent and re-assert it around the job with +//! [`ResolutionOwnerScope`]. +//! //! # No tokio //! //! This is the only async bridge perry-stdlib has. Turnloop P8 lane L split @@ -30,6 +52,8 @@ use std::sync::Mutex; use std::sync::LazyLock as Lazy; +use perry_runtime::agent::{current_agent, AgentId}; + /// Issue #859: pin a Promise so the GC can't sweep it while a pool /// job or native thread is computing its eventual resolution. /// @@ -178,6 +202,71 @@ static PENDING_DEFERRED_LEN: AtomicUsize = AtomicUsize::new(0); thread_local! { static GC_SCANNER_REGISTERED: std::cell::Cell = const { std::cell::Cell::new(false) }; + /// The agent a job running on this (agent-less) thread settles promises + /// for. Set only by [`ResolutionOwnerScope`]. + static RESOLUTION_OWNER: std::cell::Cell> = const { std::cell::Cell::new(None) }; +} + +/// The agent a resolution queued from this thread belongs to: the +/// [`ResolutionOwnerScope`] in force, or the calling agent. Also what a spawn +/// helper captures as the owner of the job it is about to hand off. +pub(crate) fn resolution_owner() -> AgentId { + RESOLUTION_OWNER + .with(std::cell::Cell::get) + .unwrap_or_else(current_agent) +} + +/// While alive, resolutions queued on this thread belong to `agent`. +/// +/// For native work that runs OFF the agent that created its promise (a +/// turnloop pool job, a fallback OS thread): capture [`current_agent`] where +/// the promise is created, and enter the scope around the job. Without it the +/// job's thread reads as the primary agent, and a Worker's promise would be +/// settled by the main thread. +pub(crate) struct ResolutionOwnerScope { + previous: Option, +} + +impl ResolutionOwnerScope { + pub(crate) fn enter(agent: AgentId) -> Self { + let previous = RESOLUTION_OWNER.with(|slot| slot.replace(Some(agent))); + Self { previous } + } +} + +impl Drop for ResolutionOwnerScope { + fn drop(&mut self) { + RESOLUTION_OWNER.with(|slot| slot.set(self.previous)); + } +} + +/// `perry_runtime::agent::retire_agent` hook: the agent's heap is going away, +/// so nothing can ever settle its queued promises. Drop them (their promise +/// addresses would dangle), and republish the lengths so the survivors' +/// keep-alive predicate stops counting them. +fn purge_agent_resolutions(agent: AgentId) { + { + let mut pending = PENDING_RESOLUTIONS.lock().unwrap(); + pending.retain(|resolution| resolution.owner != agent); + PENDING_RESOLUTIONS_LEN.store(pending.len(), Ordering::Release); + } + // Converters are dropped outside the lock: a boxed closure's captures may + // run arbitrary Drop code. + let dropped: Vec = { + let mut pending = PENDING_DEFERRED.lock().unwrap(); + let (dead, live): (Vec<_>, Vec<_>) = std::mem::take(&mut *pending) + .into_iter() + .partition(|resolution| resolution.owner == agent); + *pending = live; + PENDING_DEFERRED_LEN.store(pending.len(), Ordering::Release); + dead + }; + drop(dropped); +} + +fn ensure_retire_hook_registered() { + static REGISTER: std::sync::Once = std::sync::Once::new(); + REGISTER.call_once(|| perry_runtime::agent::register_retire_hook(purge_agent_resolutions)); } pub(crate) fn ensure_gc_scanner_registered() { @@ -189,12 +278,15 @@ pub(crate) fn ensure_gc_scanner_registered() { "stdlib:async_bridge", scan_pending_native_async_resolution_roots_mut, ); + ensure_retire_hook_registered(); registered.set(true); }); } /// A pending promise resolution (for simple values that don't need conversion) struct PendingResolution { + /// The agent whose heap `promise_ptr` (and any heap `result_bits`) is in. + owner: AgentId, /// Pointer to the Promise object (as usize for Send) promise_ptr: usize, /// True if resolved successfully, false if rejected @@ -206,6 +298,8 @@ struct PendingResolution { /// A deferred promise resolution with a conversion callback /// The converter function runs on the main thread to safely create JSValues struct DeferredResolution { + /// The agent whose heap `promise_ptr` is in, and on which `converter` runs. + owner: AgentId, /// Pointer to the Promise object (as usize for Send) promise_ptr: usize, /// True if resolved successfully, false if rejected @@ -216,21 +310,24 @@ struct DeferredResolution { } /// Mutable GC scanner for native async completions waiting in stdlib's -/// main-thread pump. Promise pointers are raw heap pointers; simple -/// result bits may be NaN-boxed heap values. +/// pump. Promise pointers are raw heap pointers; simple result bits may be +/// NaN-boxed heap values. A collection runs on the agent whose heap it +/// collects, so it visits only that agent's entries: another agent's promise +/// is in another heap, which this collector must neither mark nor rewrite. pub fn scan_pending_native_async_resolution_roots_mut( visitor: &mut perry_runtime::gc::RuntimeRootVisitor<'_>, ) { + let agent = current_agent(); { let mut pending = PENDING_RESOLUTIONS.lock().unwrap(); - for resolution in pending.iter_mut() { + for resolution in pending.iter_mut().filter(|r| r.owner == agent) { visitor.visit_usize_slot(&mut resolution.promise_ptr); visitor.visit_nanbox_u64_slot(&mut resolution.result_bits); } } { let mut pending = PENDING_DEFERRED.lock().unwrap(); - for resolution in pending.iter_mut() { + for resolution in pending.iter_mut().filter(|r| r.owner == agent) { visitor.visit_usize_slot(&mut resolution.promise_ptr); } } @@ -254,6 +351,7 @@ pub fn queue_promise_resolution(promise_ptr: usize, is_success: bool, result_bit { let mut pending = PENDING_RESOLUTIONS.lock().unwrap(); pending.push(PendingResolution { + owner: resolution_owner(), promise_ptr, is_success, result_bits, @@ -283,6 +381,7 @@ where { let mut pending = PENDING_DEFERRED.lock().unwrap(); pending.push(DeferredResolution { + owner: resolution_owner(), promise_ptr, is_success, converter: Box::new(converter), @@ -340,13 +439,20 @@ pub fn ensure_pump_registered() { pub extern "C" fn js_stdlib_process_pending() -> i32 { let mut count = 0i32; + // Only this agent's entries (#11433): another agent's promise lives in + // another heap and must be settled — and its converter run — there. + let agent = current_agent(); + // Process simple resolutions first let simple_resolutions: Vec = { let mut pending = PENDING_RESOLUTIONS.lock().unwrap(); - let n = pending.len(); - count += n as i32; - PENDING_RESOLUTIONS_LEN.store(0, Ordering::Release); - pending.drain(..).collect() + let (mine, theirs): (Vec<_>, Vec<_>) = std::mem::take(&mut *pending) + .into_iter() + .partition(|resolution| resolution.owner == agent); + *pending = theirs; + count += mine.len() as i32; + PENDING_RESOLUTIONS_LEN.store(pending.len(), Ordering::Release); + mine }; for resolution in simple_resolutions { let scope = perry_runtime::gc::RuntimeHandleScope::new(); @@ -380,10 +486,13 @@ pub extern "C" fn js_stdlib_process_pending() -> i32 { // Process deferred resolutions - these run converter functions on the main thread let deferred_resolutions: Vec = { let mut pending = PENDING_DEFERRED.lock().unwrap(); - let n = pending.len(); - count += n as i32; - PENDING_DEFERRED_LEN.store(0, Ordering::Release); - pending.drain(..).collect() + let (mine, theirs): (Vec<_>, Vec<_>) = std::mem::take(&mut *pending) + .into_iter() + .partition(|resolution| resolution.owner == agent); + *pending = theirs; + count += mine.len() as i32; + PENDING_DEFERRED_LEN.store(pending.len(), Ordering::Release); + mine }; for resolution in deferred_resolutions { @@ -461,6 +570,24 @@ pub extern "C" fn js_stdlib_process_pending() -> i32 { count } +/// Whether a queued resolution belongs to the calling agent. Asked only when +/// the length mirrors are non-zero. Another agent's entry keeps THAT agent +/// alive; counting it here would hold this one open (and spin its wait, since +/// it can never drain the entry) until the owner gets round to it. +fn has_own_pending_resolution() -> bool { + let agent = current_agent(); + PENDING_RESOLUTIONS + .lock() + .unwrap() + .iter() + .any(|resolution| resolution.owner == agent) + || PENDING_DEFERRED + .lock() + .unwrap() + .iter() + .any(|resolution| resolution.owner == agent) +} + /// Returns 1 if the stdlib has active event sources that need the event /// loop to keep running (active WS servers, pending events, etc.). /// Registered with perry-runtime via js_register_stdlib_has_active() @@ -475,8 +602,9 @@ pub extern "C" fn js_stdlib_has_active_handles() -> i32 { return 1; } // Check for pending stdlib resolutions (turnloop P0: O(1) length mirrors). - if PENDING_RESOLUTIONS_LEN.load(Ordering::Acquire) != 0 - || PENDING_DEFERRED_LEN.load(Ordering::Acquire) != 0 + if (PENDING_RESOLUTIONS_LEN.load(Ordering::Acquire) != 0 + || PENDING_DEFERRED_LEN.load(Ordering::Acquire) != 0) + && has_own_pending_resolution() { return 1; } @@ -719,6 +847,63 @@ mod tests { assert_eq!(text.as_deref(), Some("Invalid password")); } + /// #11433: a resolution queued by another agent (a Worker) is neither + /// settled, nor scanned, nor counted as keep-alive work by this agent's + /// pump, and is purged when its agent retires. The fake addresses are never + /// dereferenced — which is the point: before the owner tag, this pump would + /// have "settled" a promise living in the Worker's heap. + #[test] + fn a_pump_leaves_another_agents_resolutions_for_their_owner() { + let _queues = own_pending_queues(); + let foreign = std::thread::spawn(perry_runtime::agent::enter_worker_agent) + .join() + .unwrap(); + assert_ne!(foreign, current_agent()); + let foreign_promise = 0x2345_5000usize; + let foreign_deferred = 0x2345_6000usize; + { + let _scope = ResolutionOwnerScope::enter(foreign); + queue_promise_resolution(foreign_promise, true, TAG_UNDEFINED_BITS); + queue_deferred_resolution(foreign_deferred, true, || { + panic!("another agent's converter must never run here") + }); + } + assert_eq!(js_stdlib_process_pending_resolutions_only(), 0); + assert_eq!(PENDING_RESOLUTIONS.lock().unwrap().len(), 1); + assert_eq!(PENDING_DEFERRED.lock().unwrap().len(), 1); + assert!(!has_own_pending_resolution()); + + let mut emitted = Vec::new(); + { + let mut mark = |value: f64| emitted.push(value.to_bits()); + let mut visitor = perry_runtime::gc::RuntimeRootVisitor::for_copy(&mut mark); + scan_pending_native_async_resolution_roots_mut(&mut visitor); + } + assert!( + emitted.is_empty(), + "this agent's collector visited another heap's promise: {emitted:x?}" + ); + + purge_agent_resolutions(foreign); + assert!(PENDING_RESOLUTIONS.lock().unwrap().is_empty()); + assert!(PENDING_DEFERRED.lock().unwrap().is_empty()); + assert_eq!(PENDING_RESOLUTIONS_LEN.load(Ordering::Acquire), 0); + assert_eq!(PENDING_DEFERRED_LEN.load(Ordering::Acquire), 0); + } + + const TAG_UNDEFINED_BITS: u64 = 0x7FFC_0000_0000_0001; + + /// The resolution half of the pump, without the unrelated subsystems + /// (`worker_threads`, readline, …) whose own queues it also drains. + fn js_stdlib_process_pending_resolutions_only() -> i32 { + let before = + PENDING_RESOLUTIONS.lock().unwrap().len() + PENDING_DEFERRED.lock().unwrap().len(); + js_stdlib_process_pending(); + let after = + PENDING_RESOLUTIONS.lock().unwrap().len() + PENDING_DEFERRED.lock().unwrap().len(); + (before - after) as i32 + } + #[test] fn stdlib_bridge_does_not_hard_reference_extension_pumps() { let source = include_str!("async_bridge.rs"); @@ -751,11 +936,13 @@ mod tests { let deferred_promise_ptr = 0x1234_6000usize; let result_bits = 0x7FFD_0000_1234_7000u64; PENDING_RESOLUTIONS.lock().unwrap().push(PendingResolution { + owner: current_agent(), promise_ptr, is_success: true, result_bits, }); PENDING_DEFERRED.lock().unwrap().push(DeferredResolution { + owner: current_agent(), promise_ptr: deferred_promise_ptr, is_success: true, converter: Box::new(|| 0), diff --git a/crates/perry-stdlib/src/container/executor.rs b/crates/perry-stdlib/src/container/executor.rs index 68b9588145..58c776ee03 100644 --- a/crates/perry-stdlib/src/container/executor.rs +++ b/crates/perry-stdlib/src/container/executor.rs @@ -23,7 +23,7 @@ use std::future::Future; use crate::common::async_bridge::{ blocking_thread_stack_size, ensure_gc_scanner_registered, ensure_pump_registered, pin_promise_for_native_resolution, queue_deferred_resolution, queue_promise_resolution, - InflightGuard, + resolution_owner, InflightGuard, ResolutionOwnerScope, }; /// Reject `ptr` with `message`, building the string on the main thread. @@ -50,7 +50,11 @@ fn run_detached( // Held for exactly the operation's run, and released even when the job is // dropped unrun, so the event loop stays alive until the result is queued. let inflight = keep_alive.then(InflightGuard::new); + // The job runs on a thread with no agent; settle for the one that asked + // (#11433), or a Worker's promise would be queued for the primary agent. + let owner = resolution_owner(); let run = move || { + let _owner = ResolutionOwnerScope::enter(owner); let result = perry_container_compose::rt::try_block_on(future) .unwrap_or_else(|e| Err(format!("container runtime unavailable: {e}"))); settle(result); diff --git a/crates/perry-stdlib/src/perry_ffi_async.rs b/crates/perry-stdlib/src/perry_ffi_async.rs index d2badee1e6..96c991f288 100644 --- a/crates/perry-stdlib/src/perry_ffi_async.rs +++ b/crates/perry-stdlib/src/perry_ffi_async.rs @@ -227,7 +227,11 @@ pub extern "C" fn perry_ffi_spawn_blocking(ctx: *mut c_void, invoke: extern "C" async_bridge::ensure_pump_registered(); let ctx_addr = ctx as usize; let inflight = async_bridge::InflightGuard::new(); + // `invoke` settles its promise from the pool thread, which has no agent; + // settle for the agent that asked (#11433). + let owner = async_bridge::resolution_owner(); let run = move || { + let _owner = async_bridge::ResolutionOwnerScope::enter(owner); invoke(ctx_addr as *mut c_void); drop(inflight); }; diff --git a/crates/perry-stdlib/src/worker_threads.rs b/crates/perry-stdlib/src/worker_threads.rs index cfa9b6f7a4..e3f6b1165b 100644 --- a/crates/perry-stdlib/src/worker_threads.rs +++ b/crates/perry-stdlib/src/worker_threads.rs @@ -14,6 +14,7 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::mpsc::{self, Sender}; use std::sync::{LazyLock, Mutex}; +use perry_runtime::agent::AgentId; use perry_runtime::closure::ClosureHeader; use perry_runtime::string::{js_string_from_bytes, StringHeader}; use perry_runtime::thread::{ @@ -103,6 +104,9 @@ thread_local! { static NEXT_BROADCAST_ID: RefCell = const { RefCell::new(10_000) }; /// Non-zero while running inside an in-process Worker thread. static CURRENT_WORKER_ID: Cell = const { Cell::new(0) }; + /// The agent that created the in-process Worker running on this thread: + /// where its `parentPort.postMessage` events are delivered. + static CURRENT_PARENT_AGENT: Cell = const { Cell::new(perry_runtime::agent::PRIMARY_AGENT) }; /// workerData for the current in-process Worker. static CURRENT_WORKER_DATA: RefCell> = const { RefCell::new(None) }; /// Worker threadName for the current in-process Worker. @@ -187,7 +191,13 @@ static WORKERS: LazyLock>> = /// `WORKERS` lock by `WorkerRecord::set_liveness` (records are never removed, /// only marked dead), so the per-turn keep-alive check is an atomic load. static LIVE_REFED_WORKERS: AtomicU64 = AtomicU64::new(0); -static PARENT_EVENTS: LazyLock>> = +/// Worker → parent events, each tagged with the agent that created the Worker +/// (#11433). Every agent's pump reaches `js_worker_threads_process_pending` — +/// a Worker's own await loop included — so an untagged queue let a Worker drain +/// the events addressed to its parent: its first `postMessage` was then +/// dispatched on the Worker's own thread, against listeners that live in the +/// parent's heap, and the parent never saw it. +static PARENT_EVENTS: LazyLock>> = LazyLock::new(|| Mutex::new(VecDeque::new())); type WorkerEntry = extern "C" fn(); @@ -918,8 +928,8 @@ fn event_name(value: f64) -> Option { string_value_to_string(value) } -fn push_parent_event(event: WorkerEvent) { - PARENT_EVENTS.lock().unwrap().push_back(event); +fn push_parent_event(parent: AgentId, event: WorkerEvent) { + PARENT_EVENTS.lock().unwrap().push_back((parent, event)); perry_runtime::event_pump::js_notify_main_thread(); } @@ -1291,6 +1301,7 @@ pub extern "C" fn js_worker_threads_message_channel_new() -> f64 { #[no_mangle] pub extern "C" fn js_worker_threads_worker_new(entry_ptr: i64, options: f64) -> f64 { ensure_worker_gc_scanner(); + worker_pump::ensure_parent_event_retire_hook(); crate::worker_threads::async_shim::ensure_pump_registered(); let worker_id = NEXT_WORKER_ID.fetch_add(1, Ordering::Relaxed); @@ -1348,6 +1359,8 @@ pub extern "C" fn js_worker_threads_worker_new(entry_ptr: i64, options: f64) -> drop(workers); let thread_options = options_state.clone(); + // Captured on the creating thread: this Worker's events go to THIS agent. + let parent_agent = perry_runtime::agent::current_agent(); // #8546: the Worker re-runs its module bodies on its own thread, but it is // the SAME image as its parent (same code addresses, same class ids), so it // shares the parent's class tables instead of building a second copy. Its @@ -1381,13 +1394,14 @@ pub extern "C" fn js_worker_threads_worker_new(entry_ptr: i64, options: f64) -> let worker_agent = perry_runtime::agent::enter_worker_agent(); let previous_env = apply_worker_env(&thread_options.env); CURRENT_WORKER_ID.with(|id| id.set(worker_id)); + CURRENT_PARENT_AGENT.with(|agent| agent.set(parent_agent)); CURRENT_WORKER_DATA.with(|slot| *slot.borrow_mut() = worker_data); CURRENT_THREAD_NAME .with(|slot| *slot.borrow_mut() = thread_options.thread_name.clone()); CURRENT_RESOURCE_LIMITS.with(|slot| slot.set(thread_options.resource_limits)); CURRENT_WORKER_CLOSE_REQUESTED.with(|closed| closed.set(false)); worker_surface::install_web_worker_globals(); - push_parent_event(WorkerEvent::Online(worker_id)); + push_parent_event(parent_agent, WorkerEvent::Online(worker_id)); let entry: WorkerEntry = unsafe { std::mem::transmute(entry_ptr as usize) }; let mut exit_code = 0; @@ -1470,15 +1484,15 @@ pub extern "C" fn js_worker_threads_worker_new(entry_ptr: i64, options: f64) -> let exit_code = match result { Ok(()) => exit_code, Err(_) => { - push_parent_event(WorkerEvent::Error(worker_id)); + push_parent_event(parent_agent, WorkerEvent::Error(worker_id)); 1 } }; - push_parent_event(WorkerEvent::Exit(worker_id, exit_code)); + push_parent_event(parent_agent, WorkerEvent::Exit(worker_id, exit_code)); }); if spawned.is_err() { - push_parent_event(WorkerEvent::Error(worker_id)); - push_parent_event(WorkerEvent::Exit(worker_id, 1)); + push_parent_event(parent_agent, WorkerEvent::Error(worker_id)); + push_parent_event(parent_agent, WorkerEvent::Exit(worker_id, 1)); } object_value(worker_obj) @@ -1583,7 +1597,8 @@ pub extern "C" fn js_worker_threads_post_message(data: f64) -> f64 { let worker_id = CURRENT_WORKER_ID.with(|id| id.get()); if worker_id != 0 { let message = unsafe { serialize_nanbox_for_thread(data.to_bits()) }; - push_parent_event(WorkerEvent::Message(worker_id, message)); + let parent = CURRENT_PARENT_AGENT.with(Cell::get); + push_parent_event(parent, WorkerEvent::Message(worker_id, message)); return js_undefined(); } let str_ptr = unsafe { js_json_stringify(data, 0) }; diff --git a/crates/perry-stdlib/src/worker_threads/worker_pump.rs b/crates/perry-stdlib/src/worker_threads/worker_pump.rs index 9f54e544bb..056b18243a 100644 --- a/crates/perry-stdlib/src/worker_threads/worker_pump.rs +++ b/crates/perry-stdlib/src/worker_threads/worker_pump.rs @@ -107,9 +107,16 @@ pub(super) fn start_stdin_reader() { pub extern "C" fn js_worker_threads_process_pending() -> i32 { let mut processed = 0; + // Only the events addressed to this agent (#11433); a Worker's own pump + // must leave its parent's events where the parent will find them. + let agent = perry_runtime::agent::current_agent(); let events: Vec = { let mut q = PARENT_EVENTS.lock().unwrap(); - q.drain(..).collect() + let (mine, theirs): (VecDeque<_>, VecDeque<_>) = std::mem::take(&mut *q) + .into_iter() + .partition(|(parent, _)| *parent == agent); + *q = theirs; + mine.into_iter().map(|(_, event)| event).collect() }; for event in events { match event { @@ -191,13 +198,37 @@ pub extern "C" fn js_worker_threads_process_pending() -> i32 { processed } +/// `perry_runtime::agent::retire_agent` hook: events addressed to an agent that +/// has exited can never be delivered (a nested Worker's children outliving it). +fn purge_parent_events(agent: perry_runtime::agent::AgentId) { + let dropped: VecDeque<_> = { + let mut q = PARENT_EVENTS.lock().unwrap(); + let (dead, live): (VecDeque<_>, VecDeque<_>) = std::mem::take(&mut *q) + .into_iter() + .partition(|(parent, _)| *parent == agent); + *q = live; + dead + }; + drop(dropped); +} + +pub(super) fn ensure_parent_event_retire_hook() { + static REGISTER: std::sync::Once = std::sync::Once::new(); + REGISTER.call_once(|| perry_runtime::agent::register_retire_hook(purge_parent_events)); +} + /// Check if worker_threads has pending work (stdin reader active) #[no_mangle] pub extern "C" fn js_worker_threads_has_pending() -> i32 { let started = STDIN_READER_STARTED.with(|s| *s.borrow()); let eof = STDIN_EOF.with(|eof| *eof.borrow()); let has_messages = PENDING_MESSAGES.with(|q| !q.borrow().is_empty()); - let has_worker_events = !PARENT_EVENTS.lock().unwrap().is_empty(); + let agent = perry_runtime::agent::current_agent(); + let has_worker_events = PARENT_EVENTS + .lock() + .unwrap() + .iter() + .any(|(parent, _)| *parent == agent); // turnloop P0: O(1) (`LIVE_REFED_WORKERS`); debug builds re-derive it. let has_live_refed_worker = LIVE_REFED_WORKERS.load(Ordering::Acquire) != 0; #[cfg(debug_assertions)] diff --git a/scripts/gc_runtime_root_holders.json b/scripts/gc_runtime_root_holders.json index 1c210fe1d5..53800d29f6 100644 --- a/scripts/gc_runtime_root_holders.json +++ b/scripts/gc_runtime_root_holders.json @@ -124,12 +124,6 @@ "verdict": "not_a_gc_pointer", "why": "Per-connection PROTOCOL state, keyed by the crate's own ws id. A `WsConnection` is `transport: WsTransport`, `messages: Vec` and three bools. Rule S fires on `WsTransport::Turnloop(i64)`; that i64 is the HOST's connection id, not an address \u2014 it is minted by the turnloop client adoption path (`attach_turnloop_client`) and its only use is as the first argument of `replay_command(conn_id, ..)`, which looks the connection up in the host registry. It is never cast to a pointer and never dereferenced. The other variant, `Connecting(Vec)`, holds `Send(WsOutgoing)` / `Close(Option, String)` / `Terminate`, and `WsOutgoing` is `Text(String)` / `Binary|Ping|Pong(Vec)` \u2014 Rust-owned bytes throughout, as is `WsPayload`. No NaN-boxed value and no heap pointer enters this map. The JS callbacks for these connections live in `WS_CLIENT_LISTENERS` (`HashMap>` of closure pointers), which is exactly what the registered `scan_ws_roots` visits \u2014 the split is deliberate. Newly uncovered because turnloop replaced the transport variant that held a tokio task's command channel with this bare id, which is what made rule S look at it." }, - { - "file": "crates/perry-runtime/src/agent.rs", - "name": "CURRENT_AGENT", - "verdict": "not_a_gc_pointer", - "why": "The calling thread's JS-agent id, an integer (Cell>, AgentId = u64) set once at worker entry and never cleared (agent.rs:75). It names a heap; it never holds a pointer into one. It stopped being reached incidentally by a registered scanner when turnloop P3 moved the timer root walk behind timer::store::with_current, whose generic closure the call-graph walk does not follow -- the holder itself did not change." - }, { "file": "crates/perry-runtime/src/alloc_census.rs", "name": "CREDIT", @@ -1411,6 +1405,12 @@ "verdict": "not_a_gc_pointer", "why": "turnloop P6. A monotonic i64 draft-id counter for the SMTP C seam; it contains only a tally." }, + { + "file": "crates/perry-stdlib/src/worker_threads.rs", + "name": "CURRENT_PARENT_AGENT", + "verdict": "not_a_gc_pointer", + "why": "#11433: the JS-agent id (AgentId = u64) of the agent that created the in-process Worker running on this thread, set once at Worker entry and used only to tag Worker-to-parent events. It names a heap; it never holds a pointer into one." + }, { "file": "crates/perry-stdlib/src/worker_threads.rs", "name": "CURRENT_WORKER_ID",