nemo_gym.token_id_capture

View as Markdown

Provide the core training-token capture interfaces.

Training capture is separate from evaluation capture. Middleware sets a request-scoped token sink. The model server records a TokenEntry from its complete response. Consumers call TokenSource.freeze for an atomic snapshot. The snapshot includes entries and incomplete state. Its snapshot_id identifies the exact frozen state. TokenCaptureStore is Gym’s local sink and source implementation. Framework transports may provide their own sink and source. This leaf package avoids imports from Gym’s server stack. The rollout-record finalizer needs Gym’s server stack. It is deliberately not re-exported here. Import nemo_gym.token_id_capture.delivery from server-side code. The incomplete state prevents training on a rollout that lost a model call. Finalization freezes and rebuilds the rollout. Finalization does not retire the snapshot. The caller retires it only after durable handoff. Retirement uses the frozen snapshot_id and version. Failed or masked builds retain their capture evidence.

Subpackages

Submodules

Package Contents

Classes

NameDescription
CallRecordOne token-free call manifest row in a rollout’s capture ledger.
CaptureContextDescribe one in-flight training-token capture.
CaptureLedgerStore metadata for calls whose token data is staged externally.
CaptureLedgerCommitOne successfully staged call, as handed to CaptureLedger.record.
Chain-
FileLineageStoreResolve lineage from the token JSONL committed by TokenCaptureStore.
InMemoryLineageStoreReference resolver for in-process framework backends and tests.
IncrementalLineageStoreBase class for lineage resolvers over any committed-entry backend.
LineageIndexBound worker-local lineage by rollout and cumulative token counts.
LineageMatchDescribe a uniquely verified parent from a shared lineage store.
LineageResolutionReturn one immutable request-time parent decision.
LineageResolverResolve request-time lineage from entries committed by a token sink.
ParentResolutionStatusDescribe whether a model call has a proven captured predecessor.
RolloutLineageKeep an append-only per-rollout call index.
TokenCaptureSnapshotAn immutable view of one rollout’s frozen capture records.
TokenCaptureStoreDurable, rollout-keyed JSONL sink for TokenEntry records.
TokenEntryStore one model call’s content and token metadata.
TokenIdCaptureConfigThe capture block plus the one top-level key it falls back to.
TokenSinkReceive captured records through Gym’s file store or a framework transport.
TokenSourceWhere a trajectory builder freezes, reads, and retires records.

Functions

NameDescription
assert_prefix_contiguityRequire each generated item to extend all preceding tokens.
assistant_fingerprintFingerprint the model-authored turns of a request, in order.
capture_health_snapshotReturn worker-level capture health for metrics endpoints.
capture_tokensRecord a TokenEntry from a complete model response.
clear_token_captures_for_rolloutsRemove stale token records for rollouts about to be dispatched.
commit_entryDurably record a finished entry against the in-flight call.
compute_digestDigest of an exact token sequence.
cumulative_tokensThe full sequence a child of this call must start with.
current_capture_contextReturn the capture context for the in-flight call.
extract_token_fieldsPull the token-id fields off a served response, or None if absent.
install_lineage_storeSet (or clear) the process-wide request-time lineage store.
install_token_sinkSet (or clear, with None) the process-wide default sink.
install_token_sourceSet (or clear) the caller-owned source in this process.
installed_lineage_store-
installed_token_sink-
installed_token_source-
make_token_storeBuild the training-token file store.
mark_external_staging_committedMark the current call as durably recorded by a framework worker.
prefix_merging-
project_chain_to_output_itemsProject a chain into Responses output items with contiguous prompts.
project_main_chain_responseRebuild the main chain as a Responses object whose output items are contiguous.
register_call_intentRecord durable call intent before dispatch starts generation.
reset_token_sink-
resolve_parentResolve which recorded call this request continues.
run_builderChain frozen snapshot entries with the named strategy.
set_token_sink-
stamp_continuationAdd compact lookup metadata before the token entry is committed.
stamp_lineageFill token lineage and the request-time parent decision.
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.
validate_rollout_idReject anything that could escape the store directory or index a bad file.

Data

NG_CAPTURE_FIELD

NG_COMMIT_COORDS_FIELD

TOKEN_ENTRY_MIN_SCHEMA_VERSION

TOKEN_ENTRY_RECORD_SCHEMA_VERSION

TOKEN_FIELDS

UNCOMMITTED_CALL_REASON

UNRESOLVED_PARENT_REASON

API

class nemo_gym.token_id_capture.CallRecord()

Bases: _DigestWireModel

One token-free call manifest row in a rollout’s capture ledger.

