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. Single-response delivery rejects per_request because it can return multiple trajectories.

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-
clear_token_captures_for_rolloutsRemove stale token records for rollouts about to be dispatched.
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,
entries: list[nemo_gym.token_id_capture.records.TokenEntry],
builder: str,
model: str
) -> 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.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.consumer.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.consumer.trajectories_for_rollout(
rollout_id: str,
token_capture_dirs: list[pathlib.Path],
builder: str = 'prefix_merging',
model: str = ''
) -> 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 is unsafe and returns a masked result.

nemo_gym.token_id_capture.consumer.trajectories_from_source(
rollout_id: str,
source: nemo_gym.token_id_capture.protocols.TokenSource,
builder: str = 'prefix_merging',
model: str = ''
) -> dict | None
async

Build trajectories from a frozen TokenSource snapshot.

Missing records are unsafe and return a masked result. An incomplete snapshot is unsafe and returns a masked result.

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