# Dispatch planes

A dispatch plane is where invocation rows and frame journals live.
`Queue` / `SyncQueue` / `run_worker` take any object with the dispatch
verbs; two implementations ship:

| Plane | Backing | Use for |
| --- | --- | --- |
| `Dispatch` | stdlib sqlite3 (WAL), single file | local runs, dev loops, tests — no server, no network |
| `HttpDispatch` | the Node control plane over httpx | deployed queues, remote workers |

```python
from dreamlake.lakeshore import Dispatch, HttpDispatch, SyncQueue

q = SyncQueue("gpu")                                   # process-default plane
q = SyncQueue("gpu", dispatch=Dispatch(":memory:"))    # ephemeral, tests
q = SyncQueue("gpu", dispatch=HttpDispatch(
    "https://cp.example.com", token=TOKEN,
))
```

Omitting `dispatch` uses the process-default plane, resolved once from
the environment: `HttpDispatch` when `LAKESHORE_URL` is set, otherwise
the local SQLite file at `$LAKESHORE_HOME/dispatch.sqlite` (default
`~/.lakeshore`). Two processes on one machine pointing at the same file
share a working queue: submit in one shell, run a worker in another,
zero infrastructure.

> **Note:** `dls.Dispatch(...)` constructs a plane. There is no module-level
> `dispatch(...)` function.

## `Dispatch` — the minimal plane

The executable specification of the dispatch contract. The Node control
plane implements the same semantics at deployment scale; the conformance
suite runs against both, so the contract cannot fork.

```python
Dispatch(path: str = ":memory:", *, ns: str = "default")
```

The design rules (each one is a named regression test):

1. **One row per invocation** — queue entry, metadata, and state machine
   in one place. No pending list, no secondary index, nothing to drift.
   States are `queued`, `running`, `succeeded`, `failed`, `cancelled`.
2. **Every transition is one guarded atomic UPDATE.** A lost race is a
   clean miss, never a double-delivery.
3. **Fencing** — ack/nack require the claiming worker's id; a stale
   worker whose job was re-dispatched cannot clobber the new attempt.
4. **Idempotent submit** — a caller-supplied `key` is unique per
   (namespace, queue); resubmitting returns the existing invocation.
5. **Counts come from the same rows that get claimed** — they cannot
   disagree with the queue contents.
6. **The journal is derived, not truth** — frames feed streaming
   consumers; the terminal result lives on the invocation row.
   Trimming degrades reads to the row, never to wrong state.

### Verbs

```python
dp.submit(queue, envelope, *, id, key=None, priority=0.0, parent=None) -> str
dp.claim(queue, worker_id, *, kind="fifo", temperature=1.0) -> Claimed | None
dp.ack(inv_id, worker_id, *, result=None)
dp.nack(inv_id, worker_id, *, error, retry=True)
dp.cancel(inv_id) -> bool
dp.get(inv_id) -> dict | None          # the invocation row
dp.result(inv_id, *, timeout=30.0, poll=0.02) -> bytes
dp.count(queue, state="queued") -> int
dp.close()
```

`Claimed` is a dataclass of `id`, `envelope`, `attempt`, `queue`.
Claim ordering by `kind`: `fifo` (priority desc, oldest first; `priority`
is an alias), `filo` (priority desc, newest first), `boltzmann` (sample
proportional to `exp(priority / T)`) — all queries over the same rows
that hold the state.

### The frame journal

```python
dp.journal_append(inv_id, seq, frame_bytes, *, final=False) -> int
dp.journal_read(inv_id, *, since=-1, timeout=30.0, poll=0.02) -> Iterator[(seq, bytes)]
dp.journal_trim(inv_id, *, keep_last=0) -> int
```

An explicit `seq` makes appends idempotent on `(inv_id, seq)` — a
retried frame is absorbed, not duplicated; `seq=None` assigns the next
one. `journal_read` blocks for new frames until a final frame or a
terminal invocation state, and activity resets the timeout; the `since`
cursor makes reconnect lossless. You normally consume this through
`q.stream` and `Topic` on the [Queue API](/python-sdk/queue.md), which
decode frames for you.

Frames are `[kind, seq, body, meta]` on the wire. `FrameKind` is
`VALUE=0`, `CHUNK=1`, `BEGIN=2`, `END=3`, `ERROR=4`; `END` and `ERROR`
are terminal.

## `HttpDispatch` — the Node control plane

The same verbs, spoken over HTTP. The native `Envelope` rides
**opaquely** inside the existing submit body, so the control plane never
decodes user payloads:

```python
HttpDispatch(
    base_url: str,
    *,
    token: str | None = None,
    namespace: str = "default",
    timeout: float = 30.0,
)
```

