nemo_gym.token_id_capture.store

View as Markdown

Store training TokenEntry records by rollout.

Each rollout uses one <rollout_id>.tokens.jsonl file. Evaluation records use a separate file. Every entry line is fsynced before put returns — that is the durability guarantee. The state index is written atomically but fsynced only on lifecycle transitions (freeze, mark, drop); it is reconstructible from the JSONL tail. A per-rollout file lock serializes writers to the same rollout. Different rollouts can write concurrently.

Module Contents

Classes

NameDescription
TokenCaptureStoreDurable, rollout-keyed JSONL sink for TokenEntry records.

Functions

NameDescription
make_token_storeBuild the training-token file store.
validate_rollout_idReject anything that could escape the store directory or index a bad file.

Data

logger

API

class nemo_gym.token_id_capture.store.TokenCaptureStore(
root: str | pathlib.Path
)

Durable, rollout-keyed JSONL sink for TokenEntry records.

_root
= Path(root)
root
Path
nemo_gym.token_id_capture.store.TokenCaptureStore._begin_call(
rollout_id: str,
model_call_id: str
) -> None

Durably record that a captured call is about to be dispatched.

A lost entry leaves a dangling intent. freeze_now then masks the rollout. A failure here happens before generation.

nemo_gym.token_id_capture.store.TokenCaptureStore._dangling_intents(
rollout_id: str,
entries: tuple[nemo_gym.token_id_capture.records.TokenEntry, ...]
) -> list[str]
nemo_gym.token_id_capture.store.TokenCaptureStore._drop(
rollout_id: str,
snapshot_id: str,
version: int
) -> bool
nemo_gym.token_id_capture.store.TokenCaptureStore._entry_digest(
payload: bytes
) -> str
staticmethod
nemo_gym.token_id_capture.store.TokenCaptureStore._fsync_root() -> None
nemo_gym.token_id_capture.store.TokenCaptureStore._locked(
rollout_id: str,
shared: bool = False
)
nemo_gym.token_id_capture.store.TokenCaptureStore._mark_incomplete(
rollout_id: str,
model_call_id: str = ''
) -> None
nemo_gym.token_id_capture.store.TokenCaptureStore._read_entries_unlocked(
rollout_id: str
) -> list[nemo_gym.token_id_capture.records.TokenEntry]
nemo_gym.token_id_capture.store.TokenCaptureStore._read_state(
rollout_id: str
) -> dict[str, typing.Any]
nemo_gym.token_id_capture.store.TokenCaptureStore._sync_entry_index(
rollout_id: str,
state: dict[str, typing.Any]
) -> bool

Reconcile an entry index with any durable JSONL tail.

The JSONL write is durable before its state update. A process can therefore stop with one unindexed entry. Normal writes use the state index without parsing prior token arrays. Recovery parses only the unindexed tail.

nemo_gym.token_id_capture.store.TokenCaptureStore._write_state(
rollout_id: str,
state: dict[str, typing.Any],
durable: bool = True
) -> None
nemo_gym.token_id_capture.store.TokenCaptureStore.append(
entry: nemo_gym.token_id_capture.records.TokenEntry
) -> None

Idempotently append one entry and fsync.

nemo_gym.token_id_capture.store.TokenCaptureStore.begin_call(
rollout_id: str,
model_call_id: str
) -> None
async
nemo_gym.token_id_capture.store.TokenCaptureStore.close() -> None
async

The file store owns no persistent handles.

nemo_gym.token_id_capture.store.TokenCaptureStore.delete(
rollout_id: str
) -> None

Unconditionally remove a rollout’s records.

This compatibility helper supports administrative cleanup. Normal consumers use conditional drop. The lock file remains so concurrent callers keep using one inode.

nemo_gym.token_id_capture.store.TokenCaptureStore.drop(
rollout_id: str,
snapshot_id: str,
version: int
) -> bool
async

Delete snapshot payloads while retaining its tombstone and lock.

nemo_gym.token_id_capture.store.TokenCaptureStore.freeze(
rollout_id: str
) -> nemo_gym.token_id_capture.protocols.TokenCaptureSnapshot
async
nemo_gym.token_id_capture.store.TokenCaptureStore.freeze_now(
rollout_id: str
) -> nemo_gym.token_id_capture.protocols.TokenCaptureSnapshot

Synchronously freeze one rollout and return its stable snapshot.

nemo_gym.token_id_capture.store.TokenCaptureStore.incomplete_path_for(
rollout_id: str
) -> pathlib.Path

Sentinel marking that at least one call of this rollout failed to capture.

nemo_gym.token_id_capture.store.TokenCaptureStore.intents_path_for(
rollout_id: str
) -> pathlib.Path

Return the durable per-call intent path.

nemo_gym.token_id_capture.store.TokenCaptureStore.is_incomplete(
rollout_id: str
) -> bool
nemo_gym.token_id_capture.store.TokenCaptureStore.lock_path_for(
rollout_id: str
) -> pathlib.Path
nemo_gym.token_id_capture.store.TokenCaptureStore.mark_incomplete(
rollout_id: str,
model_call_id: str = ''
) -> None
async

Durably record that a call was lost.

nemo_gym.token_id_capture.store.TokenCaptureStore.path_for(
rollout_id: str
) -> pathlib.Path
nemo_gym.token_id_capture.store.TokenCaptureStore.put(
entry: nemo_gym.token_id_capture.records.TokenEntry
) -> None
async

Store an entry durably without blocking the event loop.

Await the append so later consumers cannot race a partial file.

nemo_gym.token_id_capture.store.TokenCaptureStore.read_entries(
rollout_id: str
) -> list[nemo_gym.token_id_capture.records.TokenEntry]
nemo_gym.token_id_capture.store.TokenCaptureStore.state_path_for(
rollout_id: str
) -> pathlib.Path
nemo_gym.token_id_capture.store.TokenCaptureStore.sweep_retired(
older_than_seconds: float
) -> int

Remove retired tombstones older than the cutoff and return the count removed.

Callers choose the retention policy. drop already removed entries and JSONL payloads. This removes state, locks, intents, and incomplete markers.

nemo_gym.token_id_capture.store.make_token_store(
global_config_dict: typing.Any
) -> nemo_gym.token_id_capture.store.TokenCaptureStore | None

Build the training-token file store.

Return None when capture is disabled. Return None when no directory resolves. Return None when a custom sink owns the records.

nemo_gym.token_id_capture.store.validate_rollout_id(
rollout_id: str
) -> str

Reject anything that could escape the store directory or index a bad file.

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