nemo_gym.token_id_capture.lineage

View as Markdown

Resolve the recorded call that a request continues.

A rollout can contain several model calls. Training consumes their exact tokens as one contiguous sequence. Request-time lineage identifies the earlier call that each request continues.

assistant_fingerprint is the lookup key. It hashes model-authored turns and ignores user and tool content added between calls. conversation_digest verifies the unchanged request context. A digest mismatch rejects the claimed lineage before any parent tokens are reused.

The shared LineageResolver resolves entries already committed by TokenSink. FileLineageStore tails the token JSONL through the token store’s lock. Each child receives its parent’s cumulative tokens. Downstream inference consumes those tokens to supply the exact prompt prefix.

Every supported record distinguishes a root, a resolved parent, and an unresolved boundary. The builder uses token-prefix matching only when a verified parent is absent from the frozen snapshot. It never uses prefix matching to cross an unresolved boundary.

A delivered chain contains exactly the tokens the policy emitted over the recorded context. The hashes ignore reasoning and selected items that a harness may omit when it echoes model output. These differences do not change the captured token sequence. Ambiguous matches remain unresolved rather than risking tokens from the wrong call.

Module Contents

Classes

NameDescription
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.
LineageNode-
RolloutLineageKeep an append-only per-rollout call index.

Functions

NameDescription
_custody_columnsReturn the ledger custody columns for one committed CallRecord.
_manifest_from_rowsBuild the token-free RolloutManifest payload from ledger rows.
stamp_continuationAdd compact lookup metadata before the token entry is committed.

Data

_CUSTODY_FIELDS

API

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.lineage.LineageNode(
call_id: str,
cum_tokens: list[int] | None,
cum_len: int,
digest: str,
entry_offset: int = -1,
context_len: int = 0,
context_digest: str = '',
parent_call_id: str | None = None,
prompt_is_delta: bool = False,
staging_key: str = '',
staging_chain: list[str] = list(),
chain_hash: str = ''
)
Dataclass
call_id
str
chain_hash
str = ''
context_digest
str = ''
context_len
int = 0
cum_len
int
cum_tokens
list[int] | None
digest
str
entry_offset
int = -1
parent_call_id
str | None = None
prompt_is_delta
bool = False
staging_chain
list[str] = field(default_factory=list)
staging_key
str = ''
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.

nemo_gym.token_id_capture.lineage._custody_columns(
staging_chain: tuple[str, ...] | list[str] = ()
) -> dict

Return the ledger custody columns for one committed CallRecord.

_manifest_from_rows rebuilds the CallRecord from these columns, so the mapping must stay a lossless round trip.

nemo_gym.token_id_capture.lineage._manifest_from_rows(
rollout_id: str,
rows: list[dict]
) -> dict

Build the token-free RolloutManifest payload from ledger rows.

Committed custody rows become CallRecord payloads; failure rows become failures entries. Lineage-only rows (local capture) carry no custody columns and are not part of a capture manifest.

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.lineage._CUSTODY_FIELDS = ('parent_call_id', 'staging_key', 'weight_version', 'prev_len', 'delta_len', 'cu...