Skip to content

test(checkpoint): add a process-level end-to-end suite driven by coordination - #3889

Draft
ananthsub wants to merge 3 commits into
ananthsub/partial-ckpt-agentfrom
ananthsub/partial-ckpt-e2e
Draft

ananthsub wants to merge 3 commits into
ananthsub/partial-ckpt-agentfrom
ananthsub/partial-ckpt-e2e

Conversation

@ananthsub

@ananthsub ananthsub commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What changed and why

A process-level end-to-end suite for partial-rollout checkpointing, in tests/e2e/checkpoint. It starts real Gym server processes (environment server, Simple Agent, resources server, policy model server) against a fake inference backend and token store that survive a Gym crash. The fake worker stages generation cuts with the real staging helpers and continues them through begin_call, so the cut record Gym restores is validated exactly as a real worker validates it.

Each scenario checkpoints mid-episode through the coordination functions, usually kills every Gym process group, restores into fresh processes, and runs the replacement attempt, asserting on the exact model calls made:

  • native and legacy /run episodes continue from their last boundary;
  • restored resources state (the counter example) is required for the right reward, with a negative control that skips the resources restore;
  • token lineage continues across the crash;
  • an in-flight generation is cut and its prefix continued;
  • a checkpoint without a crash parks and then releases the episode;
  • verification in flight is replayed or waited for, according to checkpoint_verify;
  • an episode that finishes during prepare is neither lost nor redone;
  • a second crash, after a checkpoint taken before the replacement attempt starts, still continues the episode;
  • the same checkpoint restores again after its replacement attempt made model calls and Gym crashed before the next checkpoint. The second run continues from the checkpoint's boundary, and none of the dead run's calls are in its token lineage;
  • a native episode retired mid-flight, with token capture on, stops everywhere: afterwards no server is stopping or tracking anything for it, the model makes no further call for it, and its capture ledger is retired.

Retire no longer leaves lasting fences, so a stale /run for a replaced attempt is not refused with stale_attempt any more. It runs a wasted episode that cannot touch the replacement (a different capture key and different sessions), and its capture rows are discarded because the restore retired that attempt's ledger. Idempotent /run, separate work, will refuse it.

The native, double-crash, token-lineage, generation-cut, and park scenarios run with the policy model in one process and with two uvicorn workers.

Two opt-in scenarios are skipped by default:

  • a real-model test (NEMO_GYM_CHECKPOINT_VLLM_URL) that checkpoints 16 concurrent rollouts against vLLM and checks the controller contract: every rollout either replied before the checkpoint or is exported, none replies while Gym is prepared, and every continued rollout completes after a crash and restore;
  • a scale test (NEMO_GYM_CHECKPOINT_SCALE=<rollouts>) that checkpoints thousands of in-flight rollouts with token capture and generation cuts, crashes, restores, checks the same contract, and reports per-stage prepare, commit, and restore times.

The whole suite is skipped unless NEMO_GYM_CHECKPOINT_E2E=1 is set, so it does not run in the default unit test job. It takes about 6 minutes on this 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
  T["Process-level e2e suite<br/>plays the controller"]
  T --> CO
  classDef this fill:#fde68a,stroke:#b45309,stroke-width:2px,color:#1f2937
  class T 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

How the suite runs real server processes and crashes them.

flowchart LR
  T["pytest process<br/>drives coordination"] -->|"start, then crash by killing every process group"| G["Gym servers, one process group each<br/>environment, policy model, agent, resources"]
  G --> B["Fake inference backend<br/>holds calls, serves generation cuts"]
  T -->|"/_ctl/hold and /_ctl/release"| B
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
  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 (this PR)
  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:

  • NEMO_GYM_CHECKPOINT_E2E=1 RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/e2e/checkpoint: all 23 collected tests ran, 20 passed and 3 skipped (the opt-in real-model and scale tests), in 5 minutes 31 seconds.
  • The restore-again scenario fails without the ledger fix in the policy model PR, at the second restore, with one and with two policy workers: capture ledger for again-1-a1 already holds rows from another execution.
  • RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_checkpoint_*.py: 136 passed.
  • The full core unit suite (tests/unit_tests, eight processes, without test_opensandbox_compose.py, which needs the optional opensandbox package): 6,557 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

  • Process-level crash and restore: the end-to-end suite above, with a fake inference backend: all scenarios pass.

  • Scale, on one workstation with every Gym server and the 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), run with NEMO_GYM_CHECKPOINT_E2E=1 NEMO_GYM_CHECKPOINT_SCALE=<N> pytest tests/e2e/checkpoint -k scale -s:

    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

    Every rollout either finished before the checkpoint or was exported, and every exported rollout completed after the restore. These numbers come from the lane branches this stack was cut from; the scale test was not rerun on this stack.

  • Real model: the opt-in vLLM test has not been run on this stack yet: pending. Not covered yet: a real vLLM worker serving generation cuts, and several hosts.

Compatibility

Tests only. The suite is skipped unless NEMO_GYM_CHECKPOINT_E2E=1 is set.

@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-e2e branch from 5f8b5af to 87e1f0c Compare October 1, 2026 19:38
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch from 87e1f0c to 207f588 Compare October 1, 2026 21:53
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch from 207f588 to afee845 Compare October 1, 2026 22:05
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch 2 times, most recently from 6a3c667 to f3ade85 Compare October 2, 2026 19:43
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch from f3ade85 to ca27fd2 Compare October 2, 2026 21:17
@pthombre
pthombre added this pull request to stack #3962 October 2, 2026 22:50
Real server processes against a fake inference backend and token store
that survive a Gym crash. The fake worker stages generation cuts with the
real staging helpers and continues them through begin_call, so the cut
record Gym restores is validated exactly as a real worker validates it. Each scenario checkpoints mid-episode through
the coordination functions, usually kills every Gym process, restores,
and runs the replacement attempt, asserting on the exact model calls:

- native and legacy /run episodes continue from their last boundary;
- restored resources state (counter) is required for the right reward,
  with a negative control that skips the resources restore;
- token lineage continues across the crash;
- an in-flight generation is cut and its prefix continued;
- a checkpoint without a crash parks and then releases the episode;
- verification in flight is replayed or waited for, per checkpoint_verify;
- an episode that finishes during prepare is neither lost nor redone;
- a second crash, after a checkpoint taken before the replacement attempt
  starts, still continues the episode.

An opt-in real-model test (NEMO_GYM_CHECKPOINT_VLLM_URL) checkpoints 16
concurrent rollouts against vLLM and checks the controller contract:
every rollout either replied before the checkpoint or is exported, none
replies while Gym is prepared, and every continued rollout completes
after a crash and restore.

The native, double-crash, token-lineage, generation-cut, and park
scenarios run with the policy model in one process and with two uvicorn
workers. A simulated crash kills each server's whole process group, as a
node failure takes a server's workers with it.

An opt-in scale scenario (NEMO_GYM_CHECKPOINT_SCALE=<rollouts>) checkpoints
thousands of in-flight rollouts with token capture and generation cuts,
crashes, restores, and checks the controller contract: every rollout
either finished before the checkpoint or was exported, and every
exported one completes after the restore. It reports per-stage prepare,
commit, and restore times.

Skipped unless NEMO_GYM_CHECKPOINT_E2E=1; about 3 minutes in total.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…de model calls

The replacement attempt commits model calls to its capture ledger, Gym crashes before the next checkpoint, and the same checkpoint is restored again. The second run of the attempt continues from the checkpoint's boundary, and none of the dead run's calls are in its lineage.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…ng behind

A native episode is retired mid-flight with token capture on. Afterwards no server is stopping or tracking anything for it, the model makes no further call for it, and its capture ledger is retired.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch from ca27fd2 to f98be7a 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: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