Skip to content

feat(checkpoint): continue environment server episodes from their boundaries - #3885

Draft
ananthsub wants to merge 4 commits into
ananthsub/partial-ckpt-model-worker-cutsfrom
ananthsub/partial-ckpt-environment
Draft

ananthsub wants to merge 4 commits into
ananthsub/partial-ckpt-model-worker-cutsfrom
ananthsub/partial-ckpt-environment

Conversation

@ananthsub

@ananthsub ananthsub commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What changed and why

Every environment server takes part in checkpoints once checkpointing is on. The environment server owns the episode protocol, so it records where each episode can continue.

  • An episode records a boundary after each completed step and runs each step in wait or replay mode. Prepare closes episode admission and is ready once every live episode is parked at a boundary or inside a replay step.
  • A restored boundary is handed to the replacement attempt's /run, which continues from it. The task must match the one recorded.
  • single_agent_turn runs seed, invoke agent, close agent, verify, and return as explicit stages with typed session handles. Seeding and the agent invocation replay; closing the agent waits; verification uses the mode the resources server reported on its seed reply, kept in the session handles so a restored episode still knows it.
  • The single_agent_turn_legacy adapter checkpoints too, since its rows run the same protocol.
  • A restored episode whose replacement has not started yet is exported again by the next checkpoint while the controller's commit scope still continues it, so a second crash still continues it. Once the scope leaves it out, the restored record is released.
  • Final cleanup, which closes the episode's agent and resources sessions, runs before /run replies on every path, and a retire never interrupts it. A retire cancels an episode only before cleanup starts, then waits until the episode has ended, so its reply means the sessions are closed. This matters because anyio's shielded scope does not stop a native task.cancel().
  • Checkpointing refuses num_workers > 1 on environment servers, since episodes live in one process.

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

single_agent_turn records a boundary after every step. A checkpoint exports the latest boundary, and a restored attempt starts at its next step.

flowchart LR
  B1(["next: seed"]) --> S["seed the resources session"]
  S --> B2(["next: invoke_agent"]) --> I["invoke the agent<br/>which parks its own turn loop"]
  I --> B3(["next: verify"]) --> V["verify<br/>wait or replay, per the resources server"]
  V --> B4(["next: return"]) --> R["reply to /run"]
Loading

Final cleanup closes the episode's sessions before /run replies. A retire never interrupts it: it cancels an episode only before cleanup starts, and then waits for the episode to end.

flowchart TD
  X{"retire arrives"} -->|"episode still in its protocol"| C["cancel the episode"] --> F
  X -->|"final cleanup already started"| F["wait"]
  F --> CL["final cleanup closes the agent and resources sessions"] --> D["retire replies"]
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 (this PR)
  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)

New compared with the old stack, which had no environment servers. It takes over the protocol side of #2942 (where an episode can continue), which the old stack reconstructed from agent and resources state after the fact.

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: 97 passed, including, each failing without its change:
    • a retire during an episode's final cleanup does not interrupt it;
    • a retire replies only after the episode and its cleanup have ended;
    • a commit whose scope leaves out a restored episode releases it.
  • RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider tests/unit_tests/test_environment_server.py: 16 passed.
  • RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider environment_servers/single_agent_turn/tests: 15 passed.
  • RAY_TMPDIR=/tmp .venv/bin/python -m pytest -q -p no:cacheprovider environment_servers/single_agent_turn_legacy/tests: 10 passed.
  • The full core unit suite (tests/unit_tests, eight processes, without test_opensandbox_compose.py, which needs the optional opensandbox package): 6,518 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 evidence is in the top PR (native episode continues from its last boundary; second crash before the replacement starts).
  • Real-model rollout evidence: pending.

Compatibility

  • Nothing changes unless the global checkpoint: block sets enabled: true.
  • With checkpointing on, an environment server with num_workers > 1 fails at startup.

@ananthsub ananthsub added feature New capabilities, enhancements, or enablement work area:env-infra Shared environment framework, lifecycle, registry, scaffolding, and validation 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-environment branch 2 times, most recently from 7c6475d to b9ac5f1 Compare October 1, 2026 21:53
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-environment branch from b9ac5f1 to 6c95786 Compare October 1, 2026 22:05
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-environment branch from 6c95786 to 1739ac6 Compare October 2, 2026 13:22
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-environment branch from 1739ac6 to 2812ca7 Compare October 2, 2026 19:43
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-environment branch from 2812ca7 to 5aa1530 Compare October 2, 2026 21:17
…ndaries

Every environment server takes part in checkpoints once checkpointing is
on. The environment server owns the episode protocol, so it records
where each episode can continue:

- An episode records a boundary after each completed step and runs each
  step in wait or replay mode (nemo_gym._checkpoint.steps). Prepare
  closes episode admission and is ready once every live episode is
  parked at a boundary or inside a replay step.
- A restored boundary is handed to the replacement attempt's run, which
  continues from it. The task must match the one recorded.
- single_agent_turn runs seed, invoke_agent, close_agent, verify, and
  return as explicit stages with typed session handles. Seeding and the
  agent invocation replay; closing the agent waits; verification uses
  the mode the resources server reported on its seed reply, kept in the
  session handles so a restored episode still knows it, and its result is
  recorded before the episode returns.
- The single_agent_turn_legacy adapter checkpoints too, since its rows
  run the same protocol.
- A restored episode whose replacement has not started yet is exported
  again by the next checkpoint, so a second crash still continues it.
- Checkpointing refuses num_workers > 1, since episodes live in one
  process.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
A controller retires stragglers while a checkpoint is open, and the episode
steps cancel a retired episode's task. When the episode had already ended
and was running its final cleanup, which closes its agent and resources
sessions, that cancellation interrupted the cleanup and leaked both
sessions: the cleanup's anyio shield holds against anyio cancellation, not
a native task.cancel().

The environment server now tells the participant when an episode enters
final cleanup. A retire after that point stops tracking the episode but
does not cancel its task.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…has claimed

A restored episode record waits for its replacement attempt's /run. The environment participant now reports unclaimed records as pending, so a commit that no longer continues their episode releases them instead of exporting them again at every checkpoint.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…e and its final cleanup

Retiring an episode now cancels it, unless its final cleanup already started, and waits until it has ended, so the retire replies only after the episode's agent and resources sessions are closed.

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-environment branch from 5aa1530 to a53fcad 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:env-infra Shared environment framework, lifecycle, registry, scaffolding, and validation feature New capabilities, enhancements, or enablement work

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant