nemo_automodel.components.speculative.streaming.store
nemo_automodel.components.speculative.streaming.store
Pluggable data-plane transport for the speculative-training stream.
The FeatureStore ABC abstracts where the supervision tensors actually
live — an in-process dict for build/test, a POSIX shared mount for
multi-node/colocated, or NCCL for GPU-to-GPU in later PRs. The
SampleRefQueue reads the store’s FeatureStore.health to
decide whether to back off.
The contract is deliberately small (5 methods + 1 property) so PR 2 and later have an obvious surface to extend:
put— produce-side: stash tensors for a sample.get— consume-side: materialize them; returns aStoreHandlethe consumer must hand toreleaseonce it’s done with them.release— consume-once: free / drop the materials.gc— sweep partially-released handles (a stale lease, a crashed consumer) so the store cannot leak.health— ints only; the queue uses these for backpressure.close— dispose store resources at shutdown.
Module Contents
Classes
Functions
Data
API
Abstract transport every feature-store backend implements.
Implementations must be safe to call from one thread per operation;
cross-thread concurrency is the queue’s responsibility and outside this
contract. The SampleRefQueue
uses health for backpressure; everything else is the
produce/consume pair plus release / gc for lifecycle.
Release every resource owned by the store (file handles, NCCL groups, …).
After close, all subsequent calls raise RuntimeError.
Idempotent with respect to an already-closed store.
Sweep stale entries (e.g. failed releases from a crashed consumer).
Returns the number of entries reclaimed. Called opportunistically by
the SampleRefQueue between leases and unconditionally by
close.
Materialize ref’s tensors on device and hand back a StoreHandle.
The returned tensors are detached copies (clone for cuda
tensors, plain detach for CPU views) so prefetch cannot observe
aliasing through release. Returns one tensor per key in
SampleRef.feature_keys, in the same insertion order, so the
consumer can wire them straight into the per-algorithm batch.
Return a ints-only StoreHealth snapshot for backpressure.
Must not block on tensor I/O (no .cpu(), no .to()); in-memory
counters are sufficient. Backends that do background I/O MUST serve
health from cached counters, not from the in-flight I/O
thread.
Stash tensors for sample_id and return a tensor-free SampleRef.
The producer-side metadata (run_id, algorithm, schema_version,
target_model_version, draft_weight_version, num_tokens) is
the data the SampleRef carries on the control plane, so the
store must accept it here even though it does not inspect the values
beyond building the ref. Consumers see the ref unchanged.
Implementations must validate that the tensors match what they
declare (dtype, shape per feature, numel * element_size
summed across features == ref.estimated_bytes) and reject the
put with a specific exception when they do not, before any partial
write is observable.
Free the resources backing handle; idempotent.
After this returns, the sample is no longer present in the store and
a second get on the same SampleRef raises
KeyError. gc retries releases that a previous call
rejected (e.g. transient I/O), so the queue can rely on
gc + release being idempotent + retriable.
Opaque token the consumer must return to FeatureStore.release.
Each FeatureStore.get mints a fresh handle with a unique
handle_id. Two get calls on the same sample return two
distinct handles; FeatureStore.release matches against
handle_id so releasing one handle twice (or releasing a stale
handle after a sibling has been acquired) cannot decrement a
sibling’s outstanding count.
Holds the producing store, the sample id, the originating
SampleRef, and the per-get handle identity.
Integers-only snapshot of store residency for backpressure decisions.
Configured hard cap from
LocalFeatureStore
(PR 3’s shared-dir / PR 4’s NCCL store report the analogous cap).
True iff resident_bytes >= high_watermark_bytes
on the last health call. The queue pauses a producer that
sees this transition.
True iff resident_bytes <= low_watermark_bytes
on the last health call. The queue resumes a producer
that has been paused and now sees this transition. The
hysteresis band between the two is what prevents flapping.
Bytes currently held in the store across un-acked
samples. Compared against capacity_bytes for the high/low
watermark hysteresis.
Number of un-acked samples currently held. Used for the second cap (sample count) the RFC §“Open questions” Q2 keeps alongside the byte backstop.
Mint a fresh StoreHandle.handle_id (module-level counter).