Skip to content

feat(checkpoint): persist recoverable Gym rollout state - #4265

Open
macandro96 wants to merge 2 commits into
amahishi/gym-checkpoint-participationfrom
amahishi/gym-durable-rollout-state
Open

macandro96 wants to merge 2 commits into
amahishi/gym-checkpoint-participationfrom
amahishi/gym-durable-rollout-state

Conversation

@macandro96

@macandro96 macandro96 commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Defines and persists the rollout-level state required to recover Gym executions: stable execution ownership, completed-result receipts and ACK obligations, agent continuations, resource dependencies, policy lineage/TQ references, and source-to-replacement attempt mappings.

This is 2/3 in the stacked decomposition of #4117 and depends on #4264.

Stack PR Scope
1/3 #4264 Gym participant protocol
2/3 This PR Durable rollout state and validation
3/3 #4266 Single Controller orchestration

This layer defines what a recoverable rollout snapshot contains. The next PR schedules, publishes, and restores those snapshots from Single Controller.

Why

Participant files alone are insufficient to resume a rollout. Recovery must also prove that:

  • a completed Gym result is durably owned exactly once;
  • Gym completion acknowledgements are not silently lost;
  • each saved agent continuation refers to valid policy lineage;
  • every external token reference exists in the Transfer Queue;
  • resource revisions match the continuation that depends on them;
  • an unfinished source attempt maps to one replacement attempt;
  • executions using restart_only resources restart without discarding unrelated continuable work.

Without these relationships, recovery can duplicate a rollout, acknowledge the wrong result, or restore agent state against missing model/resource data.

Durable-state flow

flowchart LR
    A[Gym execution completes] --> B[Completion receipt]
    B --> C[Rollout manager seals result]
    C --> D[Recovery ledger records ACK obligation]
    D --> E[Gym validates and accepts ACK]
    E --> F[Recovery ledger removes obligation]

    C --> G[Agent continuation index]
    G --> H[Policy lineage and TQ reference index]
    G --> I[Resource revision dependencies]
    H --> J[Cross-artifact validation]
    I --> J
    F --> J
    J --> K[Recoverable rollout state]

    K --> L[Source rollout attempt N]
    L --> M[Replacement attempt N plus 1]
    M --> N{Resource mode}
    N -->|export_restore| O[Continue from saved boundary]
    N -->|restart_only| P[Restart only dependent execution]
Loading

Main changes

  • Add a Gym execution registry with atomic registration/freeze boundaries.
  • Carry receipt-bound completion acknowledgements instead of positional or response-only identity.
  • Add a persistent RolloutManager acknowledgement sink and durable pending-ACK records.
  • Extend the rollout recovery ledger with attempt mappings, stable ownership, sealed sibling/group state, and ACK obligations.
  • Persist agent continuation roots and per-continuation resource dependencies.
  • Collect policy lineage and external TQ storage references.
  • Validate participant manifests, continuation relationships, digests, capture keys, and TQ references before restore.
  • Reject unsupported checkpoint schema versions with an actionable message; this stack writes schema v3.
  • Support selective restart of continuations that depend on restart_only resources.
  • Add APIs to discard restored continuations that are intentionally restarted.

Correctness rules

  • A receipt identifies the exact durable Gym result being acknowledged; stale or mismatched receipts are rejected.
  • Pending ACK obligations remain explicit until Gym accepts them and RL records their removal.
  • PENDING_MODEL continuations must point to real captured model lineage.
  • Continuation coordinates cannot silently move backward or disappear.
  • External TQ keys are deduplicated only when their metadata is identical; conflicting duplicates fail closed.
  • Restored source attempt N is never reused as the new physical owner; recovery creates a tracked replacement attempt.
  • Restart-only resource dependencies affect only the continuations that reference those resources.

Scope

  • Builds on the participant protocol in feat(checkpoint): add Gym participant coordination #4264.
  • Does not yet start a periodic checkpoint pump or publish snapshots from Single Controller.
  • Does not add in-generation token-prefix restoration.
  • The checkpoint schema is experimental; backward restore from earlier private schema revisions is intentionally unsupported.

