DreamLake

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.

The unprefixed name is the async one

dls.Queue is async; dls.SyncQueue is blocking. Queue is not a subclass of SyncQueue — it holds one and exposes it as q.sync.

python
import dreamlake.lakeshore as dls
from dreamlake.lakeshore import Queue, SyncQueue
python
q = SyncQueue("gpu")                          # blocking surface
inv = q.submit({"x": 1}, key="scene7/x")      # -> durable string id
val = q.result(inv, timeout=60.0)

q = Queue("gpu")                              # async surface
inv = await q.submit({"x": 1})
val = await q.result(inv)
async for frame in q.stream(inv): ...

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

python
SyncQueue(
    name: str,
    *,
    dispatch: Dispatch | None = None,   # defaults to the process plane
    kind: str = "fifo",                 # "fifo" | "priority" | "filo" | "boltzmann"
    temperature: float = 1.0,           # boltzmann sampling temperature
)

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

KindClaim order
fifo (default)priority descending, then oldest first. "priority" is an accepted alias.
filopriority descending, then newest first — autocomplete and live-preview workloads where only the latest request matters
boltzmannsample 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.

python
inv = q.submit({"seed": 42})
inv = q.submit({"seed": 42}, key="sweep/42")   # idempotent per queue
inv = q.submit({"seed": 42}, priority=9.0)
inv = q.submit(dls.thunk(lambda: 40 + 2))      # ad-hoc callable, by value

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.

Plain names here, underscored names on a UDF

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.

python
answer = q.call({"prompt": "hello"}, timeout=10.0)

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.

python
for manifest in q.stream(inv):          # SyncQueue -> Iterator
    print(manifest)

async for manifest in q.stream(inv):    # Queue -> AsyncIterator
    print(manifest)

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().

cancel and count are local-plane only today

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.

python
@dataclass
class Job:
    id: str          # the invocation id
    payload: Any     # the decoded payload / Envelope
    attempt: int
    queue: str

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.

python
with q.working() as job:
    if job is not None:
        handle(job.payload)

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.

python
t = q.topic("progress")          # q.topic() auto-names it "topic-<ULID>"

t.post({"step": 1})              # -> seq
t.post(None, final=True)         # END frame closes the topic

t.listen(since=-1, timeout=5.0)  # next value after `since`; None on timeout
for v in t.stream(since=-1, timeout=5.0):    # values until END
    ...

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

SyncQueueQueue
styleblockingawait / async for
producersubmit result call stream cancel count __len__submit result call stream cancel count
workerclaim ack nack working postclaim ack nack
topic()sync, returns Topicsync, returns the same Topic
use fromscripts, workers, testsasyncio 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:

python
ids = [await q.submit({"seed": i}) for i in range(100)]
results = await asyncio.gather(*(q.result(i, timeout=300) for i in ids))

Read next