Conversation
…ination
The core of partial-rollout checkpointing, shared by every server. It
installs nothing on its own; the model, environment, agent, and
resources changes each add a participant on top of it.
- Identity: an episode attempt is keyed by its capture key ("r" or
"r-aN"). ServerClient attributes model and resources calls to their
rollout whenever checkpointing is on, not only with observability.
- Control plane (control.py): one phase machine behind bearer-protected
routes under /ng-control/v1/checkpoint, driving one participant per
server. The controller owns checkpoint-ID fencing, attempt fencing,
the readiness wait, deadline-bounded I/O, manifest-last storage, and a
lease that resumes a participant whose controller went away.
- Episode steps (steps.py): a boundary after each completed step, and a
wait or replay mode for a step in flight when a checkpoint closes.
Time parked for a checkpoint does not count against an episode's
deadline.
A resources server reports whether its /verify may be replayed in an
x-ng-checkpoint-verify header on its /seed_session reply, so callers
learn the mode without another request; the default is wait.
- Coordination (coordination.py): the Gym-owned participant order a
training controller calls: prepare environment servers, the policy
model, agents, then resources servers; resume in reverse; commit and
retire everywhere; restore all or nothing.
- Settings: a global checkpoint block, off by default, that requires a
control token when enabled.
- Checkpoint coordination and other control calls use a reserved aiohttp
pool (global_aiohttp_control_connector_limit, and _per_host), so data
calls that hold every connection to a host, such as parked /run calls,
cannot starve prepare, commit, or resume. coordination.renew extends
every participant's lease when publishing outlasts it.
- Participants may override async export and install hooks to keep slow
I/O off the event loop, and may add fields to the commit reply.
- An episode woken by a resume stays parked if a new checkpoint closed
admission before it ran.
Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
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, and a restore installs all ledgers with one directory sync instead of one 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>
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>
…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>
Every resources server takes part in checkpoints once checkpointing is
on, in one of three modes:
- restart_only (the default): live sessions block a checkpoint until
their rollouts are retired, so an unsupported server fails closed.
- stateless: tool results depend only on the request, so nothing is
exported. The weather example server uses it.
- exported: the server implements export_session_state,
restore_session_states, and retire_session_state. The counter example
server uses it. export_session_state raises KeyError for a session the
server already dropped (for example after a failed verification); the
participant stops tracking it instead of failing the commit.
The session lifecycle comes from /seed_session, /close_session, and
/verify, keyed by the rollout attempt from the rollout context, so no
request body is inspected. Servers whose protocol creates sessions
elsewhere, such as Gymnasium-style servers, call
checkpoint_session_started and checkpoint_session_ended.
The checkpoint_verify class attribute declares whether /verify is safe
to run again after a crash ("replay") or must finish before a
checkpoint ("wait", the default). It is a property of the code, so it is
not configurable. The server reports it in an x-ng-checkpoint-verify
header on every /seed_session reply, which is how callers learn it
without another request, and in its checkpoint status.
Replay-safe requests (/seed_session, and /verify when declared replay)
are admitted even after admission closes. The checkpoint does not wait for
the episode steps that send them, so refusing one would fail its episode.
A session seeded while closed is not part of that checkpoint: it neither
blocks it nor is exported by it, and becomes an ordinary session on
resume.
Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Resources session hooks become coroutines, and export takes every session at once: export_session_states(session_ids) returns the state of each session the server still holds. A server whose session state lives outside the process, such as a sandbox per session, can then snapshot it concurrently within the commit. A session left out of the result replaces the KeyError signal for a session the server already dropped. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…e verification A resources server that declares nothing is restart-only, so every checkpoint has to retire the rollouts that use it. These five servers keep no session state, and their verification can run again after a crash, so they declare checkpoint_mode = "stateless" and checkpoint_verify = "replay": - code_gen, competitive_coding_challenges, equivalence_llm_judge, and math_with_judge verify each request on its own. - genrm_compare keys its cohorts by prompt, not by session. Its verification must replay: a member waits for its siblings, so waiting on it would deadlock a checkpoint while siblings are still generating. After a crash, the members that had not recorded their reward re-verify and rebuild the cohort. A crash in the moment between a cohort's result and every member recording it leaves the rest waiting for siblings that will not re-verify; they fail after cohort_collection_timeout_s and are retried from input. The declarations are in code, not config: they describe what each server can do. A test reads them from source, because each server's dependencies live in its own environment. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Split from the episode orchestration prototype so the checkpoint stack can build on native agent and resources sessions without the Hermes changes. Simple Agent implements seed and close for agent sessions and uses the direct HTTP tool access supplied by the environment server. The example single tool call resources server implements native seed and close. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Simple Agent's and the example Resources Server's seed and close sections only read and update dicts and sets. asyncio switches requests only at an await, so no other request can run between their checks and updates, and the per-session locks guarded nothing. Removing them also stops creating a new lock on every call and keeping one per session for the life of the server. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…urces_session hook The base Resources Server now serves /close_session and registers it first, so the example's own route never ran and its sessions were never released. The example overrides close_resources_session instead, and the route test checks that the session was released rather than only that the call succeeded. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
…tion Every agent takes part in checkpoints once checkpointing is on. Agents that implement the session hooks (export_agent_session, restore_agent_sessions, retire_agent_session) and set checkpoint_sessions_supported get an AgentSessionParticipant: - A session is idle, running, or at a boundary; a running activation parks at its next boundary when a checkpoint closes admission, and one waiting on a model call counts as parked at its last boundary. - Boundaries are lazy snapshots, so a long turn loop does not serialize its history at every step. - A legacy /run is an episode of its own: seed and the turn loop replay, and verify uses the mode the resources server reported on its seed reply, which the loop boundary carries across a restore. - Restored sessions are installed under the replacement attempt. Agents without the hooks get a restart-only participant: their in-flight /run and /v1/responses calls block prepare until the controller retires them, and those rollouts restart from input. With checkpointing on, agent work without a rollout id is refused with rollout_id_required, since a checkpoint could neither record nor retire it. Simple Agent implements the hooks for both its environment-server sessions and its legacy /run loop. A restored session whose replacement has not started yet keeps its restored boundary and legacy episode in the next checkpoint, so a second crash still continues it. Checkpointing refuses num_workers > 1 for every agent, including restart-only ones, since each worker would track only its own calls. A session woken by a resume stays parked if a new checkpoint closed admission before it ran. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Agent session hooks become coroutines, and export takes every session at once: export_agent_sessions(session_keys). An agent whose session state lives outside the process, such as a sandbox it runs tools or a harness in, can then checkpoint all of it concurrently within the commit, while every session is parked at a boundary. The same change as the resources session hooks. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
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>
The controller of a checkpoint has to be the process that dispatches episodes: only it knows which episodes it will continue, and only it can send their replacements after a restore. For an evaluation run that is rollout collection, so with `checkpoint_dir` set it drives Gym's coordination itself. A checkpoint stops starting rows, prepares every participant, commits the rows whose /run has not replied, writes the run's manifest, and publishes it by pointing LATEST at it; the two latest checkpoints are kept. SIGUSR1 checkpoints and continues, SIGTERM checkpoints and stops (Slurm can send it before preempting a job), and `checkpoint_every_s` checkpoints on a timer. A run restarted with resume_from_cache restores the latest checkpoint before dispatching: its unfinished rows continue as their next attempt, and the rest start from their input as before. If the restore fails, coordination has retired that attempt, so those rows start from their input one attempt later. A fresh run forgets the checkpoints of an earlier run in the same directory. Without `checkpoint_dir`, rollout collection is unchanged. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Checkpoints emitted nothing: a prepare that stalled, a lease that expired and silently aborted a checkpoint, or a restore that took minutes left no trace. Spans, in a new opt-in `checkpoint` span group (`default,checkpoint`): - gym.checkpoint.<operation> on every participant for prepare, commit, restore, resume, and retire, with the participant kind, checkpoint ID, outcome, and record counts; - child spans for the steps that grow with the live set: wait_ready, export, write, read, install, and the generation-cut round; - gym.checkpoint.coordinate.<operation> in the controller, with one child per prepare stage, and spans for rollout collection's own checkpoint and restore. Control calls carry the trace context, so one checkpoint is one trace across Gym's processes. Metrics, recorded whenever telemetry exports, on Gym-owned instruments like the sandbox ones: gym.checkpoint.operation_duration_ms (by operation, kind, and outcome), gym.checkpoint.records_total and bytes_total (commit and restore), and gym.checkpoint.events_total for lease_expired, prepare_not_ready, refused (by error code), and generation_cut (by disposition). Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Profiling commit and restore with the checkpoint spans showed two costs that grow with the live set: - Record payloads. Opaque server state (resources sessions, agent sessions and boundaries, environment boundaries, model ledger rows) was validated as JsonValue when a record was built and when it was read, and copied by model_dump on write. For a 350 KB session that was about 8 ms at commit and 4 ms at restore. Payload fields are now JsonPayload: orjson checks that the value is JSON (it rejects non-string keys and non-JSON types, as JsonValue did) and the value is written as is, still with the standard json module so NaN and infinities round-trip exactly. - Lineage imports. Restore wrote and fsynced one ledger file per episode, after reading each one to check it was absent. The checkpoint is the durable copy and the import is idempotent, so files are no longer synced one by one (the directory is, once), and one directory listing replaces the existence reads. A later append syncs the file, imported rows included. 1,000 sessions of 350 KB: commit 11.5 s to 3.6 s, restore 6.0 s to 3.4 s. Importing 15,686 episodes' ledgers: 3.2 s to 1.5 s. 16,000 rollouts with two policy workers: commit 1.69 s to 1.43 s, restore 2.71 s to 2.01 s. Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
|
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
This was referenced Oct 1, 2026
Draft
ananthsub
force-pushed
the
ananthsub/partial-ckpt-telemetry
branch
from
October 1, 2026 22:05
fd8308c to
9775d38
Compare
Contributor
Author
|
Folded into the stack rather than kept as a separate optimization PR. The JSON payload handling ( |
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
Profiling commit and restore with the checkpoint spans showed two costs that grow with the live set.
JsonValuewhen a record was built and again when it was read, and copied bymodel_dumpon write. That covers resources sessions, agent sessions and boundaries, environment boundaries, and model ledger rows. For a 350 KB session this cost about 8 ms at commit and 4 ms at restore, which was most of the cost.JsonPayload.orjsonchecks that the value is JSON: it rejects non-string keys and non-JSON types, asJsonValuedid.jsonmodule, so NaN and infinities round-trip exactly and key order is kept.Measurements
The remaining cost of heavy state is the JSON encode and decode itself, about 2 ms each per 350 KB session. Splitting that across processes is the next step.
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, fencing, lease, storage"] 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")] C --> CO --> K K --> E & M & A & R M --> W G --> D classDef this fill:#fde68a,stroke:#b45309,stroke-width:2px,color:#1f2937 class K,M thisOne checkpoint, a crash, and the restore, end to end:
This PR
What a record's opaque payload goes through at commit and at restore, before and after.
flowchart TD subgraph B["Before"] B1["build the record<br/>deep JsonValue validation"] --> B2["model_dump<br/>copies the payload"] --> B3["json.dumps"] B4["json.loads"] --> B5["model_validate<br/>deep JsonValue validation"] end subgraph A["After"] A1["build the record<br/>one orjson check"] --> A3["json.dumps of the payload as is"] A4["json.loads"] --> A5["model_validate<br/>one orjson check"] endImporting restored ledgers into the lineage store.
flowchart LR L["list the ledger directory once"] --> E{"ledger file exists?"} E -->|no| W["write the rows, no fsync"] E -->|yes| C{"same rows?"} C -->|yes| S["skip"] C -->|no| X["refuse the restore"] W --> D["fsync the directory once"]Where this sits in the stack
This is one PR in a stack of draft PRs for partial-rollout checkpointing. Each PR's base is the branch of the PR before it, so each diff shows only that PR's commits:
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 cutsananthsub/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-restore-perf): perf(checkpoint): cheaper record payloads and lineage imports (this PR)Server ports that build on the stack: #3894 (Workplace Assistant), #3895 (Gymnasium), #3896 (Blackjack), #3897 (indirect prompt injection), #3898 (proof refinement).
Issue
No tracking issue exists. The profile came from the telemetry PR below this one.
Validation
pytest tests/unit_tests/telemetry tests/unit_tests/test_checkpoint_*.py tests/unit_tests/test_token_capture_ledger.py: 372 passed.pre-commit run --files <changed files>: passed.Rollout evidence
The scale test runs real Gym servers against a fake inference backend. Real-model runs are pending, as for the rest of the stack.
Compatibility