nemo_gym.token_id_capture.consumer

View as Markdown

Turn a rollout’s frozen token capture into trajectories.

Gym rollout collection and trainer finalization use this consumer. Gym reads a frozen snapshot from the local token store. A trainer freezes the TokenSource provided by its transport. Both paths pass snapshot entries through the same build and projection.

This module does not import rollout-record or model-server modules. The caller supplies the rollout_id. Gym derives that ID from task, rollout, and attempt indices. The result includes metrics that describe the build.

Module Contents

Functions

NameDescription
_assemble-
_failed_build-
_incomplete_masksWhether an incomplete snapshot masks this build.
clear_token_captures_for_rolloutsRemove stale token records for rollouts about to be dispatched.
mask_incomplete_when_attributed_from_configReturn the token_id_capture.mask_incomplete_when_attributed setting.
token_id_capture_dirs_from_configReturn the enabled token store directory or an empty list.
trajectories_for_rolloutBuild trajectories from a frozen local token-store snapshot.
trajectories_from_sourceBuild trajectories from a frozen TokenSource snapshot.

Data

logger

API

nemo_gym.token_id_capture.consumer._assemble(
rollout_id: str,
builder: str,
model: str,
verified_response: dict | None = None,
explicit_terminal_call_id: str | None = None,
declared_response_id: str | None = None
) -> dict
nemo_gym.token_id_capture.consumer._failed_build(
rollout_id: str,
builder: str,
error: str,
n_calls: int = 0
) -> dict
nemo_gym.token_id_capture.consumer._incomplete_masks(
built: dict,
always: bool
) -> bool

Whether an incomplete snapshot masks this build.

Incomplete means some model call registered its capture intent and never committed a record (a harness killed at its timeout backstop leaves exactly this signature). With always the build is masked, because that call may be the trajectory’s real terminal.

Without always a delivered terminal attribution keeps the build: the attributed terminal is the verified response’s final model output, the delivered chain reaches it whole, and a call with no record cannot sit inside that chain, so the uncaptured call is outside the scored trajectory. Any other attribution outcome still masks.

nemo_gym.token_id_capture.clear_token_captures_for_rollouts(
records: list,
token_capture_dirs: list[pathlib.Path]
) -> None

Remove stale token records for rollouts about to be dispatched.

Rollout IDs are deterministic. TokenCaptureStore.append uses append mode. A reused ID would append records to a previous attempt. The builder could then combine two attempts. The caller passes only rows ready for dispatch. Retry suffixes must already be assigned.

nemo_gym.token_id_capture.mask_incomplete_when_attributed_from_config(
global_config_dict
) -> bool

Return the token_id_capture.mask_incomplete_when_attributed setting.

A caller that holds the Gym global config reads the setting here and passes it to the build functions below, which default to the strict behavior.

nemo_gym.token_id_capture.token_id_capture_dirs_from_config(
global_config_dict
) -> list[pathlib.Path]

Return the enabled token store directory or an empty list.

nemo_gym.token_id_capture.trajectories_for_rollout(
rollout_id: str,
token_capture_dirs: list[pathlib.Path],
builder: str = 'prefix_merging',
model: str = '',
verified_response: dict | None = None,
explicit_terminal_call_id: str | None = None,
mask_incomplete_when_attributed: bool = True
) -> dict | None

Build trajectories from a frozen local token-store snapshot.

Return None only when no capture directory is configured. Missing records are unsafe and return a masked result. An incomplete snapshot masks the result. With mask_incomplete_when_attributed set to False it masks only when terminal attribution did not deliver its chain, because a delivered chain places the uncaptured call off the scored path.

nemo_gym.token_id_capture.trajectories_from_source(
rollout_id: str,
builder: str = 'prefix_merging',
model: str = '',
verified_response: dict | None = None,
explicit_terminal_call_id: str | None = None,
mask_incomplete_when_attributed: bool = True,
declared_response_id: str | None = None
) -> dict | None
async

Build trajectories from a frozen TokenSource snapshot.

declared_response_id is the served response id the harness reports; a declared id that matches no entry masks the rollout.

Missing records are unsafe and return a masked result. An incomplete snapshot masks the result. With mask_incomplete_when_attributed set to False it masks only when terminal attribution did not deliver its chain, because a delivered chain places the uncaptured call off the scored path.

nemo_gym.token_id_capture.consumer.logger = logging.getLogger(__name__)