Skip to content

feat(checkpoint): add the participant control plane, episode steps, and coordination - #3882

Draft
ananthsub wants to merge 3 commits into
ananthsub/token-capture-store-retirefrom
ananthsub/partial-ckpt-core
Draft

ananthsub wants to merge 3 commits into
ananthsub/token-capture-store-retirefrom
ananthsub/partial-ckpt-core

Conversation

@ananthsub

@ananthsub ananthsub commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What changed and why

This PR adds the core that every checkpoint participant builds on. On its own it installs nothing on any server; the later PRs in the stack each add one participant on top of it.

  • Identity. An episode attempt is keyed by its capture key: r for the first attempt and r-aN after a restore. When checkpointing is on, ServerClient always attributes model and resources calls to their rollout through the rollout path prefix, not only when observability is on. main already reserves the -a<N> suffix in rollout IDs; this PR adds EpisodeId.from_capture_key, which parses a key back to its rollout and attempt.
  • Control plane (nemo_gym/_checkpoint/control.py). Each server has one participant behind bearer-protected routes under /ng-control/v1/checkpoint/ (status, prepare, renew, retire, commit, restore, resume). One controller per server owns the phase machine, checkpoint ID fencing, the readiness wait, deadline-bounded file I/O, and manifest-last storage. A lease resumes a participant whose controller stopped calling. Resumed checkpoint IDs are remembered for the last 64 checkpoints, enough to refuse a stale controller's delayed request. Status reports how many attempts a retire is stopping right now (retiring).
  • Episode steps (steps.py). An episode records a boundary after each completed step. A step in flight when a checkpoint closes runs in wait mode (the checkpoint waits for it) or replay mode (it runs again after a crash). Time parked for a checkpoint does not count against the episode's deadline. A resources server reports whether its /verify may be replayed in an x-ng-checkpoint-verify header on its /seed_session reply.
  • Coordination (coordination.py). Gym owns the participant order, so a training controller calls these functions instead of sequencing servers itself: prepare environment servers, then the policy model, then agents, then resources servers; resume in reverse; commit everywhere; restore all or nothing; retire callers before the servers they call (environment servers, then agents, then the policy model and resources servers). renew extends every lease when publishing a checkpoint takes longer than expected.
  • Retire stops and frees. A participant's retire replies only after the running work of the attempt, and of its earlier attempts, has stopped and its state is freed. While it stops, the attempt is refused; afterwards nothing about it remains, so no table grows with the job's history.
    • Retire is refused while a checkpoint is preparing, prepared or committed: stopping an episode waits for its final cleanup, whose calls wait for resume. A controller with a straggler resumes, retires, then checkpoints again.
    • Retire stays allowed right after a restore, where nothing is running yet; a failed restore uses it to discard what it installed.
    • The controller keeps its own side paused from prepare to resume. For NeMo RL that is rollout dispatch, the trainer actors, and the trajectory assembler actors.
  • A commit releases restored state the controller no longer continues. A commit's scope names every episode the controller continues, including a restored episode whose replacement has not started yet, named by that replacement attempt. After writing its records, a participant retires the restored episodes the scope leaves out, so unclaimed restored state is never carried from checkpoint to checkpoint.
  • Reserved control connections. Checkpoint and other control calls use a separate aiohttp connection pool, so data calls that hold every connection to a host (for example, /run calls parked during a checkpoint) cannot starve prepare, commit, or resume. The control pool is built like the data pool on main, with the same keepalive socket options and, when metrics export is on, the same timed connector.
  • Faithful, cheap storage.
    • Each participant's records file keeps every dict's key order, NaN and infinities exactly as exported, so a restored episode sees byte-for-byte the state it had. Tool outputs often serialize stored dicts. A test checks the round trip.
    • Opaque server state in a record (JsonPayload) is checked once with orjson for JSON-ness, then written as is, rather than deep-validated and copied. For 1,000 sessions of 350 KB, commit drops from 11.5 s to 3.6 s and restore from 6.0 s to 3.4 s.
  • Settings. A global checkpoint: block with enabled, control_auth_token, and lease_grace_seconds. It is off by default, and a control token is required when it is on.

