nemo_automodel.components.speculative.streaming.stores
nemo_automodel.components.speculative.streaming.stores
Pluggable feature-store backends for the streaming data plane.
LocalFeatureStore is the in-process reference backend.
SharedDirFeatureStore adds cross-process rendezvous on a shared
POSIX mount. A future NcclFeatureStore can plug in behind the same
FeatureStore
contract.
Submodules
nemo_automodel.components.speculative.streaming.stores.localnemo_automodel.components.speculative.streaming.stores.shared_dir
Package Contents
Classes
API
Bases: FeatureStore
In-process FeatureStore implementation.
Thread safety: every public method holds a single threading.Lock,
so concurrent puts and gets from the same Python process are safe. Async
/ cross-process safety is the queue’s responsibility and is out of scope
for PR 1.
Parameters:
Hard cap on simultaneously-stored samples. None means
unbounded sample count (still bounded by max_bytes).
Hard cap on resident bytes. None means unbounded
(still bounded by max_samples). At least one of
max_samples / max_bytes must be set, otherwise a
misconfigured store silently behaves as unbounded.
Threshold above which StoreHealth.high_watermark_hit
is True. The producer pauses here.
Threshold below which StoreHealth.low_watermark_hit
is True. The producer resumes here. Must be
strictly less than high_watermark_bytes; a hysteresis band
of zero flaps the producer on every step.
Materialize ref’s features on device and hand back a StoreHandle.
Parameters:
The reference returned by put (typically via a
queue lease). ref.store_uri MUST equal this store’s
store_uri; a mismatch raises KeyError so a
consumer cannot accidentally materialize a foreign ref.
Optional target device. None returns each feature
on the device it was put on; a non-None value
materializes every feature on that device via
Tensor.to(device) (a no-op when already in place).
Returns: dict[str, torch.Tensor]
A (tensors, handle) pair. tensors is a
Raises:
KeyError: whenref.store_uridoes not match this store, or whenref.sample_idis no longer present (released or never put).RuntimeError: when the stored tensor’s shape or dtype differs from what the ref claims, or the store has been closed.
Store tensors under sample_id and return a tensor-free SampleRef.
Parameters:
Stable identifier within run_id. Must be unique
in this store at put time; duplicates raise ValueError.
Feature-name to tensor mapping. The store detaches,
clones, and makes each tensor contiguous before stashing it,
so the producer may keep mutating its source tensors after
the put returns without disturbing what a later
get hands out. The shape and dtype of each tensor
are captured into the returned SampleRef’s
feature_specs; the consumer uses those specs to
preallocate the receive buffer at get time, so
changing tensors[name].shape or dtype between put
and get without updating the ref will surface as a
RuntimeError on materialization.
Same value on every ref of one run; surfaces on the
SampleRef.run_id so producers and consumers can
verify they are talking about the same run.
Which draft family produced this sample; gates the
SampleRef required-features check.
Bumped whenever the producer’s feature set or
layout for algorithm changes incompatibly.
Monotonically increasing identifier of the target-model weights.
Same idea for the draft model’s weights.
Sum of attended tokens; used by the consumer for empty / short loss-mask neutralization.
Returns: SampleRef
A tensor-free SampleRef carrying the per-feature
Raises:
MemoryError: if the put would exceedmax_samplesormax_bytes. The producer is expected to retry after the store drains below the low watermark (seehealth).RuntimeError: if the store has been closed.ValueError: on bad input (empty sample id, empty tensors map, duplicate sample id).
Bases: FeatureStore
FeatureStore backed by one <sample_id>.safetensors per sample.
Thread safety: every public method holds a single
threading.Lock. Cross-process / cross-node is supported
as long as distinct processes use distinct sample_id values;
the lock does not extend across processes.
Parameters:
Filesystem path used as the rendezvous. Created if it
does not exist. Concurrent producers and consumers in
separate processes / ranks coordinate via unique
sample_id values — collision is the caller’s problem.
Same residency contract as LocalFeatureStore; the
queue’s HWM/LWM hysteresis reads them off
health.