nemo_automodel.components.speculative.streaming.store

View as Markdown

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 a StoreHandle the consumer must hand to release once 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

NameDescription
FeatureStoreAbstract transport every feature-store backend implements.
StoreHandleOpaque token the consumer must return to FeatureStore.release.
StoreHealthIntegers-only snapshot of store residency for backpressure decisions.

Functions

NameDescription
_next_handle_idMint a fresh StoreHandle.handle_id (module-level counter).

Data

_handle_id_counter

API

class nemo_automodel.components.speculative.streaming.store.FeatureStore()
Abstract

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.

nemo_automodel.components.speculative.streaming.store.FeatureStore.close() -> None
abstract

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.

nemo_automodel.components.speculative.streaming.store.FeatureStore.gc() -> int
abstract

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.

nemo_automodel.components.speculative.streaming.store.FeatureStore.get(
device: torch.device | str | None = None
abstract

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.

nemo_automodel.components.speculative.streaming.store.FeatureStore.health() -> nemo_automodel.components.speculative.streaming.store.StoreHealth
abstract

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.

nemo_automodel.components.speculative.streaming.store.FeatureStore.put(
sample_id: str,
tensors: typing.Mapping[str, torch.Tensor],
run_id: str,
schema_version: int,
target_model_version: str,
draft_weight_version: str,
num_tokens: int
abstract

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.

nemo_automodel.components.speculative.streaming.store.FeatureStore.release(
) -> None
abstract

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.

class nemo_automodel.components.speculative.streaming.store.StoreHandle(
store: 'FeatureStore',
sample_id: str,
handle_id: int = _next_handle_id()
)
Dataclass

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.

handle_id
int = field(default_factory=_next_handle_id)
ref
SampleRef
sample_id
str
store
'FeatureStore'
class nemo_automodel.components.speculative.streaming.store.StoreHealth(
resident_bytes: int,
capacity_bytes: int,
sample_count: int,
high_watermark_hit: bool,
low_watermark_hit: bool
)
Dataclass

Integers-only snapshot of store residency for backpressure decisions.

capacity_bytes
int

Configured hard cap from LocalFeatureStore (PR 3’s shared-dir / PR 4’s NCCL store report the analogous cap).

high_watermark_hit
bool

True iff resident_bytes >= high_watermark_bytes on the last health call. The queue pauses a producer that sees this transition.

low_watermark_hit
bool

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.

resident_bytes
int

Bytes currently held in the store across un-acked samples. Compared against capacity_bytes for the high/low watermark hysteresis.

sample_count
int

Number of un-acked samples currently held. Used for the second cap (sample count) the RFC §“Open questions” Q2 keeps alongside the byte backstop.

nemo_automodel.components.speculative.streaming.store._next_handle_id() -> int

Mint a fresh StoreHandle.handle_id (module-level counter).

nemo_automodel.components.speculative.streaming.store._handle_id_counter = itertools.count()