Test plan

Targeted coverage includes:

  • tests/unit/environments/test_gym_checkpoint.py
  • tests/unit/environments/test_nemo_gym.py
  • tests/unit/environments/test_nemo_gym_checkpoint.py
  • tests/unit/environments/test_nemo_gym_token_capture.py
  • tests/unit/experience/test_rollout_manager.py
  • tests/unit/experience/test_rollout_recovery.py
  • tests/unit/data_plane/test_rollout_reassembler.py
  • tests/unit/experience/test_rollout_generation_failures.py

Suggested command:

uv run pytest -q \
  tests/unit/environments/test_gym_checkpoint.py \
  tests/unit/environments/test_nemo_gym.py \
  tests/unit/environments/test_nemo_gym_checkpoint.py \
  tests/unit/environments/test_nemo_gym_token_capture.py \
  tests/unit/experience/test_rollout_manager.py \
  tests/unit/experience/test_rollout_recovery.py \
  tests/unit/data_plane/test_rollout_reassembler.py \
  tests/unit/experience/test_rollout_generation_failures.py

Before review

  • Contributor conventions followed.
  • Unit coverage added for receipts, continuations, validation, and recovery mappings.
  • Full target environment test run recorded in CI.

@copy-pr-bot

copy-pr-bot Bot commented Sep 25, 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.

@terrykong terrykong left a comment •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This PR adds the rollout-level state needed to restore Gym runs: receipt-bound ACKs, ACK obligations in the recovery ledger, agent continuation roots, and the checks that run before restore. This review covers only the diff on top of #4264 (stack 1/3). #4266 (stack 3/3) was used to check who calls the new code.

Please fix before merge

  • rollout_manager.py L2434: in prompt_group mode a finished sibling is ACKed only when the whole group seals. A group that fails part-way leaves Gym results that are never ACKed, and Gym then refuses every later checkpoint prepare.
  • nemo_gym.py L1906: a Gym run that finished but whose reply RL never uses (its own rollout_s deadline, an early stop, or a failed receipt fetch) blocks every later checkpoint for good.

The rest are two test gaps, small cleanups, one comment-wording fix, and one question.

  • nemo_gym.py L1071: control calls share one connection pool with every open /run, so the ACKs (sent one POST at a time right before prepare) queue behind parked rollouts and time out. Asks for a separate pool, then one ACK request per agent.

Checkpoint levels. The #4266 review proposes one target level, rollout_recovery.target_level: turn | sibling | prompt_group, that RL lowers per environment from what Gym reports (the target-level comment). The prompt-group and lost-reply fixes above are needed under that design: group-scored environments such as GenRM compare always retry their whole group.

What works well

  • GymActorExecutionRegistry is a small, clean piece that can be tested without Ray.
  • The Single Controller side, here and in #4266, uses only public helpers.
  • The validators sit next to the wire models they check.
  • The copied Gym protocol matches the pinned Gym exactly: every model, route and feature name round-trips.
  • The ACK order is right: record the obligation, send it, and remove it only after Gym accepts it.

Tests
All 473 unit tests in the PR's test plan and the other touched test files pass on a CPU-only venv. Every new test is collected by a CI lane, and the functional sibling-recovery test is already registered. Removing the new receipt-identity checks, or the three restore-time checks on continuations, leaves every test passing — see the two test comments (rollout_recovery.py L712 and gym_checkpoint.py L1141).

The PR body's "Support selective restart" is implemented in #4266; this layer only carries the data.

CI and what was not run

  • CI has not run yet (only the copy-pr-bot notice). The Lint job should fail on import order (nemo_gym.py L60).
  • Nothing was run on GPU. The Gym findings are traced from source at the pinned Gym commit. The prompt-group and lost-reply stalls, and the retire + adopt fix, were also checked against Gym's own checkpoint code, run in-process (no live Gym server).

Visual explainer: https://terrykong.github.io/gh-pages-poc/terryk/pr-4265-durable-gym-state.html
Visual explainer: https://terrykong.github.io/gh-pages-poc/terryk/pr-4265-prompt-group-ack.html
Visual explainer: https://terrykong.github.io/gh-pages-poc/terryk/pr-4265-lost-reply-stall.html

