Skip to content

feat(checkpoint): export IPI resource session state - #3900

Draft
zyzhou5 wants to merge 20 commits into
NVIDIA-NeMo:ananthsub/partial-ckpt-e2efrom
zyzhou5:zezhou/ipi-checkpoint-session-hooks
Draft

zyzhou5 wants to merge 20 commits into
NVIDIA-NeMo:ananthsub/partial-ckpt-e2efrom
zyzhou5:zezhou/ipi-checkpoint-session-hooks

Conversation

@zyzhou5

@zyzhou5 zyzhou5 commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

What does this PR do?

Save and restore Indirect Prompt Injection resource state through the session hooks introduced in #3886. Restored tool calls use the same cookie session and continue from the saved environment, including completed mutations. Retirement removes that session's state.

This ports row 22 to the new checkpoint stack. It is separate from the legacy-stack implementation in #3548. Based on #3889 (ananthsub/partial-ckpt-e2e), including the checkpoint interfaces from #3888 and SimpleAgent continuation used to exercise this environment.

Validation

Rebased onto #3889 at 6a3c6678f7e0332718cd8434a9a4357a7710532f; the test and lint checks below were rerun on the rebased branch.

  • 292 tests passed across resources_servers/indirect_prompt_injection/tests and tests/unit_tests/test_checkpoint_resources.py; IPI app coverage: 99.35%.
  • Test invocation (using a local runner that pins imports to this worktree): pytest --import-mode=importlib -o addopts= -q --tb=short resources_servers/indirect_prompt_injection/tests tests/unit_tests/test_checkpoint_resources.py --cov=resources_servers.indirect_prompt_injection.app --cov-report=term-missing --cov-fail-under=96.
  • Tests exercise on-disk controller save/restore after seeding and at two tool boundaries, exact nested-state recovery, invalid-state rejection, independent snapshots, and session retirement.
  • pre-commit run --all-files passed.

Rollouts

The environment-specific process drivers described below are local validation scripts; they have not yet been added to the committed E2E suite.

Local process recovery passed with the #3889 harness, the actual IPI and SimpleAgent servers, and a scripted model backend:

  1. Save after a chart mutation while the next model request is pending.
  2. Resume the original execution and record its result.
  3. Stop every Gym process, start fresh processes, restore the same checkpoint, and finish the rollout.

The restored execution made only the four remaining model calls. Tool calls, tool outputs, chart contents, and safety/utility rewards matched the original continuation. A negative control that omitted resource restoration produced missing-session tool errors and different chart outputs, even though its reward remained 1; reward alone was not used as the recovery assertion.

This uses scripted responses, not real model inference or a training job. Keep this PR as a draft pending representative real-model validation.

Compatibility and benchmark impact

Existing cookie-based initialization, tool routing, verification, and YAML configuration are unchanged. The shared checkpoint layer owns episode/session bookkeeping. Verification keeps the default wait behavior because it removes session state. No intended benchmark-score changes. Native Environment Server migration is outside this PR.

Documentation: no new user-facing configuration; the implementation follows the checkpoint hook contract in the base stack.

Checklist

  • I have read the contributing guidelines.
  • The change is focused; unrelated changes are excluded.
  • Tests added or updated and pass locally.
  • Documentation assessed; no new configuration or public workflow.
  • Pre-commit checks pass locally (pre-commit run --all-files).
  • New source files include the required Apache-2.0 SPDX header.
  • All commits have DCO sign-off (git commit -s).

@zyzhou5 zyzhou5 added feature New capabilities, enhancements, or enablement work area:environment Individual environments, benchmarks, verifiers, and environment-specific resources servers labels Oct 1, 2026
@copy-pr-bot

copy-pr-bot Bot commented Oct 1, 2026

Copy link
Copy Markdown

This pull request requires additional validation before any workflows can run on NVIDIA's runners.

Pull request vetters can view their responsibilities here.

Contributors can view more details about this message here.

…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.
- Records keep each participant's opaque state (an environment, a session,
  a boundary) exactly as exported: key order, NaN, and infinities round-trip.
  Payload fields are checked once with orjson for JSON-ness, then written as
  they are, instead of a deep validation and copy of large states.

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. 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 NVIDIA-NeMo#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>
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-agent branch from cbabe4a to d25d7e6 Compare October 1, 2026 22:05
…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 discards what a dead execution left for the target attempt
and for later attempts of the same rollout: the ledger, token records, intents,
the incomplete marker, and capture state. A ledger the restoring process
itself appended to belongs to a live episode, so the import still refuses it
before writing anything.

Signed-off-by: Ananth Subramaniam <ansubramania@nvidia.com>
Ported from the prefix-recovery work (Gym NVIDIA-NeMo#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>
…over MCP

A resources server with expose_tools_over_mcp and checkpointing failed at
startup: MCP auto-exposure refused the checkpoint middleware, because direct
MCP dispatch skips middleware.

- The checkpoint middleware declares that it applies to the MCP request as a
  whole, and MCP auto-exposure accepts it. Admission, in-flight counting, and
  attempt fencing work on the one POST that carries each tool call.
- An MCP request names its session with the signed token, not the cookie,
  so the middleware reads the session from the token and fences retired
  sessions. The token key is deterministic, so a harness keeps its token
  across a restore.
- While a checkpoint is open, an MCP call waits for resume instead of being
  refused. MCP clients are third-party agent harnesses that would hand the
  refusal to the model as a tool error. A waiting call is not in flight, so
  it does not hold up prepare.
- Checkpoint control routes are no longer harvested as MCP tools.

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>
…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>
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-agent branch from d25d7e6 to a508546 Compare October 2, 2026 13:22
Signed-off-by: Zeyu Zhou <zezhou@nvidia.com>
@zyzhou5
zyzhou5 force-pushed the zezhou/ipi-checkpoint-session-hooks branch from fefe83f to 3f635ab Compare October 2, 2026 18:18
@zyzhou5
zyzhou5 changed the base branch from ananthsub/partial-ckpt-agent to ananthsub/partial-ckpt-e2e October 2, 2026 18:18
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch from 6a3c667 to f3ade85 Compare October 2, 2026 19:43
@ananthsub
ananthsub force-pushed the ananthsub/partial-ckpt-e2e branch 2 times, most recently 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:environment Individual environments, benchmarks, verifiers, and environment-specific resources servers feature New capabilities, enhancements, or enablement work

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants