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-environment
branch
2 times, most recently
from
October 1, 2026 21:53
7c6475d to
b9ac5f1
Compare
This was referenced Oct 1, 2026
ananthsub
force-pushed
the
ananthsub/partial-ckpt-environment
branch
from
October 1, 2026 22:05
b9ac5f1 to
6c95786
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-environment
branch
from
October 2, 2026 13:22
6c95786 to
1739ac6
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-environment
branch
from
October 2, 2026 19:43
1739ac6 to
2812ca7
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-environment
branch
from
October 2, 2026 21:17
2812ca7 to
5aa1530
Compare
…ndaries Every environment server takes part in checkpoints once checkpointing is on. The environment server owns the episode protocol, so it records where each episode can continue: - An episode records a boundary after each completed step and runs each step in wait or replay mode (nemo_gym._checkpoint.steps). Prepare closes episode admission and is ready once every live episode is parked at a boundary or inside a replay step. - A restored boundary is handed to the replacement attempt's run, which continues from it. The task must match the one recorded. - single_agent_turn runs seed, invoke_agent, close_agent, verify, and return as explicit stages with typed session handles. Seeding and the agent invocation replay; closing the agent waits; verification uses the mode the resources server reported on its seed reply, kept in the session handles so a restored episode still knows it, and its result is recorded before the episode returns. - The single_agent_turn_legacy adapter checkpoints too, since its rows run the same protocol. - A restored episode whose replacement has not started yet is exported again by the next checkpoint, so a second crash still continues it. - Checkpointing refuses num_workers > 1, since episodes live in one process. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
A controller retires stragglers while a checkpoint is open, and the episode steps cancel a retired episode's task. When the episode had already ended and was running its final cleanup, which closes its agent and resources sessions, that cancellation interrupted the cleanup and leaked both sessions: the cleanup's anyio shield holds against anyio cancellation, not a native task.cancel(). The environment server now tells the participant when an episode enters final cleanup. A retire after that point stops tracking the episode but does not cancel its task. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…has claimed A restored episode record waits for its replacement attempt's /run. The environment participant now reports unclaimed records as pending, so a commit that no longer continues their episode releases them instead of exporting them again at every checkpoint. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…e and its final cleanup Retiring an episode now cancels it, unless its final cleanup already started, and waits until it has ended, so the retire replies only after the episode's agent and resources sessions are closed. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
pthombre
added this pull request to stack #3962
October 2, 2026 22:50
ananthsub
force-pushed
the
ananthsub/partial-ckpt-environment
branch
from
October 2, 2026 23:32
5aa1530 to
a53fcad
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
Every environment server takes part in checkpoints once checkpointing is on. The environment server owns the episode protocol, so it records where each episode can continue.
waitorreplaymode. Prepare closes episode admission and is ready once every live episode is parked at a boundary or inside a replay step./run, which continues from it. The task must match the one recorded.single_agent_turnruns seed, invoke agent, close agent, verify, and return as explicit stages with typed session handles. Seeding and the agent invocation replay; closing the agent waits; verification uses the mode the resources server reported on its seed reply, kept in the session handles so a restored episode still knows it.single_agent_turn_legacyadapter checkpoints too, since its rows run the same protocol./runreplies on every path, and a retire never interrupts it. A retire cancels an episode only before cleanup starts, then waits until the episode has ended, so its reply means the sessions are closed. This matters because anyio's shielded scope does not stop a nativetask.cancel().num_workers > 1on environment servers, since episodes live in one process.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 thisOne checkpoint, a crash, and the restore, end to end:
This PR
single_agent_turnrecords a boundary after every step. A checkpoint exports the latest boundary, and a restored attempt starts at itsnextstep.Final cleanup closes the episode's sessions before
/runreplies. A retire never interrupts it: it cancels an episode only before cleanup starts, and then waits for the episode to end.flowchart TD X{"retire arrives"} -->|"episode still in its protocol"| C["cancel the episode"] --> F X -->|"final cleanup already started"| F["wait"] F --> CL["final cleanup closes the agent and resources sessions"] --> D["retire replies"]Where this sits in the stack
This is one PR in a stack of draft PRs that re-cut partial-rollout checkpointing onto environment servers. 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 boundaries (this PR)ananthsub/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 workersPartial-rollout checkpointing lets a training controller, such as NeMo RL, checkpoint Gym while rollouts are in flight and, after a crash, continue those rollouts from their last safe point instead of starting them over. Checkpointing is off by default; with the
checkpoint:block unset, no server installs any checkpoint routes or behavior.Relationship to the old stack (#2939 to #2946)
New compared with the old stack, which had no environment servers. It takes over the protocol side of #2942 (where an episode can continue), which the old stack reconstructed from agent and resources state after the fact.
Issue
No tracking issue exists for this re-cut. The design and the mapping from the old stack are described in this PR series, and the old stack's PRs (#2939 to #2946) carry the original discussion.
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: 97 passed, including, each failing without its change:RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_environment_server.py: 16 passed.RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider environment_servers/single_agent_turn/tests: 15 passed.RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider environment_servers/single_agent_turn_legacy/tests: 10 passed.tests/unit_tests, eight processes, withouttest_opensandbox_compose.py, which needs the optionalopensandboxpackage): 6,518 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
Compatibility
checkpoint:block setsenabled: true.num_workers > 1fails at startup.