Skip to content

feat(checkpoint): make policy model servers checkpoint participants with generation cuts - #3883

Draft
ananthsub wants to merge 7 commits into
ananthsub/partial-ckpt-corefrom
ananthsub/partial-ckpt-policy-model
Draft

ananthsub wants to merge 7 commits into
ananthsub/partial-ckpt-corefrom
ananthsub/partial-ckpt-policy-model

Conversation

@ananthsub

@ananthsub ananthsub commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What changed and why

A policy model server with checkpoint_policy: true becomes a checkpoint participant. The checkpoint never waits for generations.

  • Held responses. Closing admission holds every response that has not been delivered yet, and new policy calls wait until resume. Only a response that is already streaming blocks readiness.
  • Undelivered calls are left out. An undelivered call is excluded from the checkpoint, together with its capture-ledger rows, so after a restore the agent simply makes the call again.
  • Generation cuts (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_choice other than auto or none, a non-text response_format, guided or structured outputs) is never cut and never continues a cut, because a token prefix cannot restore its structured decoder.
  • Token lineage. Ledger rows are restored under the replacement attempt's capture key, so parent resolution continues the lineage by normal fingerprint matching, with no parent header relay.
  • Several uvicorn workers. A policy model server may run 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.
  • Scale fixes. Commit and restore read and write ledgers on a worker thread. A restore installs ledgers without syncing each file (the checkpoint is the durable copy, and importing again is a no-op), with one directory listing and one directory sync: importing 15,686 episodes drops from 3.2 s to 1.5 s. Readiness is cheap to recompute. An episode keeps only its latest cut. The commit reply lists 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.
  • A restored cut that no re-issued call has consumed yet is exported again by the next checkpoint while the controller's commit scope still continues its episode, so a second crash still continues the prefix. Once the scope leaves the episode out, the cut is released.
  • Restoring the same checkpoint again. The capture directory survives a Gym restart. If the replacement attempt makes model calls and Gym crashes before the next checkpoint, restoring the same checkpoint again finds that attempt's files.
    • The import first deletes what dead executions left for the target attempt and for every later attempt of the same rollout. It uses the delete operations 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.
    • A fence left in place would silently discard the restored episode's rows once the episode reached that attempt again. Attempt numbers after a restore are deterministic, and the training framework now retires every finished rollout, so fences are common.
    • The restore then retires the ledgers of the attempts it continues, and of earlier attempts, in one batch. Attempt N continues as attempt N+1, so attempt N's ledger would otherwise stay on disk.
    • A model server that has served any call refuses to restore, because a ledger it wrote may belong to a live episode. With several workers, each worker reports this to the coordinator. This replaces tracking every ledger the process appended to, which grew without bound.
  • Capture-ledger cleanup stays consistent with checkpoints. The training framework retires and deletes ledgers through the model server's control routes (feat(token-capture): retire and delete external-staging capture ledgers #3938).
    • While a checkpoint is open on the model server, those routes refuse with 409, and the caller retries after resume. A commit reads the ledger of every episode it continues, so a retire landing between prepare and commit would leave the checkpoint holding an episode with no lineage, with no error.
    • A checkpoint retire of 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: CaptureContext on main gained external_worker_response_seen (#3670); this PR adds admission_hook next 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 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

A model call across a checkpoint and a crash.

sequenceDiagram
  participant A as Agent
  participant P as Policy model server
  participant W as Inference worker
  A->>P: model call for rollout r, attempt 0
  P->>W: generate
  Note over P: prepare closes admission - new calls park,<br/>responses not yet delivered are held
  P->>W: generation cut, if enabled - stage the prefix so far
  W-->>P: receipt with staging keys per call
  Note over P: commit exports ledger rows minus undelivered calls,<br/>plus the latest cut of each episode
  Note over A,W: crash, restore, resume
  A->>P: the same call again, as attempt 1
  P->>W: continue from the staged prefix, or generate again
Loading

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.

flowchart LR
  CO["Coordination"] -->|"control call on the shared port"| W1["uvicorn worker 1<br/>gate: tickets, held responses"]
  CO -->|"control call on the shared port"| W2["uvicorn worker 2<br/>gate: tickets, held responses"]
  W1 <-->|"Unix socket: forward, close, reports, cut claims"| C["Coordinator in the main process<br/>phases, lease, ledger, restored cuts"]
  W2 <-->|"Unix socket"| C
Loading

The capture ledger across checkpoints, restores and the training framework's cleanup.

sequenceDiagram
  participant T as Training framework
  participant P as Policy model server
  participant L as Capture ledger
  Note over P: checkpoint open - ledger retire and delete routes answer 409
  T->>P: retry retire and delete after resume
  Note over P,L: restore - delete what dead executions left for attempt N+1 and later, fences included
  P->>L: import attempt N's rows as attempt N+1, then retire attempts 0 to N
  Note over P,L: checkpoint retire of attempt N - stop its calls, then retire the ledgers of attempts 0 to N
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
  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 (this PR)
  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)

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.
  • Among the checkpoint tests, each failing without its change:
    • restoring the same checkpoint again after the replacement made a call (without the fix: 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;
    • a restore retires the ledgers of the attempts it continues;
    • a model server that has served calls refuses to restore, with one process and with two workers;
    • ledger retire and delete return 409 while a checkpoint is open and succeed after resume;
    • a retire stops a running generation before it replies, and retires the ledgers of attempts 0 to N but not N+1;
    • a retire leaves no fence on any worker, including one that joins later.
  • The full core unit suite (tests/unit_tests, eight processes, without test_opensandbox_compose.py, which needs the optional opensandbox package): 6,503 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

  • 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_turn weather 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:

    Rollouts Policy workers In flight at checkpoint (exported) Prepare (s) Commit (s) Restore (s) Replacement failures
    8,000 1 4,201 0.67 0.60 0.60 0
    8,000 2 4,083 0.56 0.69 0.56 0
    16,000 2 14,320 7.57 (model 7.38) 4.64 4.97 0
    16,000 2 (rerun) 12,228 3.60 (model 3.53) 1.83 3.23 0

    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

  • New model server config: checkpoint_policy (default false) makes the server a participant, and checkpoint_generation_cuts (default false) turns on generation cuts. Both only take effect with the global checkpoint: block enabled.
  • Users of the old stack's token_id_capture.generation_prefix_cuts_enabled key (never on main) switch to checkpoint_generation_cuts.
  • Generation-cut models are imported from nemo_gym._checkpoint.generation_cut.
  • With checkpointing off, model servers behave as before.

@ananthsub ananthsub added feature New capabilities, enhancements, or enablement work area:model Model servers, inference providers, and model adapters 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-policy-model branch 2 times, most recently from bf327a4 to aeda6de Compare October 1, 2026 21:53
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-policy-model branch from aeda6de to efc0332 Compare October 1, 2026 22:05
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-policy-model branch 2 times, most recently from 1a3143b to 461da5b Compare October 2, 2026 21:17
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
pthombre added this pull request to stack #3962 October 2, 2026 22:50
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-policy-model branch from 461da5b to 3e94e72 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:model Model servers, inference providers, and model adapters feature New capabilities, enhancements, or enablement work

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant