From 7fb3a9d43243d1e4be7a2c278dbdd252eac2dba5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sat, 26 Sep 2026 08:08:46 +0000 Subject: [PATCH 1/6] fix(ext-net): each agent drains only its own socket events (#11340) --- crates/perry-ext-net/src/lib.rs | 29 +++++++++++++++++-- crates/perry-ffi/src/agent_post.rs | 21 ++++++++++++++ crates/perry-runtime/src/turnloop_post/abi.rs | 12 ++++++++ 3 files changed, 60 insertions(+), 2 deletions(-) diff --git a/crates/perry-ext-net/src/lib.rs b/crates/perry-ext-net/src/lib.rs index 1eeaaf82d4..1e91b809eb 100644 --- a/crates/perry-ext-net/src/lib.rs +++ b/crates/perry-ext-net/src/lib.rs @@ -175,9 +175,34 @@ pub(crate) mod statics { O.get_or_init(|| Mutex::new(HashMap::new())) } + /// The CALLING AGENT's queue of events waiting for its pump. + /// + /// #11340: this used to be one process-wide queue. Every agent's pump + /// drains through `js_ext_net_drain_pending`, and every push wakes the + /// primary thread, so a socket opened on a `worker_threads` worker had its + /// completions raced for by the primary: whichever pump ran first took + /// the worker's `'connect'` / `'data'` / `'close'`. Taken by the primary, + /// the event was dispatched against a socket whose listeners live on the + /// worker's heap — dropped (the worker then waited forever), delivered to + /// the wrong thread (`TypeError`s after the worker exited), or a crash. + /// Events are produced on the owning agent's thread (its loop's sink, or + /// an FFI call made by its JS), so keying by the agent that pushes is + /// keying by the agent that owns the socket. + /// + /// One queue per agent, created on the agent's first use and never freed: + /// the handful of bytes an exited worker's empty queue keeps is the price + /// of handing out `&'static` without a teardown hook. pub fn pending_events() -> &'static Mutex> { - static P: OnceLock>> = OnceLock::new(); - P.get_or_init(|| Mutex::new(Vec::new())) + type Queues = HashMap>>; + static QUEUES: OnceLock> = OnceLock::new(); + let agent = perry_ffi::agent_post::current_agent(); + let mut queues = QUEUES + .get_or_init(|| Mutex::new(HashMap::new())) + .lock() + .unwrap_or_else(|e| e.into_inner()); + queues + .entry(agent) + .or_insert_with(|| Box::leak(Box::new(Mutex::new(Vec::new())))) } /// HTTP Agent-owned socket handles are transport facades over the HTTP diff --git a/crates/perry-ffi/src/agent_post.rs b/crates/perry-ffi/src/agent_post.rs index 5d78c28cb6..7045724a4e 100644 --- a/crates/perry-ffi/src/agent_post.rs +++ b/crates/perry-ffi/src/agent_post.rs @@ -46,6 +46,27 @@ extern "C" { fn js_perry_agent_post_available() -> i32; fn js_perry_agent_post(run: Option, ctx: *mut c_void) -> i32; fn js_perry_agent_post_dispatched() -> u64; + fn js_perry_agent_current() -> u64; +} + +/// The agent the calling thread acts for: its own id on a `worker_threads` +/// worker, the primary agent's otherwise (including a second thread pumping +/// on the primary's behalf). +/// +/// A binding whose loop completions land in a process-wide queue keys that +/// queue by this, so each agent drains only work that names its own heap +/// (#11340). The value is an opaque identity: compare it, never interpret it. +pub fn current_agent() -> u64 { + #[cfg(any(not(test), feature = "runtime-link"))] + { + // SAFETY: a plain thread-local read in the linked runtime. + unsafe { js_perry_agent_current() } + } + #[cfg(all(test, not(feature = "runtime-link")))] + { + // No runtime is linked: every thread is the one agent there is. + 0 + } } /// Work a binding hands to the loop of the agent it is acting for. diff --git a/crates/perry-runtime/src/turnloop_post/abi.rs b/crates/perry-runtime/src/turnloop_post/abi.rs index fd11118d2f..c6958682ab 100644 --- a/crates/perry-runtime/src/turnloop_post/abi.rs +++ b/crates/perry-runtime/src/turnloop_post/abi.rs @@ -62,3 +62,15 @@ pub unsafe extern "C" fn js_perry_agent_post( pub extern "C" fn js_perry_agent_post_dispatched() -> u64 { super::dispatched() } + +/// The agent the calling thread acts for ([`crate::agent::current_agent`]): +/// its own if it is a `worker_threads` worker, otherwise the primary agent. +/// +/// For a binding that keeps a queue of work produced by several agents' loops +/// and must hand each agent only its own (#11340): a socket event produced on +/// a worker's loop names objects on the worker's heap, so the primary thread +/// must never drain it. +#[no_mangle] +pub extern "C" fn js_perry_agent_current() -> u64 { + crate::agent::current_agent() +} From f613adfb1c868511f739095bf686b1f3b4792fba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sat, 26 Sep 2026 09:00:38 +0000 Subject: [PATCH 2/6] test: gap test for worker-owned socket events (#11340) --- .../gap_11340_worker_net_events_worker.ts | 31 +++++++++++++++++ .../test_gap_11340_worker_net_events.ts | 34 +++++++++++++++++++ 2 files changed, 65 insertions(+) create mode 100644 test-files/_helpers/gap_11340_worker_net_events_worker.ts create mode 100644 test-files/test_gap_11340_worker_net_events.ts diff --git a/test-files/_helpers/gap_11340_worker_net_events_worker.ts b/test-files/_helpers/gap_11340_worker_net_events_worker.ts new file mode 100644 index 0000000000..917379d141 --- /dev/null +++ b/test-files/_helpers/gap_11340_worker_net_events_worker.ts @@ -0,0 +1,31 @@ +// The worker half of test_gap_11340_worker_net_events.ts. +import net from 'node:net'; +import { parentPort } from 'node:worker_threads'; + +const server = net.createServer((sock) => { + sock.on('data', (d: Buffer) => sock.end(d.toString().toUpperCase())); +}); +const port = await new Promise((resolve) => { + server.listen(0, '127.0.0.1', () => resolve((server.address() as { port: number }).port)); +}); + +const ROUNDS = 25; +let completed = 0; +const replies = new Set(); +for (let i = 0; i < ROUNDS; i++) { + const reply = await new Promise((resolve) => { + // pg's shape (lib/connection.js): construct, connect, THEN subscribe. + const s = new net.Socket(); + s.setNoDelay(true); + s.connect(port, '127.0.0.1'); + let got = ''; + s.once('connect', () => s.write('round' + i)); + s.on('data', (d: Buffer) => { got += d.toString(); }); + s.on('close', () => resolve(got)); + s.on('error', (e: Error) => resolve('error ' + e.message)); + }); + if (reply === 'ROUND' + i) completed++; + replies.add(reply.replace(/[0-9]+$/, '')); +} +server.close(); +parentPort?.postMessage(`rounds ${completed}/${ROUNDS} ${[...replies].sort().join(',')}`); diff --git a/test-files/test_gap_11340_worker_net_events.ts b/test-files/test_gap_11340_worker_net_events.ts new file mode 100644 index 0000000000..45388f2959 --- /dev/null +++ b/test-files/test_gap_11340_worker_net_events.ts @@ -0,0 +1,34 @@ +// #11340: socket events for a socket opened on a `worker_threads` worker must +// be delivered on that worker. perry-ext-net kept ONE process-wide queue of +// pending socket events and woke the primary thread on every push, so the +// primary's pump raced the worker's for them: a worker's `'connect'` taken by +// the primary was dispatched against listeners on the worker's heap — dropped +// (the worker waited forever: compiled `pg` hung in `connect()` on a worker), +// run on the wrong thread (a `TypeError` after the worker exited, as mysql2 +// showed), or a crash. +// +// The worker runs its own echo server and 25 sequential client sockets in +// node-postgres' shape (construct, `connect()`, then `once('connect')`), which +// lost roughly one round in five before the fix. Everything stays inside the +// worker, so the primary's only job is to keep pumping while it waits. +import { Worker } from 'node:worker_threads'; + +const worker = new Worker(new URL('./_helpers/gap_11340_worker_net_events_worker.ts', import.meta.url)); + +// A lost event presents as a hang; the watchdog turns it into a visible failure. +const watchdog = setTimeout(() => { + console.log('WORKER STOPPED'); + process.exit(3); +}, 15000); +if (typeof (watchdog as { unref?: () => void }).unref === 'function') { + (watchdog as { unref: () => void }).unref(); +} + +const result = await new Promise((resolve) => { + worker.on('message', (value: string) => resolve(String(value))); + worker.on('error', (e: Error) => resolve(`worker-error ${e.message}`)); +}); +clearTimeout(watchdog); +console.log(result); +await worker.terminate(); +console.log('done'); From c6d9bfe3185c7a38e23e79f40817ae00df78a30d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sat, 26 Sep 2026 16:41:46 +0000 Subject: [PATCH 3/6] fix(ext-net): an agent's keepalive counts only its own sockets and servers (#11340) --- crates/perry-ext-net/src/adopt.rs | 2 ++ crates/perry-ext-net/src/ipc.rs | 1 + crates/perry-ext-net/src/lib.rs | 13 +++++++++++++ crates/perry-ext-net/src/server_state.rs | 13 +++++++++++-- crates/perry-ext-net/src/turnloop_io.rs | 16 ++++++++++------ 5 files changed, 37 insertions(+), 8 deletions(-) diff --git a/crates/perry-ext-net/src/adopt.rs b/crates/perry-ext-net/src/adopt.rs index 5557bca031..60a81db622 100644 --- a/crates/perry-ext-net/src/adopt.rs +++ b/crates/perry-ext-net/src/adopt.rs @@ -56,6 +56,7 @@ pub fn adopt_upgraded_tcp_stream(stream: std::net::TcpStream) -> i64 { id, SocketState { tcp_async_id: 0, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id: 0, shutdown_async_id: 0, awaiting_connect: false, @@ -144,6 +145,7 @@ pub fn adopt_turnloop_upgrade(id: i64) -> bool { id, SocketState { tcp_async_id: 0, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id: 0, shutdown_async_id: 0, awaiting_connect: false, diff --git a/crates/perry-ext-net/src/ipc.rs b/crates/perry-ext-net/src/ipc.rs index 88259a1d80..b26a8dcc39 100644 --- a/crates/perry-ext-net/src/ipc.rs +++ b/crates/perry-ext-net/src/ipc.rs @@ -19,6 +19,7 @@ fn allocate_socket() -> i64 { id, SocketState { tcp_async_id: 0, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id: 0, shutdown_async_id: 0, awaiting_connect: false, diff --git a/crates/perry-ext-net/src/lib.rs b/crates/perry-ext-net/src/lib.rs index 1e91b809eb..6bcead6a8b 100644 --- a/crates/perry-ext-net/src/lib.rs +++ b/crates/perry-ext-net/src/lib.rs @@ -253,6 +253,10 @@ pub(crate) mod statics { /// pass instead of needing a second per-server scanner. pub(crate) struct ServerState { pub async_id: u64, + /// The agent that created this server (`perry_ffi::agent_post::current_agent`). + /// Only that agent's pump delivers its events, so only that agent's + /// keepalive may count it (#11340). + pub owner_agent: u64, /// Set by `.listen()`. Keeps `has_active_handles` answering yes from the /// `listen()` call until the server's registry entry is removed by its /// `'close'`, including the window before the bind has been published. @@ -274,6 +278,10 @@ pub(crate) struct ServerState { pub(crate) struct SocketState { pub(crate) tcp_async_id: u64, + /// The agent the socket belongs to — the one whose thread created it + /// (`perry_ffi::agent_post::current_agent`). Its events go to that agent's + /// queue, so only that agent's keepalive may count it (#11340). + pub(crate) owner_agent: u64, pub(crate) connect_async_id: u64, pub(crate) shutdown_async_id: u64, /// True only between `js_net_socket_alloc` and the first @@ -427,6 +435,7 @@ pub(crate) fn register_turnloop_socket( socket_id, SocketState { tcp_async_id: 0, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id: 0, shutdown_async_id: 0, awaiting_connect: false, @@ -469,6 +478,7 @@ impl SocketState { pub(crate) fn for_test(awaiting_connect: bool) -> Self { SocketState { tcp_async_id: 0, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id: 0, shutdown_async_id: 0, awaiting_connect, @@ -721,6 +731,7 @@ pub unsafe extern "C" fn js_net_socket_alloc() -> i64 { id, SocketState { tcp_async_id, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id: 0, shutdown_async_id: 0, awaiting_connect: true, @@ -778,6 +789,7 @@ pub unsafe extern "C" fn js_net_create_server( id, ServerState { async_id: 0, + owner_agent: perry_ffi::agent_post::current_agent(), listen_armed: false, bound_port: 0, bound_host: String::new(), @@ -1260,6 +1272,7 @@ where id, SocketState { tcp_async_id, + owner_agent: perry_ffi::agent_post::current_agent(), connect_async_id, shutdown_async_id: 0, awaiting_connect: false, diff --git a/crates/perry-ext-net/src/server_state.rs b/crates/perry-ext-net/src/server_state.rs index 03bf15bb4b..ae4ceeed70 100644 --- a/crates/perry-ext-net/src/server_state.rs +++ b/crates/perry-ext-net/src/server_state.rs @@ -395,8 +395,15 @@ pub(crate) fn has_active_handles() -> bool { if crate::turnloop_io::close_in_flight() { return true; } + // #11340: only this agent's sockets and servers. A worker's socket is + // served by the worker's pump; counting it here kept the primary alive + // forever once the worker had exited with that socket's events undrained. + let agent = perry_ffi::agent_post::current_agent(); if statics::sockets().lock().unwrap().values().any(|socket| { - socket.refed && !socket.destroyed && (socket.is_open || !socket.awaiting_connect) + socket.owner_agent == agent + && socket.refed + && !socket.destroyed + && (socket.is_open || !socket.awaiting_connect) }) { return true; } @@ -405,7 +412,9 @@ pub(crate) fn has_active_handles() -> bool { .unwrap() .iter() .any(|(id, server)| { - (server.listening || server.listen_armed) && crate::bun_tcp::server_keeps_alive(*id) + server.owner_agent == agent + && (server.listening || server.listen_armed) + && crate::bun_tcp::server_keeps_alive(*id) }) } diff --git a/crates/perry-ext-net/src/turnloop_io.rs b/crates/perry-ext-net/src/turnloop_io.rs index eea60b9ddc..387feeafb0 100644 --- a/crates/perry-ext-net/src/turnloop_io.rs +++ b/crates/perry-ext-net/src/turnloop_io.rs @@ -141,25 +141,29 @@ struct Aux { /// so a program with nothing else pending exited in between and the /// `'close'` never fired. undici's `Client.close()` awaits exactly that /// event, so it never settled. -fn closing() -> &'static Mutex> { - static CLOSING: OnceLock>> = OnceLock::new(); - CLOSING.get_or_init(|| Mutex::new(std::collections::HashSet::new())) +/// Sockets whose close was submitted and whose `'close'` is still owed, each +/// with the agent that submitted it (#11340: only that agent waits for it). +fn closing() -> &'static Mutex> { + static CLOSING: OnceLock>> = OnceLock::new(); + CLOSING.get_or_init(|| Mutex::new(std::collections::HashMap::new())) } /// Whether a submitted socket close still owes its `'close'` event; the /// keepalive gate (`server_state::has_active_handles`) waits for it. pub(crate) fn close_in_flight() -> bool { - !closing() + let agent = perry_ffi::agent_post::current_agent(); + closing() .lock() .unwrap_or_else(|e| e.into_inner()) - .is_empty() + .values() + .any(|owner| *owner == agent) } fn note_close_submitted(id: i64) { closing() .lock() .unwrap_or_else(|e| e.into_inner()) - .insert(id); + .insert(id, perry_ffi::agent_post::current_agent()); } fn submit_socket_close(id: i64) -> bool { From 58c919e3677bbcf539983a9701ce268a7ae057c7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sat, 26 Sep 2026 20:23:30 +0000 Subject: [PATCH 4/6] chore(gc): drop the stale perry-ext-net pending_events root-holder entry (#11340) --- scripts/gc_runtime_root_holders.json | 6 ------ 1 file changed, 6 deletions(-) diff --git a/scripts/gc_runtime_root_holders.json b/scripts/gc_runtime_root_holders.json index 00fde96abb..1c210fe1d5 100644 --- a/scripts/gc_runtime_root_holders.json +++ b/scripts/gc_runtime_root_holders.json @@ -69,12 +69,6 @@ "verdict": "not_a_gc_pointer", "why": "write_tokens(): HashMap. Each i64 is the same handle-band id ABORTS holds \u2014 the key used to look the socket up in crate::statics::sockets(), never a heap address. Entries are removed on write completion (bun_tcp.rs) and retained-out when a handle closes, so no id outlives its facade." }, - { - "file": "crates/perry-ext-net/src/lib.rs", - "name": "P", - "verdict": "not_a_gc_pointer", - "why": "pending_events(): Vec; every variant carries socket/server ids, Bytes, String, bool, or DropInfo (SocketAddrs) \u2014 no closures or NaN-boxed values. Listener closures live in listeners(), visited by scan_net_roots." - }, { "file": "crates/perry-ext-net/src/server_state.rs", "name": "STATE", From 166d0bb0e74bf90995a700d61980e0a3b394ea1f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sat, 26 Sep 2026 21:04:27 +0000 Subject: [PATCH 5/6] test: #11340 gap test drives the worker's sockets against the harness echo server --- .../gap_11340_worker_net_events_worker.ts | 25 +++++++++---------- .../test_gap_11340_worker_net_events.ts | 9 ++++--- 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/test-files/_helpers/gap_11340_worker_net_events_worker.ts b/test-files/_helpers/gap_11340_worker_net_events_worker.ts index 917379d141..c3acf80855 100644 --- a/test-files/_helpers/gap_11340_worker_net_events_worker.ts +++ b/test-files/_helpers/gap_11340_worker_net_events_worker.ts @@ -2,30 +2,29 @@ import net from 'node:net'; import { parentPort } from 'node:worker_threads'; -const server = net.createServer((sock) => { - sock.on('data', (d: Buffer) => sock.end(d.toString().toUpperCase())); -}); -const port = await new Promise((resolve) => { - server.listen(0, '127.0.0.1', () => resolve((server.address() as { port: number }).port)); -}); - +// The parity harness runs a plain-TCP echo server here +// (test-files/test_net_echo_server.py), the same one test_net_min.ts uses. +const PORT = 17891; const ROUNDS = 25; let completed = 0; const replies = new Set(); for (let i = 0; i < ROUNDS; i++) { const reply = await new Promise((resolve) => { - // pg's shape (lib/connection.js): construct, connect, THEN subscribe. + // node-postgres' shape (lib/connection.js): construct, connect, THEN subscribe. const s = new net.Socket(); s.setNoDelay(true); - s.connect(port, '127.0.0.1'); + s.connect(PORT, '127.0.0.1'); + const want = 'round' + i; let got = ''; - s.once('connect', () => s.write('round' + i)); - s.on('data', (d: Buffer) => { got += d.toString(); }); + s.once('connect', () => s.write(want)); + s.on('data', (d: Buffer) => { + got += d.toString(); + if (got.length >= want.length) s.end(); + }); s.on('close', () => resolve(got)); s.on('error', (e: Error) => resolve('error ' + e.message)); }); - if (reply === 'ROUND' + i) completed++; + if (reply === 'round' + i) completed++; replies.add(reply.replace(/[0-9]+$/, '')); } -server.close(); parentPort?.postMessage(`rounds ${completed}/${ROUNDS} ${[...replies].sort().join(',')}`); diff --git a/test-files/test_gap_11340_worker_net_events.ts b/test-files/test_gap_11340_worker_net_events.ts index 45388f2959..4f2a06a17e 100644 --- a/test-files/test_gap_11340_worker_net_events.ts +++ b/test-files/test_gap_11340_worker_net_events.ts @@ -7,10 +7,11 @@ // run on the wrong thread (a `TypeError` after the worker exited, as mysql2 // showed), or a crash. // -// The worker runs its own echo server and 25 sequential client sockets in -// node-postgres' shape (construct, `connect()`, then `once('connect')`), which -// lost roughly one round in five before the fix. Everything stays inside the -// worker, so the primary's only job is to keep pumping while it waits. +// The worker opens 25 sequential client sockets, in node-postgres' shape +// (construct, `connect()`, then `once('connect')`), to the harness's echo +// server (test-files/test_net_echo_server.py on 127.0.0.1:17891, the one +// test_net_min.ts uses). Before the fix the worker stalled within the 25 +// rounds; the primary does nothing but wait. import { Worker } from 'node:worker_threads'; const worker = new Worker(new URL('./_helpers/gap_11340_worker_net_events_worker.ts', import.meta.url)); From a3c2bc53e0add9a76699983879745c0fb53840fb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Sat, 26 Sep 2026 23:52:11 +0000 Subject: [PATCH 6/6] changelog: #11444 --- changelog.d/11444-net-events-per-agent.md | 7 +++++++ 1 file changed, 7 insertions(+) create mode 100644 changelog.d/11444-net-events-per-agent.md diff --git a/changelog.d/11444-net-events-per-agent.md b/changelog.d/11444-net-events-per-agent.md new file mode 100644 index 0000000000..d09a0a68f0 --- /dev/null +++ b/changelog.d/11444-net-events-per-agent.md @@ -0,0 +1,7 @@ +**fix(ext-net): each agent drains only its own socket events, so compiled pg and mysql2 inside `worker_threads` no longer hang or misdeliver events (#11340).** perry-ext-net kept one process-wide queue of pending socket events, and every push woke the primary thread. A socket opened on a worker therefore had its events drained by whichever agent's pump ran first. + +When the primary took a worker's `'connect'`, `'data'` or `'close'`, it dispatched the event against listeners on the worker's heap. Depending on timing, the event was lost and the worker hung (pg's `connect()`), it ran on the wrong thread (a `TypeError` after the worker exited, as seen with mysql2), or the process crashed. + +`pending_events()` now returns the calling agent's queue, keyed by the new `perry_ffi::agent_post::current_agent()`. Sockets, servers and in-flight closes record their owner agent, so an agent's keepalive counts only its own handles. A new gap test, `test_gap_11340_worker_net_events`, runs 25 pg-shaped client sockets from a worker: 2 of 20 runs passed on main, 20 of 20 with this fix. Verified against a live PostgreSQL 16.15 (pg 8.23.0) and MySQL 8.0.46 (mysql2 3.24.4). + +A worker that hosts its own server still hangs intermittently. That is a separate defect, filed as #11434.