Conversation
|
Auto-sync is disabled for draft pull requests in this repository. Workflows must be run manually. Contributors can view more details about this message here. |
This was referenced Oct 1, 2026
Draft
Draft
ananthsub
force-pushed
the
ananthsub/partial-ckpt-telemetry
branch
from
October 2, 2026 13:22
9775d38 to
dc49274
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-multi-worker
branch
from
October 2, 2026 13:22
d64723f to
a92580a
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-telemetry
branch
from
October 2, 2026 19:43
dc49274 to
f9affd7
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-multi-worker
branch
from
October 2, 2026 19:43
a92580a to
6647f1c
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-telemetry
branch
from
October 2, 2026 21:17
f9affd7 to
f8f2b11
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-multi-worker
branch
from
October 2, 2026 21:17
6647f1c to
ec6b311
Compare
uvicorn's workers share one listening socket, so each request goes to whichever worker accepts the connection, while session state lives in the memory of the worker that created the session. With the stateful counter resources server at 4 workers, 200 concurrent episodes through /run scored 1.0 on only 7; swe_rebench (8 workers) can lose a session's sandbox at /verify and score 0. With num_workers > 1, resources and agent servers now: - stamp each new session with the ID of the worker that created it, inside the signed session cookie; - serve the same app on a private Unix socket per worker, in a directory the main process creates under /tmp; - forward a request for a session another worker owns to that worker's private socket, unchanged, and stream the reply back unchanged. A forwarded request is always handled locally; - answer 410 when the owner's socket is gone, instead of running against empty state. The router sits outside SessionMiddleware and decodes the cookie with the same secret, so a forwarded reply carries only the owner's Set-Cookie. Single-worker servers and model servers are unchanged. The counter resources server gains the module-level app that multi-worker uvicorn imports, and an opt-in e2e test (NEMO_GYM_MULTIWORKER_E2E=1) runs it with 4 workers behind the legacy /run relay and Simple Agent. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
CLI harnesses in sandboxes reach a resources server's tools over /mcp with the signed session token minted at /seed_session, and send no Gym session cookie. With several workers, those calls ran on whichever worker accepted them: 4 workers, 200 concurrent counter sequences over MCP, 65 of 200 correct on main and 72 of 200 with cookie routing alone. With routing active, the MCP token now also carries the owning worker's ID. The router reads the owner from the token on the MCP path, and on any request without a session cookie, verifying it with the same serializer and salt. A token without an owner (minted by a single-worker server or before this change) or with a bad signature is handled locally, where the MCP endpoint's own check applies as before. Single-worker tokens are unchanged. With this change, 200 of 200 are correct. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
The policy model server already runs as one checkpoint participant across several uvicorn workers: a coordinator in the main process owns the controller, and each worker links to it over a Unix socket. Environment, agent, and resources servers need the same machinery. Move the parts that do not depend on the policy model into nemo_gym/_checkpoint/workers.py: - the framed channel between the coordinator and its workers; - CoordinatedParticipant: worker registration, closing and reopening every worker, sequenced readiness reports, fence broadcast, and the blockers for unreported or lost workers; - WorkerCoordinator, which serves it on a background-thread event loop; - WorkerLink: forwarding control calls, joining while a checkpoint is open, reporting, and stopping the worker when the coordinator is gone. The policy model keeps only its own state: the gate, the restored generation cuts and their claims, and the ledger export and import. Its behavior and its tests are unchanged. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…ral workers These servers refused num_workers > 1 when checkpointing was on: each worker would have closed, exported, and restored only its own share of the episodes and sessions. With session routing in place, they now run as one participant across workers, on the worker coordinator the policy model server uses. Each worker runs the participant a single process runs, linked to the coordinator in the main process: - Prepare closes every worker and is ready once every worker reports ready. A worker re-reports while a checkpoint is open, so readiness that changes without a notification, such as a session ending, still reaches the coordinator. - Commit asks every worker for its records and writes the participant's one records file and manifest. Resources and agent session records carry their routing owner, the worker ID in the session's cookie or MCP token. The manifest lists every owner a cookie may name. - Restore installs resources sessions and native agent sessions on live workers, one worker per pre-crash owner, spread evenly. Every worker's router gets the alias table from old owner to new worker, also for owners that exported nothing, such as a stateless resources server's. A request whose cookie or token names a pre-crash worker reaches the worker that holds its session instead of getting 410. A worker that joins later gets the table when it registers. - Environment episodes and legacy agent /run episodes stay with the coordinator. The worker that receives the replacement /run claims the record, and exactly one claim succeeds. A claim is answered before any later message to that worker, so a claimed episode is always either exported by its worker or waited for by the next checkpoint. Claims are refused while a checkpoint is open, as new episodes are. Unclaimed records are exported again by the next checkpoint. - Retire and attempt fences reach every worker; the coordinator drops unclaimed records of retired attempts. - An unregistered or lost worker blocks prepare and commit, and a worker that loses the coordinator stops itself. A legacy /run's call to the agent's own /v1/responses names its worker in the new x-ng-session-owner header, so the turn loop's activation runs where the episode is tracked. Simple Agent no longer refuses sessions with several workers, and it and the weather resources server gain the module-level app that multi-worker uvicorn imports. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
The environment, agent, and resources servers in the e2e deployment take a server_workers count, and the crash and continue scenarios (native, legacy, resources state with its negative control, a second crash before the replacement, and a checkpoint without a crash) run with one worker and with two. The scale test takes the same parameter. A new scenario starts 16 rollouts at once, native over agent sessions and legacy over counter resources sessions, checkpoints them mid-episode, kills every Gym process, restores, and runs the replacements. With two workers the checkpoint's sessions come from both workers; every replacement scores 1.0, and the backend sees only the calls each rollout still needed. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
A checkpoint taken with one worker could not restore resources sessions or native agent sessions into servers with several workers. Routing was off at commit, so their records and cookies name no owner, and the restore failed. Users will change the worker count between a checkpoint and a restore, so it must work in any direction. The coordinator now places each session record without an owner on a live worker of its own, spread evenly with the owners' groups, and installs it there; that worker owns it from then on. Every worker's router gets a second table, alongside the owner aliases, from restored session ID to worker, for these sessions only. When a cookie or MCP token names no owner, the router looks the decoded session ID up in it before handling the request locally. Workers that register later get it too. The manifest lists the placed sessions, so a session restored again before its cookie was stamped with its new owner keeps a table entry, pointed at its owner's new worker. The table is replaced at each restore and holds at most one entry per restored session; an entry outlives its session until then, since a session ends on one worker and the table is on all of them. Agent session records now carry the cookie's session ID, which the session middleware exposes to code without the request at hand. The e2e scenario with 16 concurrent rollouts now restarts the servers with a different worker count: 1 to 2, 2 to 1, and 2 to 4, as well as unchanged. Before this change, both 1 to 2 cases failed at restore. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…kers A CLI agent harness reaches the resources server only through MCP tokens. Its calls during an open checkpoint wait on the worker that owns the session, and each harness keeps its token when Gym restarts with a different worker count. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…eral workers The coordinator now reports as pending both the restored records no worker has claimed and the restored sessions each worker has not used yet, so a commit that no longer continues their episode retires them on every worker. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
ananthsub
force-pushed
the
ananthsub/partial-ckpt-telemetry
branch
from
October 2, 2026 23:32
f8f2b11 to
14f23fa
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-multi-worker
branch
from
October 2, 2026 23:32
ec6b311 to
51dc9c8
Compare
This branch has not been deployed
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
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
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.
What changed and why
The bottom two commits are #3905 (session routing). This PR carries them until #3905 merges. Review only the commits after them:
refactor(checkpoint): a worker coordinator for any participant kindfeat(checkpoint): resources, agent, and environment servers with several workerstest(checkpoint): run the e2e scenarios with two workers per serverfeat(checkpoint): restore sessions onto a different number of workersUntil now, resources, agent, and environment servers refused
num_workers > 1when checkpointing was on. uvicorn's workers share one port, so a control call reaches one arbitrary worker. If each worker ran its own participant, a checkpoint would close, export, and restore only that worker's share of the episodes and sessions. Production servers run several workers (swe_rebenchruns 8, several sandboxed agents run 4), so this guard kept them out of partial-rollout checkpoints.This PR makes each of these servers one participant across its workers:
A generic worker coordinator. The policy model server already runs as one participant across workers. The parts of it that do not depend on the policy model move to
nemo_gym/_checkpoint/workers.py:The policy model keeps only its own state, and its behavior and tests are unchanged.
Prepare closes every worker. It is ready once every worker reports ready. While a checkpoint is open, a worker re-reports its readiness on a short poll, so readiness that changes without a notification still reaches the coordinator. One example is a resources session that ends.
Commit asks every worker for its records. The coordinator writes the participant's one records file and manifest, so the store format and the controller's view do not change.
Restore of sessions that span requests (resources sessions and native agent sessions):
Restore of state that lives inside one request (environment episodes and legacy agent
/runepisodes):/runclaims its record, and exactly one claim succeeds.Retire reaches every worker. Each worker refuses the attempt only while its own retire stops it, and the coordinator drops unclaimed records of retired attempts. No lasting fence is copied to workers or broadcast on resume.
Restored state no longer continued is released on every worker. The coordinator gathers the restored records no worker has claimed and the restored sessions each worker has not used yet, so a commit whose scope leaves them out retires them.
Fail closed. Fewer registered workers than
num_workersblocks prepare. So does a worker lost while a checkpoint is open, until the controller resumes.Two smaller changes support this:
/runcalls the agent's own/v1/responsesfor its turn loop. That call now names its worker in a newx-ng-session-ownerheader, so the turn loop's activation runs on the worker that tracks the episode. The router honors the header only for a request without a session cookie or MCP token.appthat multi-worker uvicorn imports.How it works
Where this PR sits in the overall flow
The highlighted part is what this PR adds.
flowchart LR C["Controller<br/>NeMo RL, or rollout collection"] CO["Coordination<br/>prepare, commit, restore, resume, retire"] K["Control plane on every server<br/>phases, lease, storage,<br/>retire: stop, then free"] subgraph G["One participant per Gym server"] E["Environment server<br/>episode steps"] M["Policy model<br/>held responses, generation cuts"] A["Agent<br/>sessions parked at boundaries"] R["Resources server<br/>session state"] end W["Inference worker<br/>stages cut prefixes"] D[("Checkpoint directory<br/>records, then manifest")] L[("Capture ledger<br/>retire and delete from #3938, #3939")] C --> CO --> K K --> E & M & A & R M --> W M --> L G --> D classDef this fill:#fde68a,stroke:#b45309,stroke-width:2px,color:#1f2937 class E,A,R thisOne checkpoint, a crash, and the restore, end to end:
This PR
One server with several uvicorn workers is still one participant. Its coordinator, in the main process, drives each worker's own participant; the router sends a session's requests to the worker that holds it.
Restoring sessions into new workers. A cookie or MCP token names the worker that held the session before the crash; the alias table sends it to the worker that holds it now.
Whole-episode state, such as an environment episode or a legacy
/run, waits in the coordinator. Exactly one worker claims it, whichever receives the replacement attempt's/run.Where this sits in the stack
This is one PR in a stack of draft PRs for partial-rollout checkpointing. Each PR's base is the branch of the PR before it, so each diff shows only that PR's commits.
The stack is based on the token-capture cleanup PRs #3938 (capture ledger
retireanddelete) and #3939 (complete-recordretireanddelete): #3882's base is #3939's branch. The checkpoint stack uses those operations to free the ledgers of retired attempts and to clear a dead execution's capture files before a restore. The striped lock files of #3937 are independent of the stack.ananthsub/partial-ckpt-core): feat(checkpoint): add the participant control plane, episode steps, and coordinationananthsub/partial-ckpt-policy-model): feat(checkpoint): make policy model servers checkpoint participants with generation cutsananthsub/partial-ckpt-model-worker-cuts): feat(token-capture): add worker staging helpers for generation cutsananthsub/partial-ckpt-environment): feat(checkpoint): continue environment server episodes from their boundariesananthsub/partial-ckpt-resources): feat(checkpoint): add the resources server participant with asynchronous session hooksananthsub/partial-ckpt-verifier-declarations): feat(checkpoint): declare five training verifiers stateless with replayable verificationananthsub/partial-ckpt-agent): feat(checkpoint): add the agent session participant and Simple Agent continuationananthsub/partial-ckpt-e2e): test(checkpoint): add a process-level end-to-end suite driven by coordinationananthsub/partial-ckpt-rollout-collection): feat(checkpoint): checkpoint evaluation runs from rollout collectionananthsub/partial-ckpt-telemetry): feat(checkpoint): spans and metrics for partial-rollout checkpointsananthsub/partial-ckpt-multi-worker): feat(checkpoint): resources, agent, and environment servers with several workers (this PR)Server ports that build on the stack: #3894 (Workplace Assistant), #3895 (Gymnasium), #3896 (Blackjack), #3897 (indirect prompt injection), #3898 (proof refinement). Session routing for several workers is #3905, against main.
Relationship to the old stack
The old stack (#2939 to #2946) required one worker for every checkpointing server except the policy model. The policy model already had multi-worker support in #3883, as a coordinator in the main process with worker gates. This PR generalizes that coordinator rather than adding a second mechanism. It replaces the three
num_workers=1guards the re-cut stack kept for resources, agent, and environment servers.Issue
No tracking issue exists. Multi-worker support for these servers was a stack requirement. Session routing, which it depends on, is #3905.
Validation
Run on this branch, on top of #3939's branch:
RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_checkpoint_*.py tests/unit_tests/test_session_routing.py: 185 passed.test_checkpoint_participant_workers.pyruns each scenario with a real coordinator socket and two worker links:/runare each claimed by exactly one worker;RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider responses_api_agents/simple_agent/tests: 30 passed.environment_servers/single_agent_turn/tests: 14 passed, after removing the test of the one-worker refusal this PR removes.NEMO_GYM_CHECKPOINT_E2E=1 RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/e2e/checkpoint tests/e2e/session_routing, one test at a time: 43 passed, 9 skipped (the opt-in real-model and scale tests), in about 17 minutes. The scenarios run with the environment, agent and resources servers at one and two workers, and the policy model at one and two, including:tests/unit_tests, eight processes, withouttest_opensandbox_compose.py, which needs the optionalopensandboxpackage): 6,722 passed, 2 failed. The two failures, a sandbox retry test and a Slurm script test, fail the same way onmainin this development environment.pre-commit run --files <files changed by this PR>: all hooks passed, and no hook modified a file.Signed-off-byline.Rollout evidence
Real-model rollout evidence: pending, as for the rest of the stack. The opt-in real vLLM test was not run.
Scale, on one workstation with every Gym server and the fake inference backend. The run used native
single_agent_turnweather rollouts, with token capture and generation cuts on. Each rollout's second model call was held, so about half the rollouts were in flight at the checkpoint. The test checkpoints, kills every Gym process group, restarts, restores, and runs the replacement attempts. Every server ran with two workers:Compatibility
Resources, agent, and environment servers now accept
num_workers > 1with checkpointing on. Single-worker servers behave as before. The visible differences are in the records:ownerfield,nullwith one worker, in resources and agent session records;session_idfield in agent session records;session_ownersandplaced_sessionsin a multi-worker participant's manifest.The worker count may change between a checkpoint and a restore, in either direction: from one worker to several, from several to one, or between different counts.
With several workers, a restored environment or legacy episode whose replacement
/runarrives while a checkpoint is open is refused withcheckpoint_parked, as a new episode is, and the caller retries it. With one worker, it is admitted and parks at its first boundary.With several workers, commit and restore carry each worker's records over the coordinator's socket. The frame limit is now 1 GiB, up from the policy model's 64 MiB, which only rejects a corrupt length prefix. Servers with very large sessions put all their records through the main process. Per-worker record files are the alternative if that becomes a bottleneck.
New header
x-ng-session-owner: on a server with several workers, a request without a session cookie or MCP token that sets it is forwarded to that worker.