Generated by Claude Code

Comment thread nemo_rl/experience/rollout_manager.py Outdated
Comment on lines +2434 to +2438
if self._gym_acknowledgement_sink is not None:
self._recovery_ledger.record_sealed_group_acknowledgements(
cut,
group_id,
)

@terrykong terrykong Sep 25, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item — please fix in this PR.

TL;DR — In prompt_group mode, a finished sibling's ACK is recorded only when the whole group seals. If the group fails first, that Gym result is never ACKed, and Gym then refuses every later checkpoint prepare.

This is new in this PR: the ACK sender and "record the ACK when the group seals" are both added here. How it shows up:

  1. With turn-level checkpointing on, Gym holds each finished run as COMPLETED until RL ACKs its receipt (Gym agent.py#L288-L290):

    Successful terminal results remain replayable until their receipt is acknowledged. Prepare reports those results as completed_unacknowledged and must not publish while that count is nonzero.

    retire() refuses a COMPLETED run, so an ACK is the only way to release it.

  2. sibling mode records the ACK as each row seals (L2445-L2465). prompt_group mode keeps finished siblings in pending_group_results and returns early (L2426-L2442) until all N arrive. Only then does this call to record_sealed_group_acknowledgements run.

  3. If another sibling fails, _abandon_entire_group clears the finished siblings' completion_receipt. No ACK is ever recorded for them.

  4. Those runs stay COMPLETED in Gym. Every later agent prepare gets a 409, and _prepare_agent_checkpoint retries until it raises TimeoutError: Gym agent '...' missed the prepare deadline. No Gym-aware snapshot is published again in that process.

This path is meant to stay reachable, and environments that score each sibling on its own hit it most: their siblings finish at different times, so a group is often half done when a snapshot starts. Under the target-level comment on #4266, an environment runs this branch at target turn when an override forces prompt_group on it, or when it is group-scored like GenRM compare. GenRM rarely stalls here, because its siblings all finish together once scoring is done. The fix below is needed; a setup check that rejects the combination is not the answer.

AI-1

In the PROMPT_GROUP branch of _record_streamed_completion, record each sibling's ACK obligation as soon as it arrives: inside a _recovery_mutation cut, then call notify_ready(). Keep the ownership seal all-or-nothing, as it is now. This is safe: a sibling ACKed this way but not yet published is not in RL's snapshot, so on restore it is simply rerun under a new try number. No result is lost or counted twice. It needs a new ledger method that builds the ACK from the SiblingSealResult's completion_receipt plus the current attempt identity, without requiring the SEALED state. Not a suggestion block: the change spans this file and rollout_recovery.py.

Scope: this only removes the stall. What a prompt group does on restore (start over, or resume its parked siblings from their saved turns) is set by the target level in the target-level comment on #4266. This fix works with either.

Visual explainer: https://terrykong.github.io/gh-pages-poc/terryk/pr-4265-prompt-group-ack.html

Second case: a snapshot that starts while a group is half done

Gym prepare parks the siblings that are still running. _wait_prepared returns once nothing is RUNNING, and its finally block resumes only PARK_REQUESTED runs, not PARKED ones. So the group cannot finish, the finished sibling cannot be ACKed, and prepare waits out prepare_timeout_s (300 s default in #4266) with admission closed, then rolls back. This costs time on each snapshot; no data is lost.

Why the ACK sender in #4266 cannot release these results

_flush_completed_gym_acknowledgements only sends obligations that are already in the ledger, and these results never get one. Gym's readiness check is status():

"ready_to_commit": not blocking_attempts and not parked_without_boundary and not completed_unacknowledged

A ledger-only run of this path ends with after PG retry attempt indices [1, 1] pending acks 0. A real Gym agent participant, with sibling g0 finished and g1 still running at the cut, reports ready_to_commit=False running=0 parked=1 completed_unacknowledged=1 today; ACKing g0 when it arrives gives ready_to_commit=True.

Comment thread nemo_rl/environments/nemo_gym.py Outdated
Comment on lines +1906 to +1909
completion_receipt = await self._completion_receipt_for(
execution,
agent_name=nemo_gym_row["agent_ref"]["name"],
)

@terrykong terrykong Sep 25, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item — please fix in this PR.

TL;DR — When a Gym /run finishes but RL never uses its reply, Gym keeps that run as finished-and-unACKed forever, so every later checkpoint prepare fails. RL also reruns the rollout as try N+1, throwing away a finished rollout.

New in this PR (receipts and ACKs are added here). Four ways RL drops a finished reply:

In each case Gym already marked the run COMPLETED (Gym base_responses_api_agent.py#L244), and retire() refuses a COMPLETED run. Rerunning and ACKing try N+1 does not release try N. A real Gym run stays stuck:

after RL reran+ACKed try 1: ready_to_commit= False completed_unacknowledged= 1 stuck: [('rD', 0)]

Why a finished-but-unACKed run blocks Gym is explained in the comment on rollout_manager.py L2434.

AI-1

Before minting try N+1 for a row RL gave up on, call Gym's existing retire(id, N) route (Gym agent.py#L1474-L1485). If it answers retired: True, try N was still running and is now cancelled, so mint N+1 as today. If it answers completed_unacknowledged: True, try N finished: fetch its stored reply, seal it and ACK it instead of minting N+1. Re-POSTing /run for the same try returns the stored reply without running again, because begin() hands back a finished run before the admission check (Gym base_responses_api_agent.py#L217-L218, Gym agent.py#L320-L326). At prepare, do the same for every entry in the 409's completed_unacknowledged_attempts (Gym agent.py#L727-L733) so nothing is left behind. /retire is fenced by the checkpoint phase (Gym agent.py#L1479-L1482): pass the active checkpoint ID when a snapshot is in flight, and wait until Gym resumes if the phase is COMMITTING or COMMITTED_PAUSED. This needs no Gym change. Not a suggestion block: it needs a new actor method, a change to the retry path and a change to the prepare path, across several files.

Visual explainer: https://terrykong.github.io/gh-pages-poc/terryk/pr-4265-lost-reply-stall.html

Real Gym run: stuck today, recovered with retire + adopt

Today (reply lost, then RL reruns as try 1):

prepare #1: ready_to_commit=False completed_unacknowledged=1
prepare #2: ready_to_commit=False completed_unacknowledged=1
prepare #3: ready_to_commit=False completed_unacknowledged=1
retire(rD, 0) -> {'retired': False, 'tombstoned': False, 'completed_unacknowledged': True}
prepare after retire: ready_to_commit= False
after RL reran+ACKed try 1: ready_to_commit= False completed_unacknowledged= 1 stuck: [('rD', 0)]

With the fix (adopt try 0 instead of rerunning):

prepare: ready_to_commit=False running=0 parked=0 completed_unacknowledged=1 | listed: [('rD', 0)]
re-POST /run (rD, 0) during prepare -> replayed: {'id': 'r-rD', 'reward': 0.5}
ACK -> {'acknowledged': True, 'idempotent': False}
prepare after adopt+ACK: ready_to_commit=True running=0 parked=0 completed_unacknowledged=0
agent runs: {('rD', 0): 1}

agent runs: {('rD', 0): 1} means the agent ran once: the reply was replayed, not recomputed.

When try 0 is still running at RL's deadline, retire cancels it and prepare is ready:

retire(rT, 0) -> {'retired': True, 'tombstoned': True}
agent task cancelled: True
prepare: ready_to_commit=True running=0 parked=0 completed_unacknowledged=0

How /retire behaves in each checkpoint phase:

== no checkpoint in flight (IDLE)
  retire(fresh id)            -> {'retired': False, 'tombstoned': True}
== snapshot ckpt-A in flight: prepared
  fence phase=prepared active=ckpt-A
  retire(other id)            -> CheckpointConflictError: checkpoint 'ckpt-A' is active in phase 'prepared'; refusing operation for 'rl-retire-2'
  retire(ckpt-A)              -> {'retired': False, 'tombstoned': True}
== ckpt-A committed, Gym paused until RL resumes it
  fence phase=committed_paused
  retire(ckpt-A)              -> InvalidPhaseError: operation not valid in phase 'committed_paused' (valid from: ['idle', 'prepared', 'preparing'])
  /acknowledge has no phase check: route source calls only require_control_auth -> True
== after resume
  retire(ckpt-A, now finished) -> StaleCheckpointError: checkpoint 'ckpt-A' already finished with outcome 'resumed'; this call is from a stale coordinator
  retire(fresh id)             -> {'retired': False, 'tombstoned': True}

Comment thread nemo_rl/experience/rollout_recovery.py Outdated


@dataclass(frozen=True)
class PendingCompletedExecutionAcknowledgement:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item (low priority).

TL;DR — PendingCompletedExecutionAcknowledgement re-types every GymCompletionReceipt field plus agent_name, which is exactly what this PR's GymCompletedExecution already is. Storing that model instead removes a field list that is now written out in six places.

New in this PR. The two side by side:

# rollout_recovery.py L335-L370
@dataclass(frozen=True)
class PendingCompletedExecutionAcknowledgement:
    rollout_id: str
    attempt_index: int
    agent_name: str
    execution_generation: int
    result_identity: str
    result_digest: str
    manifest_capture_key: Optional[str] = None
    terminal_model_call_id: Optional[str] = None
    ...
    @property
    def receipt(self) -> GymCompletionReceipt: ...  # rebuilds the receipt
# gym_checkpoint.py L563
class GymCompletedExecution(_StrictWireModel):
    receipt: GymCompletionReceipt
    agent_name: str = Field(min_length=1)

Links: the dataclass, receipt property, GymCompletedExecution. #4266 converts one into the other: GymCompletedExecution(receipt=ack.receipt, agent_name=ack.agent_name) (#4266 single_controller.py#L1884).

The field list is repeated in: the fields, __post_init__, the receipt property, _COMPLETED_EXECUTION_ACKNOWLEDGEMENT_FIELDS, state_dict, from_state_dict (whose checks at L1311-L1320 repeat L349-L358), and _acknowledgement_for_sibling. A drift fails loudly (Gym rejects the ACK, or from_state_dict rejects an unknown field), so this is upkeep cost, not a bug.

AI-1

Store GymCompletedExecution in _pending_completed_execution_acknowledgements, keyed by (receipt.rollout_id, receipt.attempt_index). state_dict becomes model_dump(mode="json"), and from_state_dict becomes GymCompletedExecution.model_validate(raw); extra="forbid" on _StrictWireModel (gym_checkpoint.py#L74) replaces _reject_unknown_fields. Delete the dataclass, the frozenset and the per-field checks. Schema v3 is new in this PR, so the saved shape can change now at no cost. Not a suggestion block: it touches several places in the file.

Comment thread nemo_rl/environments/gym_checkpoint.py Outdated
missing_agent_checkpoint_participation.append(
contract.participant.participant_name
)
if "completed_result_acknowledgement" not in contract.features:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item (low priority).

TL;DR — The ACK feature name is a bare string here and at nemo_gym.py#L1018, while every other Gym feature this module copies has a named constant. A bare string cannot be checked against Gym.

Gym names the same string COMPLETED_RESULT_ACKNOWLEDGEMENT_FEATURE. The other feature constants sit at L53-L56.

AI-1

Add GYM_AGENT_COMPLETED_RESULT_ACKNOWLEDGEMENT_FEATURE = "completed_result_acknowledgement" next to the other feature constants, use it on this line and at nemo_gym.py#L1018, and add it to the source map from the drift-test comment on this file (nemo_gym._checkpoint.agent:COMPLETED_RESULT_ACKNOWLEDGEMENT_FEATURE). This line becomes:

if GYM_AGENT_COMPLETED_RESULT_ACKNOWLEDGEMENT_FEATURE not in contract.features:

Not a suggestion block: the line only works together with the new constant, which is defined elsewhere in the file.

Comment thread nemo_rl/environments/gym_checkpoint.py Outdated
def __init__(self) -> None:
self._live: dict[tuple[str, int], GymActorExecution] = {}
self._frozen_checkpoint_id: str | None = None
self._frozen_membership: tuple[GymActorExecution, ...] = ()

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item (low priority).

TL;DR — GymActorExecutionRegistry keeps a frozen membership, a TERMINAL state and a status() that no production code reads. Only the dispatch fence and the duplicate-key check do real work.

The class docstring says it will "retain the exact membership of a checkpoint cut", but both callers drop what freeze() returns (nemo_gym.py#L1165, #L1500), and #4266 adds no reader. mark_terminal() is called at nemo_gym.py#L1900, but only tests read the TERMINAL state it sets, through status().

AI-1

Delete _frozen_membership, GymActorExecutionState (with TERMINAL), mark_terminal() and its call, and status(). Change the docstring to "Fence actor dispatch during a checkpoint cut." Not a suggestion block: deletions in several places across two files.

Comment thread nemo_rl/environments/gym_checkpoint.py Outdated
class GymExternalStorageReference(_StrictWireModel):
"""One TQ staging row required by a parked Gym continuation."""

schema_version: Literal[1] = GYM_CHECKPOINT_SCHEMA_VERSION

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item.

TL;DR — The Gym wire version is written as Literal[1] in 7 places, and pydantic never checks a default against its Literal. A bump that misses one place ships data RL itself cannot read back.

Partly pre-existing: #4264 wrote the first three copies, and this PR adds two more (this line and L960). The seven:

  • five RL-built models, with a default of GYM_CHECKPOINT_SCHEMA_VERSION: L364, L554, L597, this line, L960;
  • two required fields whose value Gym reports: L281 (GymControlCapabilities) and L327 (GymCheckpointParticipantContract).

If the constant becomes 2 and one Literal[1] is left behind, RL builds and sends {'schema_version': 2} with no error, then its own model refuses to read that data back: round-trip: REJECTED -> Input should be 1. Whether the constant matches Gym's side is checked by the drift-test comment on this file; this comment is about RL's copies agreeing with each other.

AI-1

Write the value once, in a type alias, and build the RL-side default from the constant:

# Must equal GYM_CHECKPOINT_SCHEMA_VERSION. pyrefly does not accept a named
# constant inside Literal[...], so the value is written out once, here.
GymCheckpointSchemaVersion: TypeAlias = Literal[1]


class _VersionedWireModel(_StrictWireModel):
    """A wire model RL builds; its schema_version comes from the constant."""

    schema_version: GymCheckpointSchemaVersion = Field(
        default=GYM_CHECKPOINT_SCHEMA_VERSION, validate_default=True
    )

Make the five RL-built models (L364, L554, L597, L950, L960) subclass _VersionedWireModel and drop their own schema_version line. Keep the two Gym-reported fields (L281, L327) required, but type them as schema_version: GymCheckpointSchemaVersion: a base class with a default must not absorb them, or a missing version would be silently accepted. Not a suggestion block: it changes 7 places across the file.

Context — no action. Literal[GYM_CHECKPOINT_SCHEMA_VERSION] looks like the obvious fix, but it fails this repo's type check (see the fold below).

What a half-done bump does today (run at this head)
  • Today, with the constant at 2 and a Literal[1] left behind: RL sends {'schema_version': 2} with no error, then round-trip: REJECTED -> Input should be 1.
  • With validate_default=True, the same half-done bump fails the first time any model is built: construction: REJECTED -> ('schema_version',) Input should be 1.
  • The shape above: RL-built default: 1, Gym-reported: missing version still rejected (required, as before), Gym-reported: version 2 rejected.
Why not `Literal[GYM_CHECKPOINT_SCHEMA_VERSION]`

This file is type-checked (pyrefly.toml#L174), and the pinned pyrefly==0.24.2 rejects a named constant inside Literal[...] (also with Final):

ERROR Expected a type form, got instance of `Literal[1]` [not-a-type]

Gym writes Literal[CONTROL_SCHEMA_VERSION], but Gym does not run this checker. The sketch above passes pyrefly 0.24.2 with 0 errors.

Comment thread nemo_rl/experience/rollout_manager.py Outdated
return False


class GymAcknowledgementSink:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item (nit).

TL;DR — GymAcknowledgementSink holds no data: it stores one callback and notify_ready() calls it. In this codebase a "sink" is where payloads get written, so the name suggests it stores ACKs.

The class keeps the callback set once by bind(on_ready) (L120-L143). The ACK obligation itself is written to the recovery ledger (L2435, L2458); the call only tells the controller that ACKs are ready to send (L2442, L2465). By contrast, TQTokenSink really is a sink: it writes the heavy payloads the generation engine intercepts to TQ.

AI-1

Rename it to GymAcknowledgementNotifier, to match notify_ready(). Carry the rename to gym_acknowledgement_sink and bind_gym_acknowledgement_sink (L1651, L1719, L1742), to the setup call (setup.py#L1989) and to test_gym_acknowledgement_sink_is_explicit_and_bound_once. #4266 builds the object at that setup call (#4266 setup.py#L2161-L2165), so it follows the rename. Not a suggestion block: the rename spans several files.

from nemo_rl.data.interfaces import DatumSpec

ROLLOUT_RECOVERY_SCHEMA_VERSION = 2
ROLLOUT_RECOVERY_SCHEMA_VERSION = 3

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item.

TL;DR — Nothing says what this version covers or when to bump it. Readers accept only this version (L43, L48-L54), so a saved field changed without a bump makes old ledgers fail with a confusing error, or load with the field missing.

New in this PR: the ledger goes from version 2 to 3.

AI-1

Say above the constant which file it versions and when to bump it, with the same wording as the manifest version:

Suggested change
ROLLOUT_RECOVERY_SCHEMA_VERSION = 3
# Version of the rollout recovery ledger (rollout_recovery.pt). Readers accept only
# this version, so bump it whenever a saved field is added, removed, renamed, or
# changes meaning, then regenerate tests/unit/single_controller/checkpoint_schema_lock.json.
ROLLOUT_RECOVERY_SCHEMA_VERSION = 3

The schema-lock test that enforces this is in the manifest-version comment, and it already checks the ledger. A new ledger field without a bump fails with: rollout_recovery_ledger: saved fields changed but the version is still 3 (added=['new_field'], removed=[]). Bump the version, then regenerate.

"task_name": record.prompt_ref.task_name,
},
"task_source": record.task_source,
"resolved_agent_name": record.resolved_agent_name,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item.

TL;DR — Every saved ledger field is written by hand in four places, and the loader reads each one with .get(), so a missing key restores as None with no error. A writer that forgets a field, or a rename done on one side only, goes unnoticed.

The four places: the dataclasses (L207-L377), the dict literal in state_dict (L1197-L1259), the hand-kept key sets (L57-L113), and the .get() plus isinstance checks in from_state_dict and _group_from_state (L1263-L1500). The ACK checks exist twice (L349-L358 and L1311-L1320). _reject_unknown_fields rejects extra keys only. Reloading a real state_dict() with one key removed:

ACCEPTED  attempt without completion_receipt -> completion_receipt=False
ACCEPTED  group without resolved_agent_name  -> resolved_agent_name=None
ACCEPTED  ack without manifest_capture_key / terminal_model_call_id
ACCEPTED  group without task_source
REJECTED  attempt without reward  (only because sealed attempts check it)

The pattern is older (#3923 added the ledger), but this PR adds resolved_agent_name, completion_receipt and the whole ACK list, and already bumps the ledger version from 2 to 3, so changing the saved format now costs no extra compatibility work. gym_checkpoint.py in this PR already models its saved data with pydantic and extra="forbid".

AI-1

Describe the saved ledger with pydantic models and use them on both sides:

  • state_dict() builds a RolloutRecoveryState and returns .model_dump(mode="json"). It has to be JSON mode: the file is read with torch.load(..., weights_only=True) (single_controller.py#L1038-L1042), and the default dump keeps enum objects, which that load refuses. JSON mode also cannot dump raw bytes, so type attempt_uuid as uuid.UUID (it saves as a string).
  • from_state_dict() calls RolloutRecoveryState.model_validate(state) and then builds the dataclasses. Keep the rules that span records as plain code: the sibling count equals expected_generations, group ids and attempt uuids are unique, groups that share an admission_id agree, and a sealed attempt has a reward.
  • For the ACK row, store agent_name plus a GymCompletionReceipt, as the ACK-type comment asks.
  • A field typed Optional[X] with no default is required in pydantic v2 but may be None; that is the missing-key check. Optional fields inside the reused GymCompletionReceipt keep Gym's own None defaults.

This removes the five key sets, _reject_unknown_fields, most isinstance checks and the second copy of the ACK checks. The sidecar wrapper (_SIDECAR_STATE_FIELDS) can extend the same model, and the schema-lock test in the manifest-version comment can then read the ledger's fields straight from the model. Not a suggestion block: new models plus both functions.

class _SavedState(BaseModel):
    # Optional[...] with no default: may be None, may not be missing.
    model_config = ConfigDict(extra="forbid", frozen=True)

class AttemptState(_SavedState):
    attempt_index: NonNegativeInt
    attempt_uuid: uuid.UUID
    status: RolloutAttemptStatus
    receipt: Optional[dict]
    completion_receipt: Optional[GymCompletionReceipt]
    reward: Optional[float]
    mask_sample: Optional[bool]
    staging_keys: list[str]

class SiblingState(_SavedState):
    generation_index: NonNegativeInt
    attempts: list[AttemptState]

class GroupState(_SavedState):
    group_id: str
    admission_id: str
    prompt_id: str
    prompt_ref: PromptRefState          # sample_id: str, task_name: Optional[str]
    task_source: Optional[str]
    resolved_agent_name: Optional[str]
    recovery_granularity: RecoveryGranularity
    expected_generations: PositiveInt
    target_step: Optional[int]
    start_weight_version: int
    status: PromptGroupStatus
    phase: PromptGroupPhase
    siblings: list[SiblingState]

class AckState(_SavedState):
    agent_name: str
    receipt: GymCompletionReceipt

class RolloutRecoveryState(_SavedState):
    schema_version: Literal[3]
    groups: list[GroupState]
    pending_completed_execution_acknowledgements: list[AckState]
What the prototype showed (run at this head)
mode="python": torch.load(weights_only=True) -> UnpicklingError: Weights only load failed
mode="json"  : torch.load(weights_only=True) OK; round-trip equal: True
REJECTED  attempt without completion_receipt: missing at groups.0.siblings.0.attempts.0.completion_receipt
REJECTED  group without resolved_agent_name: missing at groups.0.resolved_agent_name
REJECTED  group without task_source: missing at groups.0.task_source

)
logical_rollout_id = self.logical_rollout_id(generation_index)
attempt_index = sibling.current_attempt.attempt_index
return gym_capture_key(logical_rollout_id, attempt_index)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 action item (low priority).

TL;DR — This line was the last real use of the attempt's random UUID. Now attempt_id has no production reader, and attempt_uuid is saved and checked on load but used for nothing.

Before this PR the gate key was built from the UUID (base rollout_recovery.py#L254, linked at the base because this PR removes it):

            f"_a{sibling.current_attempt.attempt_id}"

This line builds it from the try number instead. After that:

This PR already bumps the ledger version from 2 to 3, so dropping the field needs no extra bump.

AI-1

Remove attempt_id and attempt_uuid: the property, the field, uuid.uuid4() in _new_attempt, the save and load code, the _ATTEMPT_STATE_FIELDS entry and the duplicate check. In the tests, drop the bad-UUID case (test_rollout_recovery.py#L324, #L349-L350) and the duplicate-UUID test (#L403-L404), and delete L822 and L855 (L856 already checks attempt_index == 0, which is what L855 guards). If you take the state_dict comment, leave attempt_uuid out of AttemptState too. Not a suggestion block: several places across two files.

@macandro96
macandro96 force-pushed the amahishi/gym-durable-rollout-state branch from f381ab1 to 9c955eb Compare September 29, 2026 03:05
Signed-off-by: Anish Mahishi <amahishi@nvidia.com>
Signed-off-by: Anish Mahishi <amahishi@nvidia.com>

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

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants