nemo_automodel.components.speculative.streaming.queue

View as Markdown

Lease / ack / fail queue over SampleRef for the streaming pipeline.

The SampleRefQueue carries only references — no tensors — between a producer (target-side forward) and a consumer (draft-side trainer). Each “message” is a SampleRef and is delivered exactly once: a consumer leases a ref, materializes its tensors via the FeatureStore, and then ACKs (release the lease and let the data-plane scrub the sample from the store) or FAILs (release the lease without scrubbing so a future consumer may retry).

A lease that is never ACK’d or FAIL’d within Lease.visibility_timeout is considered orphaned and is reclaimed by SampleRefQueue.reclaim_expired. That reclaim makes the queue safe to drive against a producer that may crash mid-flight.

Backpressure is driven by the bound FeatureStore’s FeatureStore.health (ints only — the queue never touches tensors in its hot path), with a high/low watermark hysteresis band so a fast producer cannot OOM the store and a slow producer cannot starve the trainer silently. The producer-side and consumer-side pause / resume transitions are tracked on the store via the same StoreHealth.high_watermark_hit / StoreHealth.low_watermark_hit flags, so a third party (e.g. an ops dashboard) can observe which side of the pipeline is the bottleneck without inspecting the queue internals.

Module Contents

Classes

NameDescription
LeaseHandle to a leased SampleRef.
SampleRefQueueLease / ack / fail queue over SampleRef.
VisibilityTimeoutHow long an unacked Lease is allowed to live before reclaim.

Functions

NameDescription
_next_lease_idMint a fresh Lease.lease_id (module-level counter).

Data

_lease_id_counter

logger

API

class nemo_automodel.components.speculative.streaming.queue.Lease(
deadline: float,
redelivery_count: int = 0,
lease_id: int = _next_lease_id()
)
Dataclass

Handle to a leased SampleRef.

Each SampleRefQueue.acquire call mints a fresh Lease with a unique lease_id. The queue’s ack and fail verify the lease identity before mutating internal state, so a late ACK for a stale (reclaimed) lease cannot pop a newer active lease for the same sample_id.

deadline
float

Monotonic-clock timestamp at which this lease is considered orphaned. Used by SampleRefQueue.reclaim_expired to redeliver the ref.

lease_id
int = field(default_factory=_next_lease_id)

Per-acquire unique identifier. The queue uses it as the key in _outstanding and verifies it before any ack/fail mutation.

redelivery_count
int = 0

Number of times this ref has been leased and re-leased (used for retry telemetry). Starts at 0.

ref
SampleRef

The leased reference — the only sanctioned way to materialize its tensors is FeatureStore.get, which returns a StoreHandle the consumer must hand to FeatureStore.release once it is done with them.

visibility_timeout
VisibilityTimeout

The VisibilityTimeout that produced this lease, kept here so the consumer can introspect it.

class nemo_automodel.components.speculative.streaming.queue.SampleRefQueue(
high_watermark_bytes: int | None = None,
low_watermark_bytes: int | None = None,
on_pause: typing.Callable[[StoreHealth], None] | None = None,
on_resume: typing.Callable[[StoreHealth], None] | None = None
)

Lease / ack / fail queue over SampleRef.

Thread safety: a single threading.Lock protects every list / counter, so a multi-producer / multi-consumer deployment works as long as only one thread at a time calls any one of the methods.

Parameters:

store
FeatureStore

The data-plane store the consumers will materialize against. The queue reads FeatureStore.health for backpressure.

visibility_timeout
VisibilityTimeout | NoneDefaults to None

How long a leased-but-not-acked ref can live before reclaim. Defaults to 30s; production deployments normally key this off the recipe’s per-step budget.

high_watermark_bytes
int | NoneDefaults to None

Optional resident-byte threshold for pausing. When None (default), the queue defers to StoreHealth.high_watermark_hit (i.e. the store’s own configured threshold). When set, the queue pauses whenever StoreHealth.resident_bytes >= high_watermark_bytes.

low_watermark_bytes
int | NoneDefaults to None

Optional resident-byte threshold for resuming. When None (default), the queue defers to StoreHealth.low_watermark_hit. When set, the queue resumes only after StoreHealth.resident_bytes <= low_watermark_bytes. Must be strictly less than high_watermark_bytes so the hysteresis band is non-empty.

on_pause / on_resume

Optional callbacks fired when the queue transitions high-watermark-paused -> resumed and back.

_active_by_sample
dict[str, int] = {}
_lock
= threading.Lock()
_outstanding
dict[int, Lease] = {}
_pending
list[SampleRef] = []
_pending_seen
set[str] = set()
_put_cv
= threading.Condition(self._lock)
_sample_counters
dict[str, int] = {}
_vt
= visibility_timeout or VisibilityTimeout()
is_closed
bool

Whether close has been called on this queue.

Consumers that pull acquire and receive None use this to disambiguate “drained, stop” (is_closed is True) from “transient empty poll, retry” (is_closed is False). Mirrors the Python queue.Queue separation between empty() and the lifecycle-shutdown signal.

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue._gc_store() -> None
nemo_automodel.components.speculative.streaming.queue.SampleRefQueue._should_pause(
) -> bool

Whether the producer should pause against health.

When the queue ctor was given an explicit high_watermark_bytes, that threshold wins; otherwise the decision defers to StoreHealth.high_watermark_hit (i.e. the store’s own configured threshold).

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue._should_resume(
) -> bool

Whether the producer should resume against health.

When the queue ctor was given an explicit low_watermark_bytes, that threshold wins; otherwise the decision defers to StoreHealth.low_watermark_hit. Hysteresis is preserved either way: resume crosses the low threshold, pause crosses the high threshold.

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.ack(
) -> None

Mark a leased ref as successfully consumed and free its queue slot.

Verifies Lease.lease_id matches the live outstanding entry for lease.ref.sample_id: a stale ACK for a lease that has been reclaimed and re-leased is rejected (logged, ignored) so the new consumer’s live lease is not popped by accident.

Does NOT touch the store — the consumer’s FeatureStore.get return value carries a StoreHandle that the consumer must hand to FeatureStore.release to drop the tensors. The queue’s responsibility ends at “lease no longer held”.

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.acquire(
poll_interval: float = 0.05

Lease the next ref; returns None when nothing is ready.

None is returned in two situations, which consumers disambiguate with is_closed:

  • is_closed is True: the queue has been shut down and is drained. The consumer should stop iterating.
  • is_closed is False: a transient empty poll (the producer is briefly behind). The consumer should retry.

The returned Lease is the only sanctioned way to access the ref’s tensors — FeatureStore.get requires a SampleRef, and that ref must come from a lease. The consumer MUST hand back the lease via ack (on success) or fail (on error) so the queue can reclaim the slot and the store can drop the sample.

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.close() -> None

Mark the queue closed; acquire drains what remains, then returns None.

Closing does not discard already-enqueued refs: acquire keeps handing out pending refs until they are all leased, and only returns None once the queue is closed and both pending and outstanding are empty. is_closed therefore reports the closed flag, not that the queue is already drained; FeatureDataLoader polls it to know when a None from acquire is terminal.

Outstanding leases are left intact: their consumer still owns the tensors, and a leaked FeatureStore.release would push the store’s residency counter below zero. The store’s own close is the canonical place to drop residency.

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.fail(
) -> None

Return a leased ref to the pending queue, without dropping its tensors.

Verifies the lease identity before re-enqueuing: a stale fail for a lease that has been reclaimed is a no-op. The ref will be leased again (its Lease.redelivery_count increments). Re-delivery is what makes the pipeline fault-tolerant to a transient consumer error — a permanently bad ref is the consumer’s problem (drop it after a bounded retry budget).

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.outstanding_count() -> int
nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.pending_count() -> int
nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.put(
) -> None

Enqueue ref for a future acquire.

Does not block on backpressure; producers that care should call put_blocks_until_below instead, which honors the high/low watermark hysteresis from FeatureStore.health.

nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.put_blocks_until_below(
poll_interval: float = 0.05,
abort_when: typing.Callable[[], bool] | None = None
) -> None

Enqueue ref, blocking the producer while the store is over its high watermark.

The producer is paused when _should_pause returns True (resident crossed the high threshold) and only resumed when _should_resume returns True (resident dropped back below the low threshold). In the band between the two thresholds the producer’s existing paused / unpaused state is preserved — that hysteresis is what prevents flapping when the producer is sitting near the high watermark.

Parameters:

ref
SampleRef

The reference to enqueue.

poll_interval
floatDefaults to 0.05

Seconds between backpressure checks when paused. Defaults to 50ms — well below typical step times, well above the cost of a Python-level FeatureStore.health call.

abort_when
Callable[[], bool] | NoneDefaults to None

Optional callable checked on each loop iteration. When it returns True, the put aborts with RuntimeError so a shutdown signal can unblock a producer waiting on backpressure without closing the queue first.

Raises:

  • RuntimeError: if the queue is closed while the producer is blocked, or if abort_when returns True.
nemo_automodel.components.speculative.streaming.queue.SampleRefQueue.reclaim_expired() -> int

Reclaim leases whose Lease.deadline has passed.

Each reclaimed lease is re-enqueued; acquire returns it on a future call with an incremented Lease.redelivery_count. Returns the number of leases reclaimed — a queue that is healthy returns 0 most of the time.

class nemo_automodel.components.speculative.streaming.queue.VisibilityTimeout(
seconds: float = 30.0
)
Dataclass

How long an unacked Lease is allowed to live before reclaim.

Any positive value is accepted (sub-second values are useful in tests). Production deployments typically pick something an order of magnitude larger than the recipe’s per-step budget so a slow but healthy consumer does not see its leases reclaimed out from under it.

seconds
float = 30.0
nemo_automodel.components.speculative.streaming.queue.VisibilityTimeout.__post_init__() -> None
nemo_automodel.components.speculative.streaming.queue._next_lease_id() -> int

Mint a fresh Lease.lease_id (module-level counter).

nemo_automodel.components.speculative.streaming.queue._lease_id_counter = itertools.count()
nemo_automodel.components.speculative.streaming.queue.logger = logging.getLogger(__name__)