nemo_automodel.components.speculative.streaming.queue
nemo_automodel.components.speculative.streaming.queue
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
Functions
Data
API
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.
Monotonic-clock timestamp at which this lease is
considered orphaned. Used by
SampleRefQueue.reclaim_expired to redeliver the ref.
Per-acquire unique identifier. The queue uses it as
the key in _outstanding and verifies it before any
ack/fail mutation.
Number of times this ref has been leased and re-leased (used for retry telemetry). Starts at 0.
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.
The VisibilityTimeout that produced
this lease, kept here so the consumer can introspect it.
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:
The data-plane store the consumers will materialize against.
The queue reads FeatureStore.health for backpressure.
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.
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.
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.
Optional callbacks fired when the queue transitions high-watermark-paused -> resumed and back.
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.
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).
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.
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”.
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.
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.
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).
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.
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:
The reference to enqueue.
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.
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 ifabort_whenreturnsTrue.
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.
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.
Mint a fresh Lease.lease_id (module-level counter).