-
-
Notifications
You must be signed in to change notification settings - Fork 164
fix(ext-net): each agent drains only its own socket events; pg/mysql2 in worker_threads no longer hang or misdeliver (#11340) #11444
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
Merged
Merged
Changes from all commits
Commits
Show all changes
6 commits
Select commit
Hold shift + click to select a range
7fb3a9d
fix(ext-net): each agent drains only its own socket events (#11340)
f613adf
test: gap test for worker-owned socket events (#11340)
c6d9bfe
fix(ext-net): an agent's keepalive counts only its own sockets and se…
58c919e
chore(gc): drop the stale perry-ext-net pending_events root-holder en…
166d0bb
test: #11340 gap test drives the worker's sockets against the harness…
a3c2bc5
changelog: #11444
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,30 @@ | ||
| // The worker half of test_gap_11340_worker_net_events.ts. | ||
| import net from 'node:net'; | ||
| import { parentPort } from 'node:worker_threads'; | ||
|
|
||
| // 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<string>(); | ||
| for (let i = 0; i < ROUNDS; i++) { | ||
| const reply = await new Promise<string>((resolve) => { | ||
| // 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'); | ||
| const want = 'round' + i; | ||
| let got = ''; | ||
| 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++; | ||
| replies.add(reply.replace(/[0-9]+$/, '')); | ||
| } | ||
| parentPort?.postMessage(`rounds ${completed}/${ROUNDS} ${[...replies].sort().join(',')}`); |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,35 @@ | ||
| // #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 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)); | ||
|
|
||
| // 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<string>((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'); |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
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
🔎 Supported by static analysis
🏁 Script executed:
Repository: PerryTS/perry
Length of output: 8117
🏁 Script executed:
Repository: PerryTS/perry
Length of output: 42320
🏁 Script executed:
Repository: PerryTS/perry
Length of output: 45503
🏁 Script executed:
Repository: PerryTS/perry
Length of output: 41700
🏁 Script executed:
Repository: PerryTS/perry
Length of output: 42611
Reclaim each agent’s pending network queue at worker exit.
The map entry and empty queue are small, but the queue can retain
PendingNetEvent::Databuffers and error strings. EachDataevent retains a refcounted read-buffer allocation, documented as 16 KiB per read. The queue is unbounded andretire_agentdoes not remove it, so short-lived workers can retain material payload storage after exit.Add a teardown hook that removes the retiring agent’s queue and drops its pending events during
retire_agent. Do not dispatch events after the worker has exited or keep a permanent leaked reference.🤖 Prompt for AI Agents