nemo_gym.token_id_capture
nemo_gym.token_id_capture
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
nemo_gym.token_id_capture.buildernemo_gym.token_id_capture.confignemo_gym.token_id_capture.conformancenemo_gym.token_id_capture.consumernemo_gym.token_id_capture.control_routesnemo_gym.token_id_capture.deliverynemo_gym.token_id_capture.external_capturenemo_gym.token_id_capture.fingerprintnemo_gym.token_id_capture.lineagenemo_gym.token_id_capture.protocolsnemo_gym.token_id_capture.recordsnemo_gym.token_id_capture.sinknemo_gym.token_id_capture.storenemo_gym.token_id_capture.terminal
Package Contents
Classes
Functions
Data
TOKEN_ENTRY_MIN_SCHEMA_VERSION
TOKEN_ENTRY_RECORD_SCHEMA_VERSION
API
Bases: _DigestWireModel
One token-free call manifest row in a rollout’s capture ledger.
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.
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.
Return whether any ledger row (committed or failed) exists.
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.
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.
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.
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.
Require one log probability for each generated token.
A trainer cannot use a chain with mismatched token and log-probability counts.
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.
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.
Publish one committed entry to the worker-local index.
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.
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.
Load several committed entries.
Backends can override this hook to fetch a parent chain in one operation.
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.
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.
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.
Describe a uniquely verified parent from a shared lineage store.
Return one immutable request-time parent decision.
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.
Release resources. Idempotent.
The store is read-only, so there is never pending work to flush.
Return whether separate model-server workers share committed entries.
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.
Bases: enum.Enum
Describe whether a model call has a proven captured predecessor.
Keep an append-only per-rollout call index.
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.
Index lookup metadata carried by one committed token entry.
store_tokens=False keeps token arrays in the durable log.
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.
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.
Return the parent decision without touching token arrays.
Matching needs only fingerprints, digests, and lengths. The caller materializes tokens for the single winner.
An immutable view of one rollout’s frozen capture records.
Durable, rollout-keyed JSONL sink for TokenEntry records.
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.
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.
Idempotently append one entry and fsync.
The file store owns no persistent handles.
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.
Delete snapshot payloads while retaining its tombstone and lock.
Synchronously freeze one rollout and return its stable snapshot.
Sentinel marking that at least one call of this rollout failed to capture.
Return the durable per-call intent path.
Durably record that a call was lost.
Store an entry durably without blocking the event loop.
Await the append so later consumers cannot race a partial file.
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.
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.
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.
Bases: BaseModel
The capture block plus the one top-level key it falls back to.
Require a resolver whenever a custom sink stores lineage.
A missing resolver makes every continuation unresolved. Current reconstruction refuses to guess across that boundary.
Construct the configured request-time lineage store.
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.
Receive captured records through Gym’s file store or a framework transport.
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.
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.
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.
Where a trajectory builder freezes, reads, and retires records.
Release resources idempotently.
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.
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.
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.
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.
Return worker-level capture health for metrics endpoints.
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.
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.
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.
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.
The full sequence a child of this call must start with.
Delta records require parent-chain reconstruction.
Return the capture context for the in-flight call.
Return None for untagged traffic.
Framework inference workers use this identity for staged records.
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.
Set (or clear) the process-wide request-time lineage store.
Set (or clear, with None) the process-wide default sink.
Set (or clear) the caller-owned source in this process.
Gym does not close an installed source.
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.
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.
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.
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.
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.
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_inadmission. - A request with no prior assistant output creates a
textroot. - An unresolved request may create a
textroot 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.
Chain frozen snapshot entries with the named strategy.
terminal_call_id anchors chain selection for prefix_merging.
Add compact lookup metadata before the token entry is committed.
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.
Return the enabled token store directory or an empty list.
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.
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.
Reject anything that could escape the store directory or index a bad file.