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 |
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.
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.
The design rules (each one is a named regression test):
- 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. - Every transition is one guarded atomic UPDATE. A lost race is a clean miss, never a double-delivery.
- Fencing — ack/nack require the claiming worker's id; a stale worker whose job was re-dispatched cannot clobber the new attempt.
- Idempotent submit — a caller-supplied
keyis unique per (namespace, queue); resubmitting returns the existing invocation. - Counts come from the same rows that get claimed — they cannot disagree with the queue contents.
- 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
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
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, 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:
| 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= |
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.
The worker loop
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_pickledgates by-value (cloudpickle) callables — executing one is arbitrary code execution.None(the default) resolves by trust domain:Trueon the localDispatchplane,Falseon any HTTP plane.runtime_policygoverns the pre-unpickle interpreter-stamp check and the reference-transportsrc_hashstaleness check."strict"demands exact Python and cloudpickle versions,"minor"(default) matchesmajor.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:
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.
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 — the surface you normally use; these planes sit behind it.
- Invocation ids — what
submithands back.