# Queue API

A queue is a named channel on a [dispatch plane](/python-sdk/dispatch.md).
`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`.

> **Note:** `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](/python-sdk/udf.md) 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](/python-sdk/dispatch.md).

### 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`](/python-sdk/udf.md#dlsthunk--an-ad-hoc-callable),
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.

> **Warning:** `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()`.

> **Warning:** `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`](/python-sdk/dispatch.md#the-worker-loop) 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](/python-sdk/dispatch.md#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-"

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`

| | `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:

```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

- [Invocation ids](/python-sdk/invocations.md) — the no-Future model and
  the fan-out + collect patterns.
- [Dispatch planes](/python-sdk/dispatch.md) — `Dispatch`, `HttpDispatch`,
  and the worker loop.
- [`@udf` decorator](/python-sdk/udf.md) — functions bound to a queue.
