DreamLake

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:

PlaneBackingUse for
Dispatchstdlib sqlite3 (WAL), single filelocal runs, dev loops, tests — no server, no network
HttpDispatchthe Node control plane over httpxdeployed 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.

Dispatch is a class, not a verb

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, 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,
)
VerbControl-plane route
submitPOST /v1/producer/submit
result / getPOST /v1/producer/await (poll loop client-side)
claimPOST /v1/daemon/poll
ack / nackPOST /v1/daemon/ack
journal_appendPOST /v1/daemon/result-chunk
journal_readGET /v1/producer/invocations/:id/result/stream?since=
cancel and count are not implemented over HTTP

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

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:

FlagEnvDoes
--queueLAKESHORE_QUEUEqueue name (default "default")
--rootLAKESHORE_ROOTfilesystem root for dls.run.read/write (default cwd)
--urlLAKESHORE_URLcontrol-plane base URL → HttpDispatch; unset → the local SQLite plane
--namespaceLAKESHORE_NAMESPACEcontrol-plane namespace (default "default")
--tokenLAKESHORE_CLIENT_TOKENbearer token for --url
--once—drain at most one job, then exit
--allow-pickled—execute by-value callables
--runtime-policyLAKESHORE_RUNTIME_POLICYstrict | 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.

Retired entrypoints

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 submit hands back.