Skip to content

feat(checkpoint): checkpoint evaluation runs from rollout collection - #3893

Draft
ananthsub wants to merge 2 commits into
ananthsub/partial-ckpt-e2efrom
ananthsub/partial-ckpt-rollout-collection
Draft

ananthsub wants to merge 2 commits into
ananthsub/partial-ckpt-e2efrom
ananthsub/partial-ckpt-rollout-collection

Conversation

@ananthsub

@ananthsub ananthsub commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

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_dir set (and the global checkpoint: block enabled):

  • Triggers.

    • SIGUSR1 checkpoints the rows in flight and continues.
    • SIGTERM checkpoints them and stops the run. Slurm can send it ahead of a preemption, for example with --signal=B:TERM@120.
    • checkpoint_every_s checkpoints on a timer, which covers crashes that give no warning.
    • The collector writes its process ID to checkpoint_dir/collector.pid, for sending the signals.
  • A checkpoint:

    1. stops starting new rows;
    2. prepares every participant;
    3. commits the rows whose /run has not replied;
    4. writes the run's manifest, then publishes the checkpoint by pointing checkpoint_dir/LATEST at it;
    5. resumes, unless it is stopping.

    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=true restores the latest checkpoint before dispatching anything.

    • Its unfinished rows continue as their next attempt.
    • Rows that never started, or that were not in flight at the checkpoint, start from their input as before.
    • Rows that finished are already in the output file, which rollout collection writes as each row completes.
    • If the restore fails, coordination has already retired that attempt everywhere, so those rows start from their input one attempt later.
    • A fresh run (resume_from_cache=false) forgets the checkpoints of an earlier run in the same directory, as it already does with the output file.
  • Code.

    • The new nemo_gym/_checkpoint/collection.py holds the controller.
    • In 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 eval and the collection CLI print a message and exit cleanly when a preemption signal stopped the run.
  • Cleanup.

    • A row whose /run got 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.
    • Rows restored from a checkpoint but not dispatched yet are continued by the next checkpoint, at the attempt they will run as. Without this, a crash between a restore and their dispatch restarted them from their input.
    • Publishing removes partial checkpoint directories left by a failed commit and temporary files left by a crash.

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 this
Loading

One checkpoint, a crash, and the restore, end to end:

sequenceDiagram
  participant C as Controller
  participant G as Gym participants
  participant D as Checkpoint directory
  C->>G: prepare, in order environment, model, agent, resources
  Note over G: admission closes, in-flight work parks at a boundary,<br/>undelivered model responses are held
  G-->>C: prepared, or blockers at the deadline
  C->>G: commit with the episodes the controller continues
  G->>D: each participant writes its records, then its manifest
  Note over G: restored state the commit's scope leaves out is released
  C->>C: publish the checkpoint with the controller's own state
  C->>G: resume, in order resources, agent, model, environment
  Note over C,G: crash - every Gym process dies
  C->>G: restore the checkpoint in fresh processes, all or nothing
  D-->>G: records installed under attempt + 1, attempt N's capture ledger retired
  C->>G: resume
  C->>G: /run as attempt + 1 continues each episode from its boundary
  Note over C,G: dropping an episode, only while no checkpoint is open
  C->>G: retire - environment, then agent, then model and resources
  Note over G: each server stops the attempt's work, waits, frees its state, then replies
Loading

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 checkpoints
Loading

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 retire and delete) and #3939 (complete-record retire and delete): #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.

  1. feat(checkpoint): add the participant control plane, episode steps, and coordination #3882 (ananthsub/partial-ckpt-core): feat(checkpoint): add the participant control plane, episode steps, and coordination
  2. feat(checkpoint): make policy model servers checkpoint participants with generation cuts #3883 (ananthsub/partial-ckpt-policy-model): feat(checkpoint): make policy model servers checkpoint participants with generation cuts
  3. feat(token-capture): add worker staging helpers for generation cuts #3884 (ananthsub/partial-ckpt-model-worker-cuts): feat(token-capture): add worker staging helpers for generation cuts
  4. feat(checkpoint): continue environment server episodes from their boundaries #3885 (ananthsub/partial-ckpt-environment): feat(checkpoint): continue environment server episodes from their boundaries
  5. feat(checkpoint): add the resources server participant with asynchronous session hooks #3886 (ananthsub/partial-ckpt-resources): feat(checkpoint): add the resources server participant with asynchronous session hooks
  6. feat(checkpoint): declare five training verifiers stateless with replayable verification #3887 (ananthsub/partial-ckpt-verifier-declarations): feat(checkpoint): declare five training verifiers stateless with replayable verification
  7. feat(checkpoint): add the agent session participant and Simple Agent continuation #3888 (ananthsub/partial-ckpt-agent): feat(checkpoint): add the agent session participant and Simple Agent continuation
  8. test(checkpoint): add a process-level end-to-end suite driven by coordination #3889 (ananthsub/partial-ckpt-e2e): test(checkpoint): add a process-level end-to-end suite driven by coordination
  9. feat(checkpoint): checkpoint evaluation runs from rollout collection #3893 (ananthsub/partial-ckpt-rollout-collection): feat(checkpoint): checkpoint evaluation runs from rollout collection (this PR)
  10. feat(checkpoint): spans and metrics for partial-rollout checkpoints #3903 (ananthsub/partial-ckpt-telemetry): feat(checkpoint): spans and metrics for partial-rollout checkpoints
  11. feat(checkpoint): resources, agent, and environment servers with several workers #3909 (ananthsub/partial-ckpt-multi-worker): feat(checkpoint): resources, agent, and environment servers with several workers

Relationship 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.py covers the controller with coordination faked:
    • a checkpoint commits exactly the rows in flight and publishes them;
    • no row starts while a checkpoint is open;
    • a checkpoint that does not prepare publishes nothing and resumes;
    • a stop does not resume;
    • only two checkpoints are kept;
    • restore continues only unfinished rows, as their next attempt, and falls back to a later attempt when it fails;
    • a fresh run forgets old checkpoints;
    • rows restored but not dispatched yet are continued by the next checkpoint;
    • rows whose /run got no reply are retired, and rows that replied are not;
    • publishing removes partial checkpoints and temporary files;
    • the timer checkpoints only while rows are in flight;
    • the timer needs checkpoint_dir.
  • The process-level rollout collection scenarios (legacy and native SIGTERM, and SIGUSR1) pass in the full e2e suite at the top of the stack: 43 passed, 9 skipped.
  • The full core unit suite (tests/unit_tests, eight processes, without test_opensandbox_compose.py, which needs the optional opensandbox package): 6,572 passed, 2 failed. The two failures, a sandbox retry test and a Slurm script test, fail the same way on main in this development environment.
  • pre-commit run --files <files changed by this PR>: all hooks passed, and no hook modified a file.
  • Every commit carries a Signed-off-by line.

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

  • New rollout collection config:
    • checkpoint_dir (default unset);
    • checkpoint_every_s (needs checkpoint_dir);
    • checkpoint_timeout_s (default 300).
  • With checkpoint_dir set, rollout collection installs SIGTERM and SIGUSR1 handlers for the length of the run.
  • Restoring needs Gym servers that have not served sessions yet, which is the case after a job restart. Gym servers that keep running across collector restarts cannot be restored into; their rows then start from their input one attempt later.

@ananthsub ananthsub added feature New capabilities, enhancements, or enablement work area:evaluation Evaluation, rollout collection, reward profiling, comparison, and exporters labels Oct 1, 2026
@copy-pr-bot

copy-pr-bot Bot commented Oct 1, 2026

Copy link
Copy Markdown

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.

@ananthsub
ananthsub requested a review from zyzhou5 October 1, 2026 19:04
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from b5ffb47 to 22a07dd Compare October 1, 2026 19:38
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from 22a07dd to 324d650 Compare October 1, 2026 21:53
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from 324d650 to 29a2feb Compare October 1, 2026 22:05
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from 29a2feb to d1bac74 Compare October 2, 2026 13:22
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from d1bac74 to 37f360b Compare October 2, 2026 19:43
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from 37f360b to 1ae4776 Compare October 2, 2026 21:17
@pthombre
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
ananthsub force-pushed the ananthsub/partial-ckpt-rollout-collection branch from 1ae4776 to a307627 Compare October 2, 2026 23:32

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

area:evaluation Evaluation, rollout collection, reward profiling, comparison, and exporters feature New capabilities, enhancements, or enablement work

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant