From 21c677427348c2b925f3d849fdc5e9df87371b28 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sun, 27 Sep 2026 00:16:52 +0200 Subject: [PATCH 1/3] fix(runtime): settle and dispatch cross-agent async work only on the owning agent Every JS agent (the primary thread and each node:worker_threads Worker) runs the same process-wide pumps, but three queues on the worker-agent network path were untagged, so whichever agent pumped first took everything: - perry-stdlib's promise-resolution queues (async_bridge): the primary agent settled a Worker's fetch promise and ran its continuation against the main heap ("immediate 000undefined", "Invalid response handle"), and the Worker settled the primary's. - worker_threads' Worker->parent event queue: a Worker's own pump drained the events addressed to its parent, so a postMessage was dispatched on the Worker thread and the parent never saw it. - perry-ext-http's server pump: a Worker's pump drained requests queued for the primary agent's http.Server and ran its handler on the Worker thread, so the response was never written and the Worker's fetch hung. Each entry now carries the agent it belongs to (perry_runtime::agent, the model timers and thread results already use), pumps and keep-alive checks consider only their own agent's entries, the GC scanner visits only its own heap's promises, and retire_agent purges a dead agent's leftovers through a new retire-hook registry. Pool/fallback threads that settle a promise re-assert the submitting agent with ResolutionOwnerScope. Fixes #11433 --- crates/perry-ext-http/src/server/mod.rs | 24 ++ crates/perry-ext-http/src/server/server.rs | 56 ++++- .../src/server/server/in_flight.rs | 16 +- .../src/server/turnloop_serve/conn.rs | 36 ++- crates/perry-runtime/src/agent.rs | 29 +++ .../perry-stdlib/src/common/async_bridge.rs | 215 ++++++++++++++++-- crates/perry-stdlib/src/container/executor.rs | 6 +- crates/perry-stdlib/src/perry_ffi_async.rs | 4 + crates/perry-stdlib/src/worker_threads.rs | 33 ++- .../src/worker_threads/worker_pump.rs | 35 ++- 10 files changed, 411 insertions(+), 43 deletions(-) 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)] From e76fe32f0a03b7921068ccfe06b8a6482da5d8ef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sun, 27 Sep 2026 01:52:46 +0200 Subject: [PATCH 2/3] chore(gc-audit): record CURRENT_PARENT_AGENT verdict, drop stale CURRENT_AGENT entry CURRENT_AGENT is now reached by a registered scanner through the new agent-scoped async_bridge scanner, so its exemption went stale; the new worker_threads CURRENT_PARENT_AGENT holder is an agent id, not a pointer. --- scripts/gc_runtime_root_holders.json | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) 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", From 24c4b35d8f4587744f756c4bcb12c628d2a722ee Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sun, 27 Sep 2026 01:53:49 +0200 Subject: [PATCH 3/3] changelog: add fragment for #11445 --- changelog.d/11445-cross-agent-pump-ownership.md | 7 +++++++ 1 file changed, 7 insertions(+) create mode 100644 changelog.d/11445-cross-agent-pump-ownership.md 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.