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-rollout-collection
branch
from
October 1, 2026 19:38
b5ffb47 to
22a07dd
Compare
This was referenced Oct 1, 2026
ananthsub
force-pushed
the
ananthsub/partial-ckpt-rollout-collection
branch
from
October 1, 2026 21:53
22a07dd to
324d650
Compare
This was referenced Oct 1, 2026
ananthsub
force-pushed
the
ananthsub/partial-ckpt-rollout-collection
branch
from
October 1, 2026 22:05
324d650 to
29a2feb
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-rollout-collection
branch
from
October 2, 2026 13:22
29a2feb to
d1bac74
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-rollout-collection
branch
from
October 2, 2026 19:43
d1bac74 to
37f360b
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-rollout-collection
branch
from
October 2, 2026 21:17
37f360b to
1ae4776
Compare
pthombre
added this pull request to stack #3962
October 2, 2026 22:50
The controller of a checkpoint has to be the process that dispatches episodes: only it knows which episodes it will continue, and only it can send their replacements after a restore. For an evaluation run that is rollout collection, so with `checkpoint_dir` set it drives Gym's coordination itself. A checkpoint stops starting rows, prepares every participant, commits the rows whose /run has not replied, writes the run's manifest, and publishes it by pointing LATEST at it; the two latest checkpoints are kept. SIGUSR1 checkpoints and continues, SIGTERM checkpoints and stops (Slurm can send it before preempting a job), and `checkpoint_every_s` checkpoints on a timer. A run restarted with resume_from_cache restores the latest checkpoint before dispatching: its unfinished rows continue as their next attempt, and the rest start from their input as before. If the restore fails, coordination has retired that attempt, so those rows start from their input one attempt later. A fresh run forgets the checkpoints of an earlier run in the same directory. Without `checkpoint_dir`, rollout collection is unchanged. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…heckpoints - A row whose /run got no reply is now retired everywhere, in batches under the checkpoint lock. Before, state it left behind, such as a restored record nothing claimed, was exported again by every later checkpoint. A reply, even a failure, means the environment server already ended the episode and released what it held, so those rows are not retired. - Rows restored from a checkpoint but not dispatched yet are continued by the next checkpoint, at the attempt they will run as. Before, they were missing from its manifest, so a crash at that point restarted them from their input. - Publishing removes partial checkpoint directories left by a failed commit and temporary files left by a crash, and a failed write removes its own temporary file. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
ananthsub
force-pushed
the
ananthsub/partial-ckpt-rollout-collection
branch
from
October 2, 2026 23:32
1ae4776 to
a307627
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 rest of the stack lets a controller checkpoint Gym with rollouts in flight. The controller has to be the process that sends episodes: only it knows which episodes it will continue, and only it can send their replacements after a restore. A training framework such as NeMo RL plays that role for training. For an evaluation run it is rollout collection, which so far could only restart unfinished rows from their inputs (
resume_from_cache). This PR makes rollout collection the controller of its own rows, so a long evaluation task survives a preemption without starting over.With
checkpoint_dirset (and the globalcheckpoint:block enabled):Triggers.
SIGUSR1checkpoints the rows in flight and continues.SIGTERMcheckpoints them and stops the run. Slurm can send it ahead of a preemption, for example with--signal=B:TERM@120.checkpoint_every_scheckpoints on a timer, which covers crashes that give no warning.checkpoint_dir/collector.pid, for sending the signals.A checkpoint:
/runhas not replied;checkpoint_dir/LATESTat it;The two latest checkpoints are kept. A checkpoint that does not prepare in time publishes nothing, and the run continues from its previous checkpoint.
Restart. A run restarted with
resume_from_cache=truerestores the latest checkpoint before dispatching anything.resume_from_cache=false) forgets the checkpoints of an earlier run in the same directory, as it already does with the output file.Code.
nemo_gym/_checkpoint/collection.pyholds the controller.rollout_collection.py: three config fields, a gate and an in-flight count around each row's/run, and the restore step after the cache load.gym evaland the collection CLI print a message and exit cleanly when a preemption signal stopped the run.Cleanup.
/rungot no reply is retired everywhere, in batches under the checkpoint lock, so a retire never lands inside a checkpoint. Otherwise state it left behind, such as a restored record nothing claimed, would be exported again by every later checkpoint. A reply, even a failure, means the environment server already ended the episode and released what it held.Without
checkpoint_dir, rollout collection behaves exactly as before. The dispatch hook is only passed when it is set, so subclasses that override the internal dispatch method keep working.How it works
Where this PR sits in the overall flow
The highlighted part is what this PR adds.
flowchart LR C["Rollout collection, this PR<br/>SIGTERM, SIGUSR1, or a timer"] 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 C thisOne checkpoint, a crash, and the restore, end to end:
This PR
Checkpointing an evaluation run, and continuing it after the job restarts.
sequenceDiagram participant S as Slurm or an operator participant RC as Rollout collection participant G as Gym servers participant D as checkpoint_dir S->>RC: SIGTERM, SIGUSR1, or the timer fires RC->>RC: stop starting rows RC->>G: prepare RC->>G: commit the rows whose /run has not replied RC->>D: collection.json, then LATEST alt SIGTERM RC-->>S: exit, the job is requeued else SIGUSR1 or the timer RC->>G: resume, rows start again end Note over RC,G: restarted job with resume_from_cache=true RC->>D: read LATEST RC->>G: restore, then resume RC->>G: continued rows as attempt + 1, other unfinished rows from their input Note over RC: restored rows not dispatched yet stay in the next checkpoint's scope Note over RC: a row whose /run got no reply is retired, between checkpointsWhere 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 collection (this PR)ananthsub/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 workersRelationship to the old stack (#2939 to #2946)
New. The old stack left the controller to the training framework and had no evaluation path.
Issue
No tracking issue exists. This came out of design discussion of the stack: how a single long evaluation task would be checkpointed without a training controller.
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_collection.py tests/unit_tests/test_rollout_collection.py tests/unit_tests/test_cli_eval.py tests/unit_tests/test_cli_eval_submit.py: 425 passed.RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_checkpoint_*.py: 151 passed.test_checkpoint_collection.pycovers the controller with coordination faked:/rungot no reply are retired, and rows that replied are not;checkpoint_dir.SIGTERM, andSIGUSR1) pass in the full e2e suite at the top of the stack: 43 passed, 9 skipped.tests/unit_tests, eight processes, withouttest_opensandbox_compose.py, which needs the optionalopensandboxpackage): 6,572 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
The process-level tests above run real Gym server processes against the fake inference backend. A real-model preemption run is pending.
Compatibility
checkpoint_dir(default unset);checkpoint_every_s(needscheckpoint_dir);checkpoint_timeout_s(default 300).checkpoint_dirset, rollout collection installsSIGTERMandSIGUSR1handlers for the length of the run.