-
-
Notifications
You must be signed in to change notification settings - Fork 166
fix(runtime): settle and dispatch cross-agent async work only on the owning agent (#11433) #11445
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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. |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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,24 +1041,25 @@ 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::<HttpServer, _>(|s| { | ||
| if server_is_active(s) { | ||
| if s.owned_here() && server_is_active(s) { | ||
| active = 1; | ||
| } | ||
| }); | ||
| if active != 0 { | ||
| return 1; | ||
| } | ||
| iter_handles_of::<crate::server::https_server::HttpsServer, _>(|s| { | ||
| if server_is_active(&s.base) { | ||
| if s.base.owned_here() && server_is_active(&s.base) { | ||
| active = 1; | ||
| } | ||
| }); | ||
| if active != 0 { | ||
| return 1; | ||
| } | ||
| iter_handles_of::<crate::server::http2_server::Http2SecureServer, _>(|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::<HttpServer>(server_handle) { | ||
| return s.owned_here(); | ||
| } | ||
| if let Some(s) = get_handle::<crate::server::https_server::HttpsServer>(server_handle) { | ||
| return s.base.owned_here(); | ||
| } | ||
| get_handle::<crate::server::http2_server::Http2SecureServer>(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<i64> = Vec::new(); | ||
| perry_ffi::iter_handle_ids_of::<HttpServer, _>(|id| http_handles.push(id)); | ||
| http_handles.retain(|h| get_handle::<HttpServer>(*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::<crate::server::https_server::HttpsServer, _>(|id| { | ||
| https_handles.push(id) | ||
| }); | ||
| https_handles.retain(|h| { | ||
| get_handle::<crate::server::https_server::HttpsServer>(*h) | ||
| .is_some_and(|s| s.base.owned_here()) | ||
| }); | ||
| for h in https_handles { | ||
| count += drain_deferred_listen_for::<crate::server::https_server::HttpsServer, _>(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::<crate::server::http2_server::Http2SecureServer, _>(|id| { | ||
| h2_handles.push(id) | ||
| }); | ||
| h2_handles.retain(|h| { | ||
| get_handle::<crate::server::http2_server::Http2SecureServer>(*h) | ||
| .is_some_and(|s| s.base.owned_here()) | ||
| }); | ||
|
Comment on lines
+1233
to
+1236
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift Scope HTTP/2 events to their owning agent. This filter limits HTTP/2 requests to owned servers, but 🤖 Prompt for AI Agents |
||
| for h in h2_handles { | ||
| count += drain_deferred_listen_for::<crate::server::http2_server::Http2SecureServer, _>( | ||
| h, | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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<Vec<InFlightRequest>> = 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; | ||
|
Comment on lines
+122
to
+125
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift Clean up HTTP queues when an agent retires. The new owner filters preserve foreign entries, but a retired Worker has no future pump. Its queued handles and events can remain in process-wide storage. Register HTTP retirement cleanup and release handles where required. (raw.githubusercontent.com)
📍 Affects 3 files
🤖 Prompt for AI Agents |
||
| } | ||
| 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 | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift
Apply agent ownership to the HTTP/2 client keep-alive check.
After the server checks become agent-scoped,
has_active_h2_clients()still returns true for any pending HTTP/2 event, pending transport work, or active client session in the process. A Worker’s HTTP/2 client can therefore keep the primary agent’s loop active after its own work finishes. Scope those checks to the calling agent as well. (raw.githubusercontent.com)🤖 Prompt for AI Agents