feat(checkpoint): persist recoverable Gym rollout state - #4265
macandro96 wants to merge 2 commits into
Conversation
There was a problem hiding this comment.
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.pyL2434: inprompt_groupmode 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.pyL1906: a Gym run that finished but whose reply RL never uses (its ownrollout_sdeadline, 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.pyL1071: 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
GymActorExecutionRegistryis 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.pyL60). - 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
| if self._gym_acknowledgement_sink is not None: | ||
| self._recovery_ledger.record_sealed_group_acknowledgements( | ||
| cut, | ||
| group_id, | ||
| ) |
There was a problem hiding this comment.
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:
-
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_unacknowledgedand must not publish while that count is nonzero.retire()refuses a COMPLETED run, so an ACK is the only way to release it. -
siblingmode records the ACK as each row seals (L2445-L2465).prompt_groupmode keeps finished siblings inpending_group_resultsand returns early (L2426-L2442) until all N arrive. Only then does this call torecord_sealed_group_acknowledgementsrun. -
If another sibling fails,
_abandon_entire_groupclears the finished siblings'completion_receipt. No ACK is ever recorded for them. -
Those runs stay COMPLETED in Gym. Every later agent prepare gets a 409, and
_prepare_agent_checkpointretries until it raisesTimeoutError: 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.
| completion_receipt = await self._completion_receipt_for( | ||
| execution, | ||
| agent_name=nemo_gym_row["agent_ref"]["name"], | ||
| ) |
There was a problem hiding this comment.
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:
- RL's own
rollout_sdeadline (_Deadline) fires after Gym finished but before the reply arrived; - the result stream closes early;
- a Gym snapshot that outlasts
rollout_timeout_s(fixed in feat(checkpoint): orchestrate Gym turn-level recovery #4266, see the parked-deadline comment on #4266); - this receipt GET fails after
await task— it sits outside the drain loop (L1884-L1892), so a 60 s_controltimeout or aGymControlRequestErrorraises out of the generator.
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}
|
|
||
|
|
||
| @dataclass(frozen=True) | ||
| class PendingCompletedExecutionAcknowledgement: |
There was a problem hiding this comment.
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.
| missing_agent_checkpoint_participation.append( | ||
| contract.participant.participant_name | ||
| ) | ||
| if "completed_result_acknowledgement" not in contract.features: |
There was a problem hiding this comment.
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.
| def __init__(self) -> None: | ||
| self._live: dict[tuple[str, int], GymActorExecution] = {} | ||
| self._frozen_checkpoint_id: str | None = None | ||
| self._frozen_membership: tuple[GymActorExecution, ...] = () |
There was a problem hiding this comment.
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.
| class GymExternalStorageReference(_StrictWireModel): | ||
| """One TQ staging row required by a parked Gym continuation.""" | ||
|
|
||
| schema_version: Literal[1] = GYM_CHECKPOINT_SCHEMA_VERSION |
There was a problem hiding this comment.
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, thenround-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.
| return False | ||
|
|
||
|
|
||
| class GymAcknowledgementSink: |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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:
| 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, |
There was a problem hiding this comment.
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 aRolloutRecoveryStateand returns.model_dump(mode="json"). It has to be JSON mode: the file is read withtorch.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 typeattempt_uuidasuuid.UUID(it saves as a string).from_state_dict()callsRolloutRecoveryState.model_validate(state)and then builds the dataclasses. Keep the rules that span records as plain code: the sibling count equalsexpected_generations, group ids and attempt uuids are unique, groups that share anadmission_idagree, and a sealed attempt has a reward.- For the ACK row, store
agent_nameplus aGymCompletionReceipt, as the ACK-type comment asks. - A field typed
Optional[X]with no default is required in pydantic v2 but may beNone; that is the missing-key check. Optional fields inside the reusedGymCompletionReceiptkeep Gym's ownNonedefaults.
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) |
There was a problem hiding this comment.
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:
RolloutAttemptRecord.attempt_idis read only by one test (test_rollout_recovery.py#L822, #L855).attempt_uuidis created (L385), saved (L1220), listed in_ATTEMPT_STATE_FIELDS(L93) and loaded (L1480-L1489, L1554). The only thing that reads it back is the load step's own duplicate-UUID check (L1342-L1346, L1487-L1489). feat(checkpoint): orchestrate Gym turn-level recovery #4266 adds no reader either.
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.
f381ab1 to
9c955eb
Compare
Signed-off-by: Anish Mahishi <amahishi@nvidia.com>
Signed-off-by: Anish Mahishi <amahishi@nvidia.com>
9c955eb to
a4e2306
Compare
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.
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:
restart_onlyresources 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]Main changes
restart_onlyresources.Correctness rules
PENDING_MODELcontinuations must point to real captured model lineage.Nis never reused as the new physical owner; recovery creates a tracked replacement attempt.Scope
Test plan
Targeted coverage includes:
tests/unit/environments/test_gym_checkpoint.pytests/unit/environments/test_nemo_gym.pytests/unit/environments/test_nemo_gym_checkpoint.pytests/unit/environments/test_nemo_gym_token_capture.pytests/unit/experience/test_rollout_manager.pytests/unit/experience/test_rollout_recovery.pytests/unit/data_plane/test_rollout_reassembler.pytests/unit/experience/test_rollout_generation_failures.pySuggested command:
Before review