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
ananthsub
force-pushed
the
ananthsub/partial-ckpt-policy-model
branch
2 times, most recently
from
October 1, 2026 21:53
bf327a4 to
aeda6de
Compare
This was referenced Oct 1, 2026
ananthsub
force-pushed
the
ananthsub/partial-ckpt-policy-model
branch
from
October 1, 2026 22:05
aeda6de to
efc0332
Compare
ananthsub
force-pushed
the
ananthsub/partial-ckpt-policy-model
branch
2 times, most recently
from
October 2, 2026 21:17
1a3143b to
461da5b
Compare
A policy model server (checkpoint_policy: true) takes part in checkpoints. The checkpoint does not wait for generations: - Closing admission holds every response that has not been delivered, and new policy calls wait until resume. Only a response already streaming blocks readiness. - An undelivered call is left out of the checkpoint, including its capture-ledger row, so after a restore the agent simply calls again. - With checkpoint_generation_cuts, prepare asks each generation backend to cut its in-flight calls at a durable token prefix. The replacement call carries the cut in its capture admission and the backend continues that prefix instead of regenerating it. - Ledger rows are restored under the replacement attempt's capture key, so parent resolution continues the lineage by normal fingerprint matching. Every model server binds its capture context to a policy ticket, and the restored-cut hook is per call, so two model apps in one process do not share state. A restored generation cut that no re-issued call has consumed yet is exported again by the next checkpoint, so a second crash still continues the prefix. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
A policy model server may now run num_workers > 1 under checkpointing. Its workers share one port, so a control call reaches an arbitrary worker; each running its own participant would close admission on one worker, export other workers' undelivered calls as delivered, and restore cuts into one worker only. It is still one participant to the controller: - A coordinator in the main process, which serves no requests, owns the controller (phases, fencing, lease, storage), the attempt fence, the restored generation cuts, and the ledger export and import through the process-shared capture ledger. - Each worker runs a PolicyGate, the per-process data plane split out of the participant: admission, held responses, and its own generation-cut round. It forwards checkpoint control calls to the coordinator, closes, reopens, and retires when told to, and while a checkpoint is open reports which responses are streaming, which calls are undelivered, and their cuts. A re-issued call claims its restored cut from the coordinator before it runs, and gives it back if it is the wrong call. - Workers reach the coordinator over a Unix socket with length-prefixed JSON frames and request and reply messages in both directions. Failures close the checkpoint instead of weakening it: a worker that has not reported for the current checkpoint, or that disconnected while it was open, blocks prepare and commit until the controller retires or resumes. A worker that starts or restarts during a checkpoint closes at once and adopts the current fences and restored cuts. A worker that loses the coordinator shuts itself down: the main process, and uvicorn's supervisor with it, is gone, and a worker left behind would hold the port against the server's restart. With one worker nothing changes: the participant owns its gate in process. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…ge cases Found by the scale test and the review of the old stack's follow-ups. - Commit and restore read and write each continued episode's ledger on a worker thread instead of the event loop. A restore installs all ledgers without syncing each file (the checkpoint is the durable copy, and importing it again is a no-op) and syncs the directory once, after one directory listing instead of an existence read per episode. At 12,000 continued rollouts these were the checkpoint costs that grew with scale. - The commit reply lists staging_keys: every staged token-store key the checkpointed episodes still refer to, from ledger rows and generation cuts. The training framework keeps those rows and may clear the rest. - A call whose decoding is constrained (tool_choice other than auto or none, a non-text response_format, guided or structured outputs) is not cut and never continues a restored cut: a token prefix cannot restore its structured decoder. It regenerates instead. Ported from the prefix recovery work (Gym #3412); vLLM reports each request to the policy gate. - Generation-cut requests to inference workers use the reserved control connection pool, since the worker's data connections are busy with the very generations being cut. - Readiness is cheap to recompute: it counts held, cut, failed, and skipped calls, and each cut record is built once, when the worker acknowledges the cut. The full list of undelivered calls and their cuts is taken once, at commit; with several workers the coordinator asks each worker for it then. At 16k rollouts, rebuilding that list on every readiness check and every report stretched the policy prepare stage from 2 s to 25 s. - An episode keeps one cut: when a client retried a call the server was still running, both calls are undelivered and cut, and only the latest admitted one is still awaited. Exporting both made restore reject the episode at 16k rollouts. - Readiness counts cut_failed and cut_skipped, and a worker that fails to cut some calls is logged, so cuts that silently degrade to regeneration are visible. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…e calls Restore installs each checkpointed episode's ledger rows as the ledger of the next attempt. The capture directory survives a Gym restart, so after one restore, a model call by the replacement attempt, and a second crash, restoring the same checkpoint found the replacement's ledger, refused it, and failed the whole restore. Every rollout then restarted from its input. The import now deletes what dead executions left for the target attempt and for later attempts of the same rollout, using the ledger's and the complete-record store's delete: ledgers, token records, intents, incomplete markers, capture state, and retire fences. A fence left in place would silently discard the restored episode's rows once it reached that attempt again. A model server that has served any call refuses to restore, since a ledger it wrote may belong to a live episode. This replaces tracking every ledger the process appended to, which grew without bound. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
The training framework now retires and deletes capture ledgers through the model server's control routes. Two interactions with checkpoints needed care: - A commit reads the ledger of every episode it continues. A retire or delete that landed between prepare and commit would leave the checkpoint holding the episode with no lineage, with no error. While a checkpoint is open on the model server, the routes refuse with 409 and the caller retries after resume. - A checkpoint retire of attempt N fences attempt N and every earlier attempt. The participant now also retires those attempts' ledgers, after their calls are cancelled, so their files are freed and no late row recreates them. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…sed restored cuts - A restore continues attempt N as attempt N+1, but left attempt N's ledger on disk with nothing to remove it unless the training framework retired that key. The restore now retires the ledgers of every attempt it continues, and earlier attempts, in one batch. - A restored generation cut waits for the replacement's re-issued call. The model participant now reports unused restored cuts as pending, so a commit that no longer continues their episode retires them. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…s retire replies The gate's retire now waits until the cancelled calls have exited, and with several workers each worker refuses the attempt only while its own calls stop. Workers no longer receive a copy of a lasting fence when they join or reopen. 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-policy-model
branch
from
October 2, 2026 23:32
461da5b to
3e94e72
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
A policy model server with
checkpoint_policy: truebecomes a checkpoint participant. The checkpoint never waits for generations.checkpoint_generation_cuts: true). Prepare asks each inference worker to cut its in-flight calls at a durable token prefix. The replacement call carries the cut in its capture admission, and the worker continues the prefix instead of regenerating it. The wire models (GenerationCutInventory,GenerationCutReceipt,GenerationCutPrefixAck) keep the field layout NeMo RL already uses. A call with constrained decoding (tool_choiceother thanautoornone, a non-textresponse_format, guided or structured outputs) is never cut and never continues a cut, because a token prefix cannot restore its structured decoder.num_workers > 1. It is still one participant to the controller: a coordinator in the main process owns the phases, lease, storage, and restored cuts, and each worker runs a gate (admission, held responses, its own cut round) that reports to the coordinator over a Unix socket. A worker that has not reported, or that disconnects during a checkpoint, blocks the checkpoint instead of weakening it. With one worker nothing changes.staging_keys, the staged token-store keys the checkpoint still refers to, so the training framework knows which rows to keep. Cut requests use the reserved control connection pool, and cut failures are counted in readiness and logged.deleteoperations of feat(token-capture): retire and delete external-staging capture ledgers #3938 and feat(token-capture): retire and delete complete-record captures #3939: ledgers, token records, intents, incomplete markers, capture state, and retire fences.retireof attempt N cancels the calls of attempt N and every earlier attempt and waits for them, then retires those attempts' ledgers, so their files are freed and no late row recreates them. With several workers, each worker refuses the attempt only while its own calls stop.Conflict note:
CaptureContextonmaingainedexternal_worker_response_seen(#3670); this PR addsadmission_hooknext to it, and both fields are kept.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 M thisOne checkpoint, a crash, and the restore, end to end:
This PR
A model call across a checkpoint and a crash.
With several uvicorn workers, one coordinator in the main process is the participant; every worker keeps its own gate and forwards control calls to the coordinator.
The capture ledger across checkpoints, restores and the training framework's cleanup.
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 cuts (this PR)ananthsub/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)
main's capture ledger plus rows re-imported under the next attempt; admission no longer waits for generations.nemo_gym._checkpoint.generation_cutinstead ofmodel_control_contracts.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: 86 passed.RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_token_capture_*.py tests/unit_tests/test_token_id_capture.py tests/unit_tests/test_base_responses_api_model.py: 541 passed.RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider responses_api_models/vllm_model/tests: 328 passed.capture ledger for roll-1-a1 already holds rows from another execution), with the files and fences of the target and later attempts deleted and look-alike rollouts untouched;tests/unit_tests, eight processes, withouttest_opensandbox_compose.py, which needs the optionalopensandboxpackage): 6,503 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
Real-model rollout evidence: pending. The top PR of the stack has an opt-in real vLLM test that has not been run on this re-cut yet.
Scale, on one workstation with every Gym server and a fake inference backend (native
single_agent_turnweather rollouts, token capture and generation cuts on, each rollout's second model call held so most rollouts are in flight at the checkpoint; checkpoint, kill every Gym process group, restart, restore, run the replacement attempts). After the scale fixes in this PR, measured on the branch this stack was cut from, with the same commits:In the first 16,000-rollout run, the policy model logged 8 uvloop "File descriptor N is used by transport" errors and 16 requests failed with "ledger lineage resolution failed". The cause is a libuv accept race while clients disconnect, not checkpoint code; the ledger failed closed, clients retried, and no replacement failed. The rerun and every 8,000-rollout run had none.
Compatibility
checkpoint_policy(defaultfalse) makes the server a participant, andcheckpoint_generation_cuts(defaultfalse) turns on generation cuts. Both only take effect with the globalcheckpoint:block enabled.token_id_capture.generation_prefix_cuts_enabledkey (never onmain) switch tocheckpoint_generation_cuts.nemo_gym._checkpoint.generation_cut.