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-core
branch
from
October 1, 2026 19:38
0a24eb7 to
c2b5103
Compare
This was referenced Oct 1, 2026
ananthsub
force-pushed
the
ananthsub/partial-ckpt-core
branch
from
October 1, 2026 21:53
c2b5103 to
9040654
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-core
branch
from
October 1, 2026 22:05
9040654 to
846c721
Compare
…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
force-pushed
the
ananthsub/partial-ckpt-core
branch
from
October 2, 2026 19:43
846c721 to
f4092dc
Compare
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
added this pull request to stack #3962
October 2, 2026 22:50
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
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.
rfor the first attempt andr-aNafter a restore. When checkpointing is on,ServerClientalways attributes model and resources calls to their rollout through the rollout path prefix, not only when observability is on.mainalready reserves the-a<N>suffix in rollout IDs; this PR addsEpisodeId.from_capture_key, which parses a key back to its rollout and attempt.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).steps.py). An episode records a boundary after each completed step. A step in flight when a checkpoint closes runs inwaitmode (the checkpoint waits for it) orreplaymode (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/verifymay be replayed in anx-ng-checkpoint-verifyheader on its/seed_sessionreply.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).renewextends every lease when publishing a checkpoint takes longer than expected./runcalls parked during a checkpoint) cannot starve prepare, commit, or resume. The control pool is built like the data pool onmain, with the same keepalive socket options and, when metrics export is on, the same timed connector.JsonPayload) is checked once withorjsonfor 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.checkpoint:block withenabled,control_auth_token, andlease_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
mainand 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 thisOne checkpoint, a crash, and the restore, end to end:
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 inidleandrestored.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."]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.
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 coordination (this PR)ananthsub/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 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)
x-nemo-gym-rollout-id,x-nemo-gym-attempt-index, andx-nemo-gym-model-call-idheaders are gone.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:Each fails without its change.
tests/unit_tests, eight processes, withouttest_opensandbox_compose.py, which needs the optionalopensandboxpackage): 6,460 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
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
checkpoint:block withenabled(defaultfalse),control_auth_token(required when enabled), andlease_grace_seconds(default 600).global_aiohttp_control_connector_limitandglobal_aiohttp_control_connector_limit_per_host.checkpoint:block setsenabled: true. With it on, model and resources calls always carry the rollout path prefix.