Skip to content

Shared executor pools: fix the thread-per-node explosion behind the red macOS CI leg - #75

Merged
ptesavol merged 2 commits into
mainfrom
claude/shared-executor-pools
Jul 10, 2026
Merged

ptesavol merged 2 commits into
mainfrom
claude/shared-executor-pools

Conversation

@ptesavol

Copy link
Copy Markdown
Collaborator

Fixes the macOS CI leg that has been red since milestone A landed (#73), and removes a scaling hazard for real applications.

The failure

Layer1ScaleTest.MultipleLayer1Dht and KademliaCorrectnessTest.CanFindCorrectNeighbors abort on the macOS runners with:

libc++abi: terminating due to uncaught exception of type std::__1::system_error:
thread constructor failed: Resource temporarily unavailable

Every DhtNode instance owned ~9 dedicated executor threads (PeerManager pings, DhtNode + PeerDiscovery recovery, Router routing, StoreManager replication, plus a server-side and a 20-thread client-side pool per RPC communicator). The 200–250-node tests peak at ~1900–2100 threads — beyond what the 3-core CI runners tolerate (a 10-core dev machine allows 16384 threads per process, so it never fired locally). The same profile would hurt real applications: trackerless-network creates one layer-1 node per stream partition, so thread count scaled linearly with subscriptions. The TS original runs all of this on one event loop; fixed-size shared pools are the C++ analogue.

New in streamr-utils

  • SharedExecutors — two process-wide pools: worker() (general detached work: pings, recovery, replication, RPC dispatch) and background() (low OS priority; keeps the former per-Router nice-10 "routing runs whenever we have time" design from phase A4).
  • SharedSerialExecutor — a per-owner serial view of a shared pool, replacing the former {1}-sized pools whose single thread WAS the ordering guarantee (and, in Router's case, the synchronization of its maps). All folly::SerialExecutor template machinery instantiates inside the SharedExecutors TU — instantiating it from an importing module TU trips the known BMI name-lookup pitfall (DistributedMutex internals).
  • GuardedAsyncScope — AsyncScope with a close() gate for owners whose event handlers can still fire while stop() drains (folly forbids add() racing joinAsync()).

Conversions (teardown + ordering semantics preserved per class)

Class Change Teardown guarantee
PeerManager pings → serial view of worker() scope join in stop() outside the mutex — the phase-AA pattern, verbatim
DhtNode rejoins → serial view GuardedAsyncScope drained in stop() (strictly earlier than the former destructor join; a KBucketEmpty firing mid-drain is dropped by the gate)
PeerDiscovery, StoreManager, Router/RoutingSession executor swap only their detached tasks pin self via sharedFromThis, exactly as before
RpcCommunicatorServerApi serial view kept its existing CancellableAsyncScope destructor drain
RpcCommunicatorClientApi restructured request()/notify() now run the this-touching work as a scope task; the caller awaits only a promise-contract future that holds no this. The new destructor's scope drain replaces the former 20-thread pool's join — and a timed-out request can no longer leave a detached task referencing a dead communicator (detachOnCancel previously detached a this-capturing task whose only safety net was the pool destructor).

Verification

  • Peak thread count for the failing tests: 1870 → 18 (Kademlia, 200 nodes) and 2075 → 21 (MultipleLayer1Dht, ~250 nodes), measured by sampling ps -M over full runs.
  • No performance change: isolated pre-refactor runs took 229 s / 254 s; post-refactor 229 s / 256 s. (An apparent "43 s baseline" turned out to be a warm full-suite artifact — isolated single-process runs, as ctest/CI executes them, are the valid comparison.)
  • Identical convergence: avg 6.84/8 correct neighbours (6.86 before).
  • All unit + integration suites green in streamr-utils, streamr-proto-rpc, streamr-dht; the phase-AA-sensitive SimultaneousConnections/ConnectionLocking/MultipleEntryPointJoining tests pass repeated isolated runs; lint green in all three packages.

Notes

🤖 Generated with Claude Code

…ed macOS CI leg

The macOS CI leg has been red since milestone A landed:
Layer1ScaleTest.MultipleLayer1Dht and KademliaCorrectnessTest abort with
"thread constructor failed: Resource temporarily unavailable"
(pthread_create EAGAIN). Every DhtNode instance owned ~9 dedicated
executor threads (PeerManager pings, DhtNode/PeerDiscovery recovery,
Router routing, StoreManager replication, and 2 per RPC communicator), so
the 200-250 node tests peaked at ~1900-2100 threads — beyond what the
3-core CI runners allow, and the same profile would hurt a real
application running many layer-1 nodes (trackerless-network creates one
per stream partition). TS runs all of this on one event loop.

New in streamr-utils:
- SharedExecutors: two process-wide pools — worker() (general detached
  work) and background() (low-priority, keeps the former per-Router
  nice-10 "routing runs whenever we have time" design).
- SharedSerialExecutor: a per-owner serial view of a shared pool,
  replacing the former {1}-sized pools whose single thread WAS the
  ordering (and, in Router's case, the synchronization) guarantee. All
  folly::SerialExecutor template machinery instantiates inside the
  SharedExecutors TU (the known BMI instantiation pitfall).
- GuardedAsyncScope: AsyncScope with a close() gate for owners whose
  event handlers can still race stop()'s drain.

Converted (semantics preserved per class):
- PeerManager: pings on a serial view + scope join in stop() outside the
  mutex — the phase-AA teardown pattern, unchanged.
- DhtNode: k-bucket-empty rejoins in a GuardedAsyncScope drained in
  stop() (strictly earlier than the former destructor join).
- PeerDiscovery, StoreManager, Router/RoutingSession: executor swap only;
  their detached tasks pin `self` (sharedFromThis), exactly as before.
- RpcCommunicatorServerApi: serial view of the shared pool (kept its
  existing CancellableAsyncScope drain).
- RpcCommunicatorClientApi: request/notify restructured so the
  `this`-touching work runs as a scope task and the caller awaits only a
  promise-contract future holding no `this`; the scope drain in the new
  destructor replaces the former 20-thread pool's join, and a timed-out
  request can no longer leave a detached task referencing a dead
  communicator.

Verified: peak threads for the two failing tests drop 1870/2075 -> 18/21
with identical durations (isolated pre-refactor runs: 229 s/254 s;
post-refactor: 229 s/256 s — the apparent "43 s baseline" was a warm
full-suite artifact) and identical convergence (avg 6.84/8 correct
neighbours). All unit+integration suites green in streamr-utils,
streamr-proto-rpc and streamr-dht; the phase-AA-sensitive
SimultaneousConnections/ConnectionLocking/MultipleEntryPointJoining tests
pass repeated isolated runs; lint green in all three packages.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@cursor

cursor Bot commented Jul 10, 2026

Copy link
Copy Markdown

Bugbot is not enabled for this team, so this pull request was not reviewed.

Enable Bugbot in the Cursor dashboard to get automatic reviews on future PRs.

…eout

The linux-arm64 leg caught a teardown hang the refactor introduced: when
a CALLER abandons a request (ConnectionManager::gracefullyDisconnect
wraps its RPC in a 2 s outer timeout), the scope task kept awaiting the
request future with no bound of its own — and once the peer is torn down
nothing ever resolves that future, so the communicator destructor's
scope drain blocked forever (ctest killed the test at 300 s on both
attempts; the fast ubuntu runner never hit it because its graceful
disconnects completed in time).

The RPC timeout now wraps the future await INSIDE the scope task, where
it survives caller abandonment: every scope task settles within the
request timeout, and the destructor's cancelAndJoin resolves it promptly
before that.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@ptesavol
ptesavol merged commit d929be5 into main Jul 10, 2026
6 checks passed
ptesavol added a commit that referenced this pull request Jul 10, 2026
… shared pool

Applies the PR #75 executor architecture to the one per-instance pool B1
added: the detached requestConnection notifications now run on a serial
view of the shared worker pool, tracked by a GuardedAsyncScope that
destroy() drains after abort (outside mMutex, the PeerManager::stop()
pattern); the scope's gate drops a connect() racing the drain.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
ptesavol added a commit that referenced this pull request Jul 10, 2026
* Phase B1 (part 1): connectivity checker, request handler, signed peer descriptors

Ports the first half of milestone B1 from v103.8.0-rc.3:

- createPeerDescriptorSignaturePayload + createPeerDescriptor: the local
  node's signed peer descriptor built from a ConnectivityResponse; the node
  id derives from the reported ip (last 13 bytes of keccak(ip) + last 7
  bytes of sign(ip)) and the descriptor is signed over the payload — both
  through SigningUtils' Ethereum-magic keccak/secp256k1, matching the TS
  EcdsaSecp256k1Evm defaults bit-exactly (wire-visible). Like the TS TODO,
  the key pair is throwaway (random private key, 20 random publicKey bytes).
- connectivityChecker (client side): connectSync opens a raw websocket
  client connection with a connect timeout; sendConnectivityRequest runs
  the connectivity check against an entry point and returns its
  ConnectivityResponse (5 s response timeout, protocol-version check).
  Deliberately synchronous-blocking: both call sites run on dedicated
  worker threads, and the listeners are registered before connect()/send()
  so a fast completion cannot be lost.
- connectivityRequestHandler (server side): parses ConnectivityRequests off
  a connection's data events, probes the requester's advertised websocket
  server (action=connectivityProbe) unless the port is 0, and replies with
  a ConnectivityResponse. Runs on a dedicated magic-static executor so the
  probe never blocks the websocket dispatch thread. GeoIP lookup omitted
  (deferred to milestone E). Template over the connection type so the unit
  test drives it with a mock (the TS test uses a bare EventEmitter).
- New error types ConnectionFailed / ConnectivityResponseTimeout.

Tests ported (all green, unit suite 177/177):
- unit/createPeerDescriptor.test.ts -> CreatePeerDescriptorTest (3 tests)
- unit/connectivityRequestHandler.test.ts -> ConnectivityRequestHandlerTest
  (2 tests; the happy path exercises a real websocket probe round trip
  against a local WebsocketServer)

Also records in the plan the owner's manual follow-up: compare the
KademliaCorrectness statistics against the TypeScript implementation (the
TS benchmark is bit-rotted at the pin; C++ numbers and the diagnosis are in
the note).

Still to come in B1 (part 2): wiring the handler into
WebsocketServerConnector's connectivityRequest/connectivityProbe actions,
the real client-side checkConnectivity flow in the connector facade
(honoring externalIp/websocketHost/port-range options), and the
ConnectivityChecking/Websocket/WebsocketConnectionManagement integration
tests.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Phase B1 (part 2): real connectivity checking and websocket reverse connections

WebsocketServerConnector now runs the ported TS flows instead of stubs:
- start() wires attachConnectivityRequestHandler into the
  connectivityRequest action and pins connectivityProbe sockets until the
  prober closes them (C++ has no GC to keep an unowned server socket
  alive; the pin is released by onClosed dropping all listeners).
- checkConnectivity() does real connectivity checking against shuffled
  entry points with 2 s abortable waits between attempts and
  WebsocketServerStartError after the list is exhausted; the
  options-info branch tolerates a serverless connector (empty host/port,
  like the TS undefined fields).
- connect()/isPossibleToFormConnection()/requestConnectionFromPeer():
  the reverse-connection path — ask a serverless peer over the signaling
  transport to open a websocket back to this node's server; the RPC
  notification runs detached on a dedicated executor that destroy()
  joins after abort (the PeerManager::stop() pattern).

DefaultConnectorFacade routes createConnection through the server
connector when the client connector cannot form the connection, and
passes the new entryPoints option through (externalIp and the TLS/
autocertify options stay deferred to B2/E).

Fixed in passing (found by the new tests):
- WebsocketClientConnectorRpcRemote::requestConnection returned the lazy
  notify() task whose const& parameters referenced already-destroyed
  locals (SEGV in Any::PackFrom on first use); it now co_awaits so the
  locals live in its own coroutine frame.
- attachConnectivityRequestHandler captured the connection weakly, so a
  connectivityRequest server socket had no owner and was destroyed
  before the request arrived; the Data listener now holds it strongly
  (cycle broken by onClosed's removeAllListeners).

Tests ported (v103.8.0-rc.3): integration/ConnectivityChecking.test.ts,
integration/Websocket.test.ts (round-trip case added to the existing
WebsocketClientServerTest), integration/
WebsocketConnectionManagement.test.ts (all five cases, including the
serverless reverse-connection ones). lint.sh: three new files added to
the documented clangd std-type-unification exclusion list.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* CI: raise ctest per-test timeout to clear the heavy DHT convergence tests

The macOS CI leg has been red since PR #73 landed: ctest's --timeout 300
(sized when the whole suite ran in ~35 s) kills
Layer1ScaleTest.MultipleLayer1Dht (~214 s on a fast 10-core dev machine,
longer on the 3-core macOS runners) and KademliaCorrectnessTest, then
fails the until-pass retry the same way. The faster Linux runners stay
under the limit, which is why only macOS failed. 1200 s clears the heavy
tests with headroom while still converting a genuine hang into a failure.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* WebsocketServerConnector: move requestConnection notifications to the shared pool

Applies the PR #75 executor architecture to the one per-instance pool B1
added: the detached requestConnection notifications now run on a serial
view of the shared worker pool, tracked by a GuardedAsyncScope that
destroy() drains after abort (outside mMutex, the PeerManager::stop()
pattern); the scope's gate drops a connect() racing the drain.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* DEBUG: instrument the reverse-connection notify path (temporary)

The ubuntu-latest leg SIGSEGVs (address 0x20) deterministically in the
three WebsocketConnectionManagement reverse-connection tests while macOS
and local arm64 runs (including under ASan) are green. These INFO logs
bracket every step of requestConnectionFromPeer -> GuardedAsyncScope ->
ClientApi::notify so the next CI run shows the exact crash point.
Will be reverted once the cause is fixed.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* Revert "DEBUG: instrument the reverse-connection notify path (temporary)"

This reverts commit 2f5bcbe.

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant