Skip to content

feat(token-capture): add worker staging helpers for generation cuts - #3884

Draft
ananthsub wants to merge 1 commit into
ananthsub/partial-ckpt-policy-modelfrom
ananthsub/partial-ckpt-model-worker-cuts
Draft

ananthsub wants to merge 1 commit into
ananthsub/partial-ckpt-policy-modelfrom
ananthsub/partial-ckpt-model-worker-cuts

Conversation

@ananthsub

@ananthsub ananthsub commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What changed and why

An inference worker that serves generation cuts needs three helpers from RolloutTokenCapture. This PR adds them:

  • begin_call(generation_cut=..., generation_cut_staging_keys=...) validates the durable prefix a replacement call continues against its admission: the staging keys, source call, digest, lineage, generated-token count, and that the prefix's policy version is not newer than the current one.
  • build_prefix_record snapshots a live call without completing it, so a cut can stage the prefix while the request keeps decoding. For a continued call it produces the old prefix plus the new tail, stamped with the oldest contributing policy version.
  • build_generation_chunk_record stages only newly generated tokens, for workers that flush a long generation in chunks.

complete_call now builds its record through build_prefix_record and claims completion only after the record is valid.

Why: the policy model PR sends cut requests and attaches restored cuts, but a real worker (for example NeMo RL's vLLM worker) could not serve them, because these helpers existed only on the old prefix-recovery branch.

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 W 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

What an inference worker does with these helpers when the policy model asks for a cut, and when a replacement call continues it.

sequenceDiagram
  participant P as Policy model server
  participant W as Inference worker
  participant S as Token staging
  P->>W: POST /ng-control/v1/generation-cut with the in-flight calls
  W->>S: build_prefix_record stages the tokens so far, the call keeps decoding
  W-->>P: receipt - staging keys and digest per call
  Note over P,W: after a crash and restore
  P->>W: replacement call carrying the continuation
  W->>S: begin_call checks keys, source call, digest, lineage, token count, policy version
  W->>W: decode only the rest of the generation
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
  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 (this PR)
  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)

None of #2939 to #2946 had these helpers. They are ported from #3412 (prefix token-level recovery), which was built on the old stack; the wire contract is unchanged.

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: 545 passed.
  • The full core unit suite (tests/unit_tests, eight processes, without test_opensandbox_compose.py, which needs the optional opensandbox package): 6,507 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 end-to-end suite in the top PR uses these helpers in its fake inference worker, so the cut record Gym restores is validated exactly as a real worker validates it (generation-cut scenario).
  • A real vLLM worker serving generation cuts: pending.

Compatibility

Additive keyword arguments and new methods on RolloutTokenCapture; existing callers are unaffected. complete_call produces the same records 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-model-worker-cuts branch from 55f0706 to f99bcb5 Compare October 1, 2026 19:38
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-model-worker-cuts branch from f99bcb5 to 3409be3 Compare October 1, 2026 21:53
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-model-worker-cuts branch from 3409be3 to 78bf698 Compare October 1, 2026 22:05
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-model-worker-cuts branch 2 times, most recently from 4a6c0f8 to 8c661c4 Compare October 2, 2026 19:43
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-model-worker-cuts branch from 8c661c4 to b6a8734 Compare October 2, 2026 21:17
Ported from the prefix-recovery work (Gym #3412) onto the re-cut policy
model participant. An inference worker that serves generation cuts needs
three things from RolloutTokenCapture:

- begin_call(generation_cut=, generation_cut_staging_keys=) validates the
  durable prefix a replacement call continues against its admission: the
  staging keys, source call, digest, lineage, generated-token count, and
  that the prefix's policy version is not newer than the current one.
- build_prefix_record snapshots a live call without completing it, so a
  cut can stage the prefix while the request keeps decoding. For a
  continued call it produces the old prefix plus the new tail, stamped
  with the oldest contributing policy version.
- build_generation_chunk_record stages only newly generated tokens, for
  workers that flush a long generation in chunks.

complete_call now builds its record through build_prefix_record and
claims completion only after the record is valid. The wire contract
(GenerationCutContinuation on CaptureAdmission) was already in place.

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-model-worker-cuts branch from b6a8734 to a432062 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