admitted_at
StrictFloat | None = None
chain_hash
DigestHex
continuation_fingerprint
DigestHex | None = None
cum_len
NonNegativeInt
cumulative_hash
DigestHex
delta_len
NonNegativeInt
digest
DigestHex
extras_digest
DigestHex
fingerprint_version
NonNegativeInt = 0
mode
CaptureMode = 'token_in'
model_call_id
Identifier
output_fingerprint
DigestHex | None = None
parent_call_id
Identifier | None = None
prev_len
NonNegativeInt
response_id
Identifier
staging_key
Identifier
weight_version
NonNegativeInt
nemo_gym.token_id_capture.CallRecord._validate_lengths() -> typing.Self
class nemo_gym.token_id_capture.CaptureContext(
rollout_id: str,
model_call_id: str,
model: str = '',
committed: bool = False,
delta_records: bool = False,
prefix_requested: bool = False,
prefix_supplied: bool = False,
external_staging: bool = False,
admitted_at: float | None = None,
parent_staging_chain: list[str] = list(),
parent_chain_hash: str = '',
request_items: list[dict] | None = None,
external_commit_coords: dict[str, typing.Any] | None = None,
external_worker_response_seen: bool = False
)
Dataclass

Describe one in-flight training-token capture.

The context identifies the rollout and model call. token_sink receives the resulting record. A framework may provide any TokenSink implementation. Parent resolution runs once for each call. Prefix supply and token capture read the same immutable decision.

admitted_at
float | None = None
capture_admission
CaptureAdmission | None = None
committed
bool = False
delta_records
bool = False
external_commit_coords
dict[str, Any] | None = None
external_staging
bool = False
external_worker_response_seen
bool = False
lineage_store
LineageResolver | CaptureLedger | None = None
model
str = ''
model_call_id
str
parent_call_id
str | None
parent_chain_hash
str = ''
parent_resolution
LineageResolution | None = None
parent_staging_chain
list[str] = field(default_factory=list)
parent_tokens
list[int]
prefix_requested
bool = False
prefix_supplied
bool = False
request_items
list[dict] | None = None
rollout_id
str
token_sink
TokenSink | None
class nemo_gym.token_id_capture.CaptureLedger()
Protocol

Bases: LineageResolver

Store metadata for calls whose token data is staged externally.

record publishes a successfully staged call for parent resolution. record_failure records a call that did not commit. manifest returns both kinds of rows without including token arrays. Records must be visible to every serving worker before record returns.

nemo_gym.token_id_capture.CaptureLedger.has_rows(
rollout_id: str
) -> bool
async

Return whether any ledger row (committed or failed) exists.

nemo_gym.token_id_capture.CaptureLedger.manifest(
rollout_id: str
) -> dict
async

Return the rollout’s token-free ledger as plain wire data.

The shape validates as staging.records.RolloutManifest: committed rows under records (each a CallRecord payload) and poison rows under failures. Cumulative token IDs never appear in the manifest.

nemo_gym.token_id_capture.CaptureLedger.record(
) -> None
async

Publish a completed call for later request-time resolution.

Repeating a model call ID with the same payload is a no-op. Reusing a model call ID with different data must fail. Return only after every serving worker can read the record.

nemo_gym.token_id_capture.CaptureLedger.record_failure(
rollout_id: str,
model_call_id: str,
reason: str
) -> None
async

Record a call whose token capture did not complete.

Failure rows do not participate in parent resolution. manifest includes them under failures so finalization rejects the rollout.

class nemo_gym.token_id_capture.CaptureLedgerCommit()

Bases: _WireModel

One successfully staged call, as handed to CaptureLedger.record.

record is the manifest row. The remaining fields feed parent resolution and are not part of the manifest.

record
CallRecord
request_items
list[dict]
response_items
list[dict]
rollout_id
Identifier
staging_chain
tuple[Identifier, ...] = ()
class nemo_gym.token_id_capture.Chain(
chain_id: str,
root_prompt: list[int] = list()
)
Dataclass
chain_id
str
links
list[ChainLink] = field(default_factory=list)
root_prompt
list[int] = field(default_factory=list)
nemo_gym.token_id_capture.Chain.validate() -> None

Require one log probability for each generated token.

A trainer cannot use a chain with mismatched token and log-probability counts.

class nemo_gym.token_id_capture.FileLineageStore(
root: str | pathlib.Path,
max_cached_rollouts: int = 65536,
max_cached_tokens: int = 8000000
)

Bases: IncrementalLineageStore

Resolve lineage from the token JSONL committed by TokenCaptureStore.

The reference IncrementalLineageStore backend: cursor = (inode, offset), ref = byte offset, reads under the store’s shared flock so a committed put is immediately visible.

_ledger_cache
dict[str, tuple[int, int, list[dict]]] = {}
_ledger_root
= Path(root)
_store
= TokenCaptureStore(root)
nemo_gym.token_id_capture.FileLineageStore._append(
rollout_id: str,
record: dict,
records: list[dict]
) -> None
nemo_gym.token_id_capture.FileLineageStore._fetch_new_entries(
rollout_id: str,
cursor: typing.Any
) -> tuple[list[tuple[nemo_gym.token_id_capture.records.TokenEntry, typing.Any]], typing.Any]
nemo_gym.token_id_capture.FileLineageStore._has_rows(
rollout_id: str
) -> bool
nemo_gym.token_id_capture.FileLineageStore._ledger_path(
rollout_id: str
) -> pathlib.Path
nemo_gym.token_id_capture.FileLineageStore._load_entries(
rollout_id: str,
refs: list[typing.Any]
nemo_gym.token_id_capture.FileLineageStore._load_entry(
rollout_id: str,
ref: typing.Any
nemo_gym.token_id_capture.FileLineageStore._locked(
rollout_id: str
)
nemo_gym.token_id_capture.FileLineageStore._manifest(
rollout_id: str
) -> dict
nemo_gym.token_id_capture.FileLineageStore._read(
rollout_id: str
) -> list[dict]
nemo_gym.token_id_capture.FileLineageStore._read_locked(
rollout_id: str
)
nemo_gym.token_id_capture.FileLineageStore._record(
) -> None
nemo_gym.token_id_capture.FileLineageStore._record_failure(
rollout_id: str,
model_call_id: str,
reason: str
) -> None
nemo_gym.token_id_capture.FileLineageStore._resolve(
rollout_id: str,
request_items: list[dict]
nemo_gym.token_id_capture.FileLineageStore._resolve_row(
rollout_id: str,
request_items: list[dict]
nemo_gym.token_id_capture.FileLineageStore.has_rows(
rollout_id: str
) -> bool
async
nemo_gym.token_id_capture.FileLineageStore.manifest(
rollout_id: str
) -> dict
async
nemo_gym.token_id_capture.FileLineageStore.record(
) -> None
async
nemo_gym.token_id_capture.FileLineageStore.record_failure(
rollout_id: str,
model_call_id: str,
reason: str
) -> None
async
class nemo_gym.token_id_capture.InMemoryLineageStore(
max_rollouts: int = 512,
max_tokens: int = 8000000
)

Reference resolver for in-process framework backends and tests.

Production wiring uses FileLineageStore when a token store exists. Its index is memory-only. Eviction or restart leaves affected continuations unresolved. That failure mode is safe but can mask otherwise usable rollouts. Multi-worker deployments require a shared LineageResolver. The resolution index evicts rollouts under memory bounds, so this store cannot serve as an external-staging capture ledger (completeness would break); it remains for unit tests and single-worker development. Its ledger rows are kept in a separate unbounded map so ledger unit tests see file-store semantics.

_ledgers
dict[str, list[dict]] = {}
index
nemo_gym.token_id_capture.InMemoryLineageStore.close() -> None
async
nemo_gym.token_id_capture.InMemoryLineageStore.has_rows(
rollout_id: str
) -> bool
async
nemo_gym.token_id_capture.InMemoryLineageStore.is_process_shared() -> bool
nemo_gym.token_id_capture.InMemoryLineageStore.manifest(
rollout_id: str
) -> dict
async
nemo_gym.token_id_capture.InMemoryLineageStore.put(
) -> None
async

Publish one committed entry to the worker-local index.

nemo_gym.token_id_capture.InMemoryLineageStore.record(
) -> None
async
nemo_gym.token_id_capture.InMemoryLineageStore.record_failure(
rollout_id: str,
model_call_id: str,
reason: str
) -> None
async
nemo_gym.token_id_capture.InMemoryLineageStore.resolve(
rollout_id: str,
request_items: list[dict]
async
class nemo_gym.token_id_capture.IncrementalLineageStore(
max_cached_rollouts: int = 65536,
max_cached_tokens: int = 8000000
)

Base class for lineage resolvers over any committed-entry backend.

An external backend implements two hooks. It inherits Gym’s matcher, bounded index, locking, and token materialization. Hash-for-hash agreement is the wire contract. The backend remains the source of truth when cache rows are evicted. A resolved match loads only the winning call’s token chain.

Required hooks: _fetch_new_entries(rollout_id, cursor) -> (items, new_cursor) where items is [(TokenEntry, ref), ...] in commit order since cursor (None means from the beginning) and ref is any handle that _load_entry can use later (byte offset, KV key, …). Raise CursorReset when the cursor no longer describes the backend (file rotated, namespace recreated); the base refetches from the beginning. _load_entry(rollout_id, ref) -> TokenEntry for one committed record.

Optional hooks: _load_entries(rollout_id, refs) — batch-load one parent chain (default: call _load_entry for each reference). _read_locked(rollout_id) — context manager held around fetch+resolve for backends with a read-lock discipline (default: no lock). is_process_shared() — default True; an external backend exists to be shared, and the multi-worker startup check trusts this answer.

_cache
dict[str, tuple[Any, dict[str, Any], RolloutLineage]] = {}
_cache_guard
= threading.Lock()
_materialized
dict[str, tuple[str, tuple[int, ...]]] = {}
_materialized_tokens
= 0
_rollout_locks
= tuple((threading.Lock()) for _ in (range(256)))
nemo_gym.token_id_capture.IncrementalLineageStore._cache_put(
rollout_id: str,
value: tuple[typing.Any, dict[str, typing.Any], nemo_gym.token_id_capture.lineage.RolloutLineage]
) -> None

Insert or touch a cache row with LRU semantics.

Reinsert a touched row so dictionary order tracks recency. Eviction only requires a later backend refetch.

nemo_gym.token_id_capture.IncrementalLineageStore._cached_materialized(
rollout_id: str
) -> tuple[str, tuple[int, ...]] | None
nemo_gym.token_id_capture.IncrementalLineageStore._fetch_new_entries(
rollout_id: str,
cursor: typing.Any
) -> tuple[list[tuple[nemo_gym.token_id_capture.records.TokenEntry, typing.Any]], typing.Any]
nemo_gym.token_id_capture.IncrementalLineageStore._load_entries(
rollout_id: str,
refs: list[typing.Any]

Load several committed entries.

Backends can override this hook to fetch a parent chain in one operation.

nemo_gym.token_id_capture.IncrementalLineageStore._load_entry(
rollout_id: str,
ref: typing.Any
nemo_gym.token_id_capture.IncrementalLineageStore._materialize(
rollout_id: str,
refs: dict[str, typing.Any],
) -> tuple[int, ...]

Load one RESOLVED parent’s cumulative tokens from the backend.

Read the chain in one batch and append each token segment once. Digest verification makes stale references fail closed.

nemo_gym.token_id_capture.IncrementalLineageStore._read_locked(
rollout_id: str
)
nemo_gym.token_id_capture.IncrementalLineageStore._refresh(
rollout_id: str
) -> tuple[dict[str, typing.Any], nemo_gym.token_id_capture.lineage.RolloutLineage]
nemo_gym.token_id_capture.IncrementalLineageStore._remember_materialized(
rollout_id: str,
call_id: str,
tokens: tuple[int, ...]
) -> None
nemo_gym.token_id_capture.IncrementalLineageStore._resolve(
rollout_id: str,
request_items: list[dict]
nemo_gym.token_id_capture.IncrementalLineageStore._rollout_lock(
rollout_id: str
)
nemo_gym.token_id_capture.IncrementalLineageStore.close() -> None
async
nemo_gym.token_id_capture.IncrementalLineageStore.is_process_shared() -> bool
nemo_gym.token_id_capture.IncrementalLineageStore.resolve(
rollout_id: str,
request_items: list[dict]
async
class nemo_gym.token_id_capture.LineageIndex(
max_rollouts: int = 512,
max_tokens: int = 8000000
)

Bound worker-local lineage by rollout and cumulative token counts.

This index backs the single-worker fallback. Shared stores provide cross-worker visibility. Eviction removes the oldest rollout. An evicted parent leaves later continuations unresolved and the builder masks them. The only live rollout is never evicted.

_rollouts
dict[str, RolloutLineage] = {}
total_tokens
int
nemo_gym.token_id_capture.LineageIndex.__len__() -> int
nemo_gym.token_id_capture.LineageIndex._evict() -> None
nemo_gym.token_id_capture.LineageIndex.clear() -> None
nemo_gym.token_id_capture.LineageIndex.drop(
rollout_id: str
) -> None

Release a rollout’s lineage early.

Gym’s model server has no rollout-completion signal. An in-process framework can call this when it retires the records.

nemo_gym.token_id_capture.LineageIndex.for_rollout(
rollout_id: str
class nemo_gym.token_id_capture.LineageMatch(
model_call_id: str,
cumulative_token_ids: tuple[int, ...],
digest: str,
staging_chain: tuple[str, ...] = (),
prev_len: int = 0,
chain_hash: str = ''
)
Dataclass

Describe a uniquely verified parent from a shared lineage store.

chain_hash
str = ''
cumulative_token_ids
tuple[int, ...]
digest
str
model_call_id
str
prev_len
int = 0
staging_chain
tuple[str, ...] = ()
class nemo_gym.token_id_capture.LineageResolution(
reason: str = ''
)
Dataclass

Return one immutable request-time parent decision.

match
LineageMatch | None = None
reason
str = ''
status
ParentResolutionStatus
nemo_gym.token_id_capture.LineageResolution.__post_init__() -> None
class nemo_gym.token_id_capture.LineageResolver()
Protocol

Resolve request-time lineage from entries committed by a token sink.

This is a read-only view over sink-committed records. After TokenSink.put returns, a later resolve on any worker must see the entry. Visibility is required per rollout key only, so the store may shard by rollout. Implementations must never guess among candidates. resolve must return UNRESOLVED when it cannot prove one parent. The offline builder independently re-verifies every claimed link by digest. External implementations should embed Gym’s RolloutLineage matcher rather than reimplement the hashing. nemo_gym.token_id_capture.lineage is importable without the server stack. Callers of commit_entry outside capture_tokens must run stamp_continuation themselves. Unstamped entries are invisible to resolution.

nemo_gym.token_id_capture.LineageResolver.close() -> None
async

Release resources. Idempotent.

The store is read-only, so there is never pending work to flush.

nemo_gym.token_id_capture.LineageResolver.is_process_shared() -> bool

Return whether separate model-server workers share committed entries.

nemo_gym.token_id_capture.LineageResolver.resolve(
rollout_id: str,
request_items: list[dict]
async

Return whether the request is a root, resolved, or unresolved.

request_items are the unmodified harness items. The implementation must verify the recorded request context. A conflicting set of committed payloads for one call id must count as zero candidates. UNRESOLVED is always a safe answer; a wrong RESOLVED is caught later by digest verification.

class nemo_gym.token_id_capture.ParentResolutionStatus

Bases: enum.Enum

Describe whether a model call has a proven captured predecessor.

RESOLVED
= 'resolved'
ROOT
= 'root'
UNRESOLVED
= 'unresolved'
class nemo_gym.token_id_capture.RolloutLineage(
by_fingerprint: dict[str, list[str]] = dict(),
by_call_id: dict[str, nemo_gym.token_id_capture.lineage.LineageNode] = dict(),
total_tokens: int = 0
)
Dataclass

Keep an append-only per-rollout call index.

by_call_id
dict[str, LineageNode] = field(default_factory=dict)
by_fingerprint
dict[str, list[str]] = field(default_factory=dict)
total_tokens
int = 0
nemo_gym.token_id_capture.RolloutLineage._continues(
messages: list[dict]
) -> bool
staticmethod

Return whether this request extends the node’s recorded context.

The leading context_len items must match the recorded request. A rewritten or summarized context fails verification. Verification excludes the model response because dialects can echo it as different item counts.

nemo_gym.token_id_capture.RolloutLineage.add_entry(
store_tokens: bool = True,
entry_offset: int = -1
) -> None

Index lookup metadata carried by one committed token entry.

store_tokens=False keeps token arrays in the durable log.

nemo_gym.token_id_capture.RolloutLineage.record(
call_id: str,
messages: list[dict],
cum_tokens: list[int],
digest: str,
context_len: int | None = None,
staging_key: str = '',
parent_staging_chain: list[str] | None = None,
cum_len: int | None = None,
chain_hash: str = ''
) -> None

Index a completed call by its continuation fingerprint.

context_len counts the request items before the model response. The default assumes one synthesized response item. cum_len must be passed explicitly for token-free custody rows, where cum_tokens is empty; a child’s prev_len reads it.

nemo_gym.token_id_capture.RolloutLineage.resolve(
messages: list[dict]

Return the immutable parent decision for this request.

A request without model-authored history is a root. A request with unverified history is unresolved. Never guess among calls with identical output.

nemo_gym.token_id_capture.RolloutLineage.resolve_node(
messages: list[dict]
) -> tuple[nemo_gym.token_id_capture.records.ParentResolutionStatus, 'LineageNode | None', str]

Return the parent decision without touching token arrays.

Matching needs only fingerprints, digests, and lengths. The caller materializes tokens for the single winner.

class nemo_gym.token_id_capture.TokenCaptureSnapshot(
rollout_id: str,
incomplete: bool,
snapshot_id: str,
version: int
)
Dataclass

An immutable view of one rollout’s frozen capture records.

entries
tuple[TokenEntry, ...]
incomplete
bool
rollout_id
str
snapshot_id
str
version
int
class nemo_gym.token_id_capture.TokenCaptureStore(
root: str | pathlib.Path
)

Durable, rollout-keyed JSONL sink for TokenEntry records.

_root
= Path(root)
root
Path
nemo_gym.token_id_capture.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.TokenCaptureStore._dangling_intents(
rollout_id: str,
) -> list[str]
nemo_gym.token_id_capture.TokenCaptureStore._drop(
rollout_id: str,
snapshot_id: str,
version: int
) -> bool
nemo_gym.token_id_capture.TokenCaptureStore._entry_digest(
payload: bytes
) -> str
staticmethod
nemo_gym.token_id_capture.TokenCaptureStore._fsync_root() -> None
nemo_gym.token_id_capture.TokenCaptureStore._locked(
rollout_id: str,
shared: bool = False
)
nemo_gym.token_id_capture.TokenCaptureStore._mark_incomplete(
rollout_id: str,
model_call_id: str = ''
) -> None
nemo_gym.token_id_capture.TokenCaptureStore._read_entries_unlocked(
rollout_id: str
nemo_gym.token_id_capture.TokenCaptureStore._read_state(
rollout_id: str
) -> dict[str, typing.Any]
nemo_gym.token_id_capture.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.TokenCaptureStore._write_state(
rollout_id: str,
state: dict[str, typing.Any],
durable: bool = True
) -> None
nemo_gym.token_id_capture.TokenCaptureStore.append(
) -> None

Idempotently append one entry and fsync.

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

The file store owns no persistent handles.

nemo_gym.token_id_capture.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.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.TokenCaptureStore.freeze(
rollout_id: str
async
nemo_gym.token_id_capture.TokenCaptureStore.freeze_now(
rollout_id: str

Synchronously freeze one rollout and return its stable snapshot.

nemo_gym.token_id_capture.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.TokenCaptureStore.intents_path_for(
rollout_id: str
) -> pathlib.Path

Return the durable per-call intent path.

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

Durably record that a call was lost.

nemo_gym.token_id_capture.TokenCaptureStore.path_for(
rollout_id: str
) -> pathlib.Path
nemo_gym.token_id_capture.TokenCaptureStore.put(
) -> 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.TokenCaptureStore.read_entries(
rollout_id: str
nemo_gym.token_id_capture.TokenCaptureStore.state_path_for(
rollout_id: str
) -> pathlib.Path
nemo_gym.token_id_capture.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.

class nemo_gym.token_id_capture.TokenEntry()

Bases: BaseModel

Store one model call’s content and token metadata.

The rollout id identifies the training sample. The model call id joins evaluation context. output_items preserves assistant text and tool calls. Text-based penalties require that content. Token arrays are stored once at the top level. token_item_index identifies their original output item. A trajectory builder can restore chain-correct token fields there.

continuation_context_digest
str = ''
continuation_context_len
int = 0
continuation_fingerprint
str = ''
created_at
float = 0.0
cum_len
int | None = None
digest
str | None = None
fingerprint_version
int | None = None
generation_log_probs
list[float]
generation_token_ids
list[int]
model
str = ''
model_call_id
str
model_config
= ConfigDict(extra='allow')
output_items
list[dict] = Field(default_factory=list)
parent_call_id
str | None = None
parent_resolution
ParentResolutionStatus | None = None
parent_resolution_reason
str = ''
prefix_requested
bool = False
prefix_supplied
bool = False
prompt_is_delta
bool = False
prompt_token_ids
list[int]
response_id
str | None = None
rollout_id
str
routed_experts
Any | None = None
schema_version
int = TOKEN_ENTRY_RECORD_SCHEMA_VERSION
token_item_index
int | None = None
nemo_gym.token_id_capture.TokenEntry._refuse_a_newer_record() -> 'TokenEntry'

Accept older records and reject newer records.

Missing older fields use their defaults. Unknown newer fields may change token semantics. Rejecting them prevents silent training corruption.

class nemo_gym.token_id_capture.TokenIdCaptureConfig()

Bases: BaseModel

The capture block plus the one top-level key it falls back to.

enabled
bool
model_call_capture_dir
Path | None = None
model_config
= ConfigDict(extra='ignore')
token_id_capture
TokenIdCaptureSettings = TokenIdCaptureSettings()
nemo_gym.token_id_capture.TokenIdCaptureConfig._build_endpoint(
target: str,
kwargs: dict[str, typing.Any],
protocol: type,
kind: str
)
staticmethod
nemo_gym.token_id_capture.TokenIdCaptureConfig._require_resolver(
) -> None
staticmethod

Require a resolver whenever a custom sink stores lineage.

A missing resolver makes every continuation unresolved. Current reconstruction refuses to guess across that boundary.

nemo_gym.token_id_capture.TokenIdCaptureConfig._validate() -> 'TokenIdCaptureConfig'
nemo_gym.token_id_capture.TokenIdCaptureConfig.build_lineage_store() -> nemo_gym.token_id_capture.protocols.LineageResolver | None

Construct the configured request-time lineage store.

nemo_gym.token_id_capture.TokenIdCaptureConfig.build_sink() -> nemo_gym.token_id_capture.protocols.TokenSink | None

Construct the configured sink.

Return None when the file store is in use. Call this once in each server process. Launcher-installed sinks do not reach spawned workers.

nemo_gym.token_id_capture.TokenIdCaptureConfig.resolved_dir() -> pathlib.Path | None
class nemo_gym.token_id_capture.TokenSink()
Protocol

Receive captured records through Gym’s file store or a framework transport.

nemo_gym.token_id_capture.TokenSink.close() -> None
async

Release resources. Idempotent.

There is never buffered unwritten data here. put guaranteed durability before it returned. A close that must flush records means put broke its contract.

nemo_gym.token_id_capture.TokenSink.mark_incomplete(
rollout_id: str,
model_call_id: str = ''
) -> None
async

Durably record that a call of this rollout failed to capture.

The rollout is now missing a turn. A consumer must mask the sample instead of training on a chain with a hole. The model call itself still succeeds. This marker is therefore the durable signal that capture failed. It must succeed after freeze. It must change the observable version; that is what invalidates a stale retirement. Make it more available than put, for example through a local spill. put and mark_incomplete failing together is the silent-loss case.

nemo_gym.token_id_capture.TokenSink.put(
) -> None
async

Durably store one record before returning.

The entry carries its continuation lookup metadata. Durability and resolver visibility are the return condition, not an eventual goal. Any worker’s paired lineage resolver must see the entry once this method returns. Repeating the same call id with the same payload is a no-op. “Same payload” means the identical serialized entry, byte-for-byte, timestamps included. A retry must resend the same bytes; rebuilding the entry produces a conflict, not a retry. Reusing a call id with a different payload must fail. A transport without compare-and-swap may delegate that conflict to the reader. A resolver must then treat conflicting committed payloads for one call id as zero candidates. Writing after freeze must fail or bump the frozen version; see TokenSource.freeze.

This method may raise. The caller marks the rollout incomplete. A capture error never fails the model call.

class nemo_gym.token_id_capture.TokenSource()
Protocol

Where a trajectory builder freezes, reads, and retires records.

nemo_gym.token_id_capture.TokenSource.close() -> None
async

Release resources idempotently.

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

Conditionally retire the exact frozen snapshot that was consumed.

Return False if state changed after the snapshot. Transports without delete return True and own retention.

nemo_gym.token_id_capture.TokenSource.freeze(
rollout_id: str
async

Freeze a rollout and return one atomic snapshot.

Freezing is idempotent. Entry order carries no meaning. Entries are unique per model_call_id. An at-least-once transport must dedupe identical copies before snapshotting. The fence is relaxed: “no successful write after freeze” is not required cluster-wide. A strict fence is unimplementable without compare-and-swap. A write racing freeze may therefore succeed durably. It must then bump the version. A conditional drop of the consumed snapshot then fails, and the evidence is retained. Attempt-scoped rollout ids are the sanctioned strategy for retirement without compare-and-swap.

nemo_gym.token_id_capture.assert_prefix_contiguity(
response: dict
) -> None

Require each generated item to extend all preceding tokens.

The preceding tokens include the prompt and all prior generations. Raise AssertionError when the response is not contiguous.

nemo_gym.token_id_capture.assistant_fingerprint(
messages: list[dict]
) -> str

Fingerprint the model-authored turns of a request, in order.

The fingerprint identifies the call that produced the last model-authored turn. User and tool content is excluded from the lookup key. Dialect-specific tool-call shapes normalize to the same hash input.

nemo_gym.token_id_capture.capture_health_snapshot() -> dict

Return worker-level capture health for metrics endpoints.

nemo_gym.token_id_capture.capture_tokens(
response: typing.Any,
request_messages: list | None = None
) -> None
async

Record a TokenEntry from a complete model response.

Accept a Pydantic model or dictionary. Return without work when no capture context exists. Mark local capture incomplete when required token ids are absent. Await the write before the model call returns.

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.commit_entry(
) -> None
async

Durably record a finished entry against the in-flight call.

capture_tokens extracts arrays from a served response. Engine-side capture may already have those arrays. Engine-side callers can use this method directly. Return without work when no capture context exists. Capture failures mark the rollout incomplete. This method never fails the model call.

nemo_gym.token_id_capture.compute_digest(
token_ids: list[int]
) -> str

Digest of an exact token sequence.

The builder verifies a claimed parent by hashing the corresponding prompt prefix. A mismatch quarantines the call. This prevents stale or interleaved records from merging silently.

nemo_gym.token_id_capture.cumulative_tokens(
) -> list[int]

The full sequence a child of this call must start with.

Delta records require parent-chain reconstruction.

nemo_gym.token_id_capture.current_capture_context() -> nemo_gym.token_id_capture.sink.CaptureContext | None

Return the capture context for the in-flight call.

Return None for untagged traffic. Framework inference workers use this identity for staged records.

nemo_gym.token_id_capture.extract_token_fields(
response_json: dict
) -> dict | None

Pull the token-id fields off a served response, or None if absent.

Handle Responses output items and Chat Completions messages. Exactly one item may carry token metadata. Return None when no item carries token ids.

nemo_gym.token_id_capture.install_lineage_store(
) -> None

Set (or clear) the process-wide request-time lineage store.

nemo_gym.token_id_capture.install_token_sink(
) -> None

Set (or clear, with None) the process-wide default sink.

nemo_gym.token_id_capture.install_token_source(
) -> None

Set (or clear) the caller-owned source in this process.

Gym does not close an installed source.

nemo_gym.token_id_capture.installed_lineage_store() -> nemo_gym.token_id_capture.protocols.LineageResolver | None
nemo_gym.token_id_capture.installed_token_sink() -> nemo_gym.token_id_capture.protocols.TokenSink | None
nemo_gym.token_id_capture.installed_token_source() -> nemo_gym.token_id_capture.protocols.TokenSource | None
nemo_gym.token_id_capture.make_token_store(
global_config_dict: typing.Any

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.mark_external_staging_committed(
rollout_id: str,
model_call_id: str
) -> None

Mark the current call as durably recorded by a framework worker.

Call this only after the external staging sink has acknowledged the call. Identity validation prevents a delayed or cross-request acknowledgement from suppressing normal capture for a different request.

nemo_gym.token_id_capture.prefix_merging(
terminal_call_id: str | None = None
nemo_gym.token_id_capture.project_chain_to_output_items(
) -> list[dict]

Project a chain into Responses output items with contiguous prompts.

Preserve captured assistant text and tool calls. Put the contiguous prompt on the item that carries each generation. Each generated item’s prompt extends the previous generated item. Preserve text for downstream scoring. Create a token-only item only when a call has no captured content items.

nemo_gym.token_id_capture.project_main_chain_response(
rollout_id: str,
model: str = ''
) -> dict

Rebuild the main chain as a Responses object whose output items are contiguous.

The result is a Gym-native Responses payload. It contains object: "response", output items, and usage. Token fields describe one unbroken sequence across the rollout. The sequence combines items from multiple model calls.

nemo_gym.token_id_capture.register_call_intent() -> None
async

Record durable call intent before dispatch starts generation.

begin_call is an optional sink extension. A dangling intent identifies a lost entry. Failure happens before generation and propagates to the caller. The harness can retry without spending inference compute. Sinks without begin_call cannot report a missing final entry this way.

nemo_gym.token_id_capture.reset_token_sink(
token: contextvars.Token
) -> None
nemo_gym.token_id_capture.resolve_parent(
request_messages: list | None
) -> None
async

Resolve which recorded call this request continues.

Use the request representation received from the harness. Resolve once before dialect conversion or dispatch. Prefix supply and capture then share one parent decision. Return without work for untagged traffic. Every attempted resolution records a root, resolved, or unresolved decision. An unresolved decision includes its reason.

For external staging, parent resolution determines whether the worker may capture the call:

  • A unique parent creates a token_in admission.
  • A request with no prior assistant output creates a text root.
  • An unresolved request may create a text root only when the rollout has no ledger rows.
  • Every other result records a failure and leaves the call unadmitted.

An unresolved continuation cannot become a new root. Doing so would train the earlier generated tokens as prompt tokens.

nemo_gym.token_id_capture.run_builder(
builder: str = 'prefix_merging',
terminal_call_id: str | None = None

Chain frozen snapshot entries with the named strategy.

terminal_call_id anchors chain selection for prefix_merging.

nemo_gym.token_id_capture.set_token_sink(
) -> contextvars.Token
nemo_gym.token_id_capture.stamp_continuation(
request_items: list[dict]

Add compact lookup metadata before the token entry is committed.

nemo_gym.token_id_capture.stamp_lineage(
parent_call_id: str | None,
cumulative: list[int] | None = None

Fill token lineage and the request-time parent decision.

cum_len and digest describe the full sequence. A delta entry must pass cumulative explicitly. parent_resolution=None preserves records built by compatibility callers.

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
) -> 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.trajectories_from_source(
rollout_id: str,
builder: str = 'prefix_merging',
model: str = '',
verified_response: dict | None = None,
explicit_terminal_call_id: str | None = None,
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 is unsafe and returns a masked result.

nemo_gym.token_id_capture.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.sink.NG_CAPTURE_FIELD = 'ng_capture'
nemo_gym.token_id_capture.sink.NG_COMMIT_COORDS_FIELD = 'ng_commit_coords'
nemo_gym.token_id_capture.records.TOKEN_ENTRY_MIN_SCHEMA_VERSION = 1
nemo_gym.token_id_capture.records.TOKEN_ENTRY_RECORD_SCHEMA_VERSION = 1
nemo_gym.token_id_capture.records.TOKEN_FIELDS = ('prompt_token_ids', 'generation_token_ids', 'generation_log_probs', 'routed_exp...
nemo_gym.token_id_capture.records.UNCOMMITTED_CALL_REASON = 'request_finished_without_staged_coordinates'
nemo_gym.token_id_capture.records.UNRESOLVED_PARENT_REASON = 'unresolved_parent'