| Verb | Control-plane route |
| --- | --- |
| `submit` | `POST /v1/producer/submit` |
| `result` / `get` | `POST /v1/producer/await` (poll loop client-side) |
| `claim` | `POST /v1/daemon/poll` |
| `ack` / `nack` | `POST /v1/daemon/ack` |
| `journal_append` | `POST /v1/daemon/result-chunk` |
| `journal_read` | `GET /v1/producer/invocations/:id/result/stream?since=` |

> **Warning:** `HttpDispatch.cancel` and `HttpDispatch.count` raise
> `NotImplementedError` — those control-plane routes have not landed. Both
> work on the local `Dispatch` plane.

`HttpDispatch` also carries `push_code()`, which archives and uploads the
current git tree (deduped per commit) and returns the code-mount stamp.
`UDF.submit` calls it for reference-transport functions so a fresh remote
worker can import your modules — see
[code mount](/python-sdk/happy-paths/10-code-mount.md).

## The worker loop

```python
from dreamlake.lakeshore import SyncQueue, run_worker

q = SyncQueue("gpu")
run_worker(q, root="/data")        # loop until stopped
run_worker(q, once=True)           # drain one job (tests, cron)
```

```python
run_worker(
    q: SyncQueue,
    *,
    root: Path | str | None = None,     # filesystem root for dls.run (default cwd)
    registry: dict[str, Callable] | None = None,
    once: bool = False,
    poll: float = 0.05,
    retries: int = 1,
    stop: Callable[[], bool] | None = None,
    allow_pickled: bool | None = None,
    runtime_policy: str = "minor",      # "strict" | "minor" | "off"
) -> int                                # jobs executed
```

For each claimed job the worker decodes the `Envelope`, resolves the
function (a `registry` hit first, else the `module:qualname` import or
the pickled blob), installs a `RunContext` from the envelope's prefix and
the worker's `root`, runs the body by its kind, and acks the packed
result — or appends an ERROR frame and nacks on failure, retrying while
`attempt <= retries`. Generator yields stream live: one VALUE frame per
yield followed by an END frame, so consumers iterate `q.stream(inv)`
while the job runs. The return-your-keys contract is enforced by the
`UDF` wrapper itself, so a body fails remotely with the same message it
fails with locally.

Two safety knobs:

- **`allow_pickled`** gates by-value (cloudpickle) callables — executing
  one is arbitrary code execution. `None` (the default) resolves by trust
  domain: `True` on the local `Dispatch` plane, `False` on any HTTP
  plane.
- **`runtime_policy`** governs the pre-unpickle interpreter-stamp check
  and the reference-transport `src_hash` staleness check. `"strict"`
  demands exact Python and cloudpickle versions, `"minor"` (default)
  matches `major.minor`, `"off"` skips both — same-machine dev only.

`stop` is polled between jobs for graceful shutdown. `run_worker` returns
the number of jobs it executed.

## The worker entrypoint

The module path `python -m dreamlake.lakeshore.daemon` is the stable
spelling. The TypeScript CLI's `lakeshore worker` (npm
`@dreamlake/lakeshore`) spawns and supervises exactly that module:

```bash
lakeshore worker start --queue gpu --root /data
lakeshore worker start --queue gpu --url https://cp.example.com --token "$TOK"
lakeshore worker once  --queue gpu      # drain at most one job, then exit

python -m dreamlake.lakeshore.daemon --queue gpu --root /data
python -m dreamlake.lakeshore.daemon --queue gpu --once
```

Flags win; env fills the gaps:

| Flag | Env | Does |
| --- | --- | --- |
| `--queue` | `LAKESHORE_QUEUE` | queue name (default `"default"`) |
| `--root` | `LAKESHORE_ROOT` | filesystem root for `dls.run.read/write` (default cwd) |
| `--url` | `LAKESHORE_URL` | control-plane base URL → `HttpDispatch`; unset → the local SQLite plane |
| `--namespace` | `LAKESHORE_NAMESPACE` | control-plane namespace (default `"default"`) |
| `--token` | `LAKESHORE_CLIENT_TOKEN` | bearer token for `--url` |
| `--once` | — | drain at most one job, then exit |
| `--allow-pickled` | — | execute by-value callables |
| `--runtime-policy` | `LAKESHORE_RUNTIME_POLICY` | `strict` \| `minor` \| `off` (default `minor`) |

With no `--url` the entrypoint uses the local SQLite plane and forces
`allow_pickled=True` — the same machine and user is one trust domain.
With `--url` it builds an `HttpDispatch` and accepts pickled callables
only when `--allow-pickled` is passed.

> **Warning:** The `lakeshore-daemon` console script no longer exists — the SDK
> installs no console scripts at all.
> `dreamlake.lakeshore.daemon.entry` still exists, but only as a
> compatibility shim for nymph's hard-coded per-job exec contract; it is
> not a way to run a worker.

## Read next

- [Queue API](/python-sdk/queue.md) — the surface you normally use; these
  planes sit behind it.
- [Invocation ids](/python-sdk/invocations.md) — what `submit` hands back.