Why: the old stack spread identity across headers and path prefixes, gave each participant kind its own control contracts, and left the participant order to the training controller. Environment servers now exist on main and own the episode protocol, which makes a smaller and fail-closed design possible.

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 CO,K 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

Each participant's phases, enforced by the control plane. A second checkpoint ID while one is open gets checkpoint_conflict. Retire is accepted only in idle and restored.

stateDiagram-v2
  [*] --> idle
  idle --> preparing: prepare closes admission
  preparing --> prepared: every live execution is at a boundary
  prepared --> committed: commit writes records, then the manifest
  committed --> idle: resume reopens admission
  preparing --> idle: resume aborts, or the lease expires
  prepared --> idle: resume aborts, or the lease expires
  idle --> restored: restore, in a fresh process only
  restored --> idle: resume
Loading

Episode steps: what happens to a step that is running when a checkpoint closes.

flowchart TD
  S["A step is in flight when the checkpoint closes"] --> Q{"Step mode"}
  Q -->|wait| W["The checkpoint waits for the step.<br/>The episode parks at the boundary after it."]
  Q -->|replay| P["The checkpoint does not wait.<br/>The boundary before the step is exported,<br/>and the step runs again after a crash."]
Loading

Retire stops and frees. Each server refuses the attempt only while it stops it; when retire returns, nothing of the attempt is running or held anywhere, so nothing needs to be remembered.

sequenceDiagram
  participant C as Controller
  participant E as Environment servers
  participant A as Agents
  participant MR as Model and resources servers
  C->>E: retire(r, N)
  Note over E: refuse r up to N, cancel the episode, wait for it and its final cleanup
  E-->>C: stopped and freed
  C->>A: retire(r, N)
  Note over A: cancel activations and wait, then the retire hook
  A-->>C: stopped and freed
  C->>MR: retire(r, N)
  Note over MR: cancel calls and wait, drain session requests, retire ledgers and hooks
  MR-->>C: stopped and freed
Loading

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 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 (this PR)
  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
  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

Partial-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)

  • Supersedes feat(checkpoint): propagate stable execution identity #2939 (execution identity). Only the rollout path prefix remains; the x-nemo-gym-rollout-id, x-nemo-gym-attempt-index, and x-nemo-gym-model-call-id headers are gone.
  • Supersedes the control-plane part of feat(checkpoint): add fenced control contracts #2940 (control contracts). The per-kind request models and capability endpoint become one route set and one request shape for every participant kind. The generation-cut models move to the policy model PR. The multi-worker admission coordinator is not carried; multi-worker policy model servers return in a different, smaller form in the policy model PR.
  • New compared with the old stack: episode steps and Gym-owned coordination. The completed-result acknowledgements are gone: no episode completes while Gym is prepared, so there is nothing to acknowledge.

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 (which is on the current main):

  • RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_checkpoint_*.py: 43 passed. New since the first version of this PR:
    • a retire replies only after the participant stopped the attempt, and refuses the attempt meanwhile;
    • a straggler is retired after resume, not while a checkpoint is open;
    • coordination retires environment servers, then agents, then the model and resources servers;
    • a commit retires restored episodes its scope leaves out;
    • the attempt fence holds an attempt only while it is being stopped.
      Each fails without its change.
  • The full core unit suite (tests/unit_tests, eight processes, without test_opensandbox_compose.py, which needs the optional opensandbox package): 6,460 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

N/A for this PR on its own: it adds no participant, so no server behaves differently. The process-level end-to-end suite and the scale results in the top PR of the stack exercise this core through every participant.

Compatibility

  • New global config: a checkpoint: block with enabled (default false), control_auth_token (required when enabled), and lease_grace_seconds (default 600).
  • New config for the reserved control connection pool: global_aiohttp_control_connector_limit and global_aiohttp_control_connector_limit_per_host.
  • Nothing changes unless the global checkpoint: block sets enabled: true. With it on, model and resources calls always carry the rollout path prefix.

@ananthsub ananthsub added feature New capabilities, enhancements, or enablement work area:core Shared APIs, servers, telemetry, health, and registries 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-core branch from 0a24eb7 to c2b5103 Compare October 1, 2026 19:38
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-core branch from c2b5103 to 9040654 Compare October 1, 2026 21:53
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-core branch from 9040654 to 846c721 Compare October 1, 2026 22:05
…ination

The core of partial-rollout checkpointing, shared by every server. It
installs nothing on its own; the model, environment, agent, and
resources changes each add a participant on top of it.

- Identity: an episode attempt is keyed by its capture key ("r" or
  "r-aN"). ServerClient attributes model and resources calls to their
  rollout whenever checkpointing is on, not only with observability.
- Control plane (control.py): one phase machine behind bearer-protected
  routes under /ng-control/v1/checkpoint, driving one participant per
  server. The controller owns checkpoint-ID fencing, attempt fencing,
  the readiness wait, deadline-bounded I/O, manifest-last storage, and a
  lease that resumes a participant whose controller went away.
- Episode steps (steps.py): a boundary after each completed step, and a
  wait or replay mode for a step in flight when a checkpoint closes.
  Time parked for a checkpoint does not count against an episode's
  deadline.
  A resources server reports whether its /verify may be replayed in an
  x-ng-checkpoint-verify header on its /seed_session reply, so callers
  learn the mode without another request; the default is wait.
- Coordination (coordination.py): the Gym-owned participant order a
  training controller calls: prepare environment servers, the policy
  model, agents, then resources servers; resume in reverse; commit and
  retire everywhere; restore all or nothing.
- Settings: a global checkpoint block, off by default, that requires a
  control token when enabled.

- Checkpoint coordination and other control calls use a reserved aiohttp
  pool (global_aiohttp_control_connector_limit, and _per_host), so data
  calls that hold every connection to a host, such as parked /run calls,
  cannot starve prepare, commit, or resume. coordination.renew extends
  every participant's lease when publishing outlasts it.
- Participants may override async export and install hooks to keep slow
  I/O off the event loop, and may add fields to the commit reply.
- An episode woken by a resume stays parked if a new checkpoint closed
  admission before it ran.
- Records keep each participant's opaque state (an environment, a session,
  a boundary) exactly as exported: key order, NaN, and infinities round-trip.
  Payload fields are checked once with orjson for JSON-ness, then written as
  they are, instead of a deep validation and copy of large states.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-core branch from 846c721 to f4092dc Compare October 2, 2026 19:43
@ananthsub
ananthsub changed the base branch from main to ananthsub/token-capture-store-retire October 2, 2026 19:44
…no longer continues

A restore installs state for every episode the controller continues, and
that state waits for the replacement attempt to claim it. If the controller
never re-dispatches an episode, nothing released its restored state, and
every later checkpoint exported it again.

A commit's scope is the controller's statement of what it continues. After
writing its records, a participant now retires every restored episode whose
replacement has not started and that the scope leaves out. Participants
report those episodes through a new restored_pending method; retire already
releases everything else.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…g it forever

A retire used to discard an attempt and then remember it for the life of the
process, so its late requests were refused. That record grew with every
rollout retired or restored.

A retire now replies only after the participant has stopped the attempt's
running work and freed its state. While it stops, the attempt is refused;
afterwards nothing remains of it, so the fence holds only retires in
progress. Coordination retires callers before the servers they call
(environment servers, then agents, then the model and resources servers), so
nothing upstream can still call a server for an attempt it has stopped.

A restore and a commit's scope release no longer raise fences: nothing they
replace is running.

Retire is refused while a checkpoint is preparing, prepared or committed.
Stopping an episode waits for its final cleanup, whose calls wait for
resume, so a retire during a checkpoint could deadlock. A controller with a
straggler resumes first. Rollout collection already does.

Resumed checkpoint IDs are kept for the last 64 checkpoints instead of all
of them.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
@pthombre
pthombre added this pull request to stack #3962 October 2, 2026 22:50

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:core Shared APIs, servers, telemetry, health, and registries feature New capabilities, enhancements, or enablement work

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant