Queue API
A queue is a named channel on a dispatch plane.
SyncQueue is the blocking class and holds all the verbs; Queue is the
async sibling that wraps one and offloads its calls with
asyncio.to_thread.
dls.Queue is async; dls.SyncQueue is blocking. Queue is not
a subclass of SyncQueue — it holds one and exposes it as q.sync.
Payload values ride the wire codec (msgpack + ExtType), so ndarray / tensor / image / tuple round-trip. File bytes never do — data UDFs return string keys and files stay on the filesystem.
Constructor
Queue(name, **kw) forwards the same arguments. dispatch selects the
plane: omit it for the process default (the local SQLite plane under
~/.lakeshore, or HttpDispatch when LAKESHORE_URL is set), pass
Dispatch(":memory:") for tests, or an explicit HttpDispatch(url) —
see Dispatch planes.
Kinds
| Kind | Claim order |
|---|---|
fifo (default) | priority descending, then oldest first. "priority" is an accepted alias. |
filo | priority descending, then newest first — autocomplete and live-preview workloads where only the latest request matters |
boltzmann | sample proportional to exp(priority / T); T -> 0 collapses to strict priority, large T to uniform random |
Anything else raises
DispatchError: unknown queue kind '...'; use fifo|filo|boltzmann.
Ordering is what the worker asks for at claim time; submit only
attaches priority.
Producer verbs
submit(payload=None, *, envelope=None, key=None, priority=0.0) -> str
Enqueue any wire-native value, a Thunk,
or a pre-built envelope. Returns the durable invocation id — a plain
string, the whole handle.
key makes the submit idempotent: the same key on the same queue
returns the existing invocation id, so resubmit-after-crash is safe by
construction.
SyncQueue.submit takes key= / priority= / envelope=.
UDF.submit takes _key= / _priority= / _dispatch= — underscored,
because its other keyword arguments belong to your function.
result(inv_id, *, timeout=30.0) -> Any
Block (or await) until the invocation is terminal, then return the
decoded result. Re-attachable from any process that holds the id.
Raises DispatchError if the invocation failed or was cancelled, and
TimeoutError when the window elapses.
call(payload, *, timeout=30.0, **kw) -> Any
RPC in one line: submit then result, keyed by the invocation id — no
side-channel reply topic.
stream(inv_id, *, since=-1, timeout=30.0)
Iterate an invocation's journal frames as decoded values, live while the
job runs. since is the reconnect cursor — pass the last seq you saw to
resume without loss or duplication. An ERROR frame raises
DispatchError; the END frame stops iteration cleanly; timeout
seconds of silence raises TimeoutError.
cancel(inv_id) -> bool · count(state="queued") -> int · len(q)
cancel moves a queued invocation to cancelled and returns whether
it did; a running job finishes its attempt. count counts rows in one
state — queued, running, succeeded, failed, or cancelled — from
the same rows that get claimed, so counts cannot disagree with the
queue's contents. len(q) is count().
HttpDispatch.cancel and HttpDispatch.count raise
NotImplementedError pending the control-plane routes. They work on the
local Dispatch plane.
Worker verbs
The consuming side of the same queue. Most workers just call
run_worker and never touch
these; they exist for custom loops.
claim() -> Job | None
Claim one job, or None when the queue is empty. Atomic — exactly one
worker wins a given row.
ack(job, result=None) / nack(job, *, error, retry=True)
Terminal success / failure. Both are fenced: only the claiming
worker may settle the job, so a stale worker whose job was re-dispatched
cannot clobber the new attempt. Settling a job you do not hold raises
DispatchError.
working() — the common loop, packaged
Claims one job, auto-acks on success, auto-nacks with retry=True on an
exception (and re-raises). Yields None when the queue is empty.
post(inv_id, value, *, final=False) -> int
Worker-side streaming: append one VALUE frame (or, with final=True, the
terminal END frame) to an invocation's journal. Returns the assigned seq.
The worker loop does this for you
on every generator yield.
Topic — pub/sub over the frame journal
Topics are journal channels: the same mechanism as invocation streaming, cursor reconnect included.
Note the shorter default: Topic.listen and Topic.stream time out
after 5 seconds of silence, where q.stream waits 30.
The journal is derived data — truth is the invocation row. Trimming a journal degrades streaming reads to the terminal result, never to wrong state.
Queue vs SyncQueue
SyncQueue | Queue | |
|---|---|---|
| style | blocking | await / async for |
| producer | submit result call stream cancel count __len__ | submit result call stream cancel count |
| worker | claim ack nack working post | claim ack nack |
topic() | sync, returns Topic | sync, returns the same Topic |
| use from | scripts, workers, tests | asyncio apps, notebooks with a loop |
Queue offloads blocking calls so slow polls never stall the event
loop; q.sync exposes the blocking sibling sharing the same binding —
that is where working(), post(), and len() live.
Fan-out composes with stdlib asyncio; there is no bespoke gather:
Read next
- Invocation ids — the no-Future model and the fan-out + collect patterns.
- Dispatch planes —
Dispatch,HttpDispatch, and the worker loop. @udfdecorator — functions bound to a queue.