Skip to content

feat(mesh): elastic worker pools and the event-loop fixes that let them scale (RFC 0002, PR 2) - #262

Open
pd-admin-jc wants to merge 5 commits into
feat/provider-governorfrom
feat/elastic-pools
Open

pd-admin-jc wants to merge 5 commits into
feat/provider-governorfrom
feat/elastic-pools

Conversation

@pd-admin-jc

Copy link
Copy Markdown

RFC 0002, PR 2: elastic worker pools, plus the event-loop fixes that let them scale. Stacked on the provider-governor PR.

Pools

  • Usage: mesh.add(AgentClass, pool=Pool(min, max, idle_seconds, lane)) runs one logical agent across as many members as its queued work needs. Members share the agent id.
  • Dispatch: one dispatcher per pool claims ready steps and grows the pool up to max. Members idle past idle_seconds retire, down to min.
  • Diagnostics: get_diagnostics() reports each pool's state.

What stopped pools scaling, and the fixes

Each of these blocked the event loop, so every in-flight agent stalled with it:

Commit Problem Fix
63317c0 Each AutoAgent setup reloaded the 1,228-function registry and re-checked the atom catalogue, about 7 s of blocking work per pool member One registry per directory per process, built off the event loop
a83ce10 Step outputs were listed with SCAN MATCH over the whole database on every dispatch, about 1,750 round trips per call at 17k keys Outputs are indexed per workflow; loops wake on state change instead of sleeping out the 2 s poll
b7cfa81 Each worker pass made several round trips per step and tried to recover claims whose leases were still valid Active graphs are fetched in one pipelined call and the scan runs off the loop. A committed step requeues its blocked dependents before waking waiters, which fixes a done-plus-blocked race
b4d15ea Every goal waiter woke on every state change on the node A goal waiter wakes only for its own workflow

Measured

Evidence Engine claim verification on one laptop, 196 claims:

Run Wall clock
Fixed workers 202 min
Pools, before these fixes 97 min
16 slots, after 11.1 min
32 slots, after 9.4 min

The handoffs now take seconds: dispatch 0.4 s, handoff to the final step 2.5 s and commit 0.6 s at the median. The verifier's model turns are nearly all of the remaining time.

These numbers measure the framework overhead removed. Verifier research depth is being fixed in the product separately and will change the per-claim time.

Tests

  • New tests in tests/test_mesh_pools.py and tests/test_redis_store.py, including a test that fails if listing a workflow's outputs ever scans the keyspace.
  • The full suite passes.

Muyukani Ephraim Kizito added 5 commits October 8, 2026 19:12
… work queues

mesh.add(Agent, pool=Pool(min, max, idle_seconds, lane)) registers one logical agent. A dispatcher per pool scans for ready steps once and hands each to an idle member, creating members up to max and retiring idle ones above min after idle_seconds. Members are separate instances (an agent's kernel, subagent cache and sandbox serve one step at a time) sharing the registered agent's id, so checkpoints and resume_agent_id pinning work across members and the step-claim lease remains the only arbiter. Members inherit infrastructure, peer clients (not registered for inbound P2P routing) and auth; they are torn down with the mesh. Pool lane 'bulk' runs member LLM calls in the governor's bulk lane. Step discovery is factored into _ready_steps, shared by unpooled workers unchanged. Pools appear in get_diagnostics().
… off the event loop

Every AutoAgent setup reloaded the 1,228-function registry (~2 s) and re-checked the shipped catalogue (~5-8 s) synchronously, so each pool member spawned froze the whole process, including the agents it was added to relieve. The registry is now built once per directory per process on a worker thread, and the empty-catalogue warning fires only when the registry is actually empty.
…ng the keyspace

Every step dispatch listed a workflow's outputs with SCAN MATCH over the whole
database, synchronously on the event loop. With 17k keys that was ~1,750 round
trips a call, 720k SCANs in 50 minutes and one core pinned. Outputs are now
indexed per workflow.

Planner, workers, pools and goal waiters also wake as soon as this node changes
workflow state instead of sleeping out the poll interval, which remains the
cross-node fallback.
Each worker pass made several Redis round trips per step, and a WATCH/MULTI
recovery attempt on every in-progress step whose lease was still valid. Under
16 concurrent workflows a pass held the event loop ~0.44 s, and waking on each
state change multiplied the passes. Active graphs are now fetched in one
pipelined call and filtered in memory, recovery runs only for expired leases,
and the scan runs on a worker thread.

A committed step now requeues its blocked dependents before waking waiters, so
a waiter can no longer see done-plus-blocked and close the workflow before the
dependents run.
Every execute_goal waiter re-read its definition and every step on any state change
in the node, so 16 concurrent goals multiplied each change sixteenfold on the event
loop. Wakes now carry the workflow id; waiters listen for theirs, and workers and
pools still wake on any change.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant