DreamLake

Function protocol

This is the wire-level reference for a remote @udf invocation. It covers the bytes, not the Python API — for the producer surface read @udf decorator and Queue API.

Every body on every hop is msgpack. Three hops carry one invocation:

HopRoutePayload
Producer → control planePOST /v1/producer/submitthe Envelope, opaque inside args_blob
Control plane → daemonPOST /v1/daemon/poll → a run commandthe same bytes, echoed as args_blob
Daemon → control planePOST /v1/daemon/ackthe wire-packed return value

The control plane never decodes user values. It routes on the outer JSON-ish fields (queue, priority, state) and treats the Envelope as bytes.

Value codec

dreamlake.lakeshore.wire is the codec. It is msgpack with integer ExtType codes, not {"ztype": ...} marker dicts — a tagged leaf can never collide with a user dict, and msgpack does the container recursion.

CodeNameExt body
1numpy.ndarraymsgpack([dtype_str, shape_list, raw_bytes])
2torch.Tensorsame as ndarray (detached, moved to CPU)
3PIL.Imagemsgpack([format_str, encoded_bytes])
5tuplewire.pack(list(t)) — recursive

Everything else must be msgpack-native (str, bytes, int, float, bool, None, list, dict). Register more with wire.register(wire.Codec(...)).

python
from dreamlake.lakeshore import wire

blob = wire.pack({"x": np.zeros(4), "pair": (1, 2)})
wire.unpack(blob)          # {"x": array([0., 0., 0., 0.]), "pair": (1, 2)}

Three properties worth knowing:

  • strict_types=True — a tuple round-trips as a tuple (code 5), not silently coerced to a list.
  • Unknown type at pack time raises TypeError with the fix in the message, instead of failing deep inside msgpack.
  • Unknown ext code at unpack time is preserved verbatim as msgpack.ExtType(code, data), so a newer peer's values pass through an older reader unharmed.

numpy, torch, and PIL imports are all deferred inside the codecs — importing wire drags in none of them.

Files never ride this wire

There is deliberately no file/blob codec. A data UDF returns string keys; the bytes move worker ↔ storage on a side channel. See @udf decorator for the return-your-keys contract.

The Envelope

One remote invocation. Built by UDF.submit() (or Thunk.to_envelope()), serialized with Envelope.to_bytes(), and handed to the dispatch plane as opaque bytes.

wire.Envelope.to_bytes()python
{
    "v": 1,                     # WIRE_VERSION
    "id": "01HX2K3M4N5P…",      # ULID — the durable handle
    "fn": {...},                # tagged union, see below
    "args": b"…",               # wire.pack(list(args)) — pre-packed
    "kwargs": b"…",             # wire.pack(dict(kwargs)) — pre-packed
    "prefix": "scenes/0007",    # dls.scope prefix for dls.run key resolution
    "key": None,                # optional idempotency key
    "cfg": {},                  # {"code": {...}} when a code archive shipped
    "parent": None,             # parent invocation id (pipeline lineage)
}
FieldTypeNotes
vintWIRE_VERSION, currently 1. from_bytes rejects v < 1 and any v newer than the reader.
idstringULID generated by the producer. This is the handle — there is no Future type.
fndictHow the callable crosses the wire. See Function transport.
argsbytesPre-packed, so routers never decode user values.
kwargsbytesSame.
prefixstringItem root the worker installs into its RunContext. "" when no dls.scope.
keystring | nullIdempotency key. See the caveat under Submit.
cfgdictOnly code is written today — the code-archive stamp {remote, commit, archiveKey, storageName}.
parentstring | nullSet on chained submits.

The header (id / fn / prefix / cfg) decodes without touching args / kwargs. Local calls never build an Envelope at all — this is strictly the remote hop.

Function transport

Envelope.fn is a tagged union with two arms.

fn (kind=ref)python
{
    "kind": "ref",
    "module": "user_pkg.search",
    "qualname": "topk",
    "src_hash": "sha256:5a4f0c…",   # 16 hex chars, or None if source unavailable
}
fn (kind=pickle)python
{
    "kind": "pickle",
    "blob": b"…",                   # cloudpickle, ≤ 32 KiB inline
    "hash": "sha256:…",             # full hex digest of the blob
    "source": "def f(): …",         # best-effort, for audit/display — never executed
    "qualname": "<lambda>",
    "where": "<lambda> (/path/x.py:12)",
    "runtime": {"python": "3.12.4", "impl": "cpython",
                "cloudpickle": "3.0.0", "pickle_protocol": 5,
                "platform": "linux-x86_64"},
}

transport="auto" (the @udf default) picks pickle exactly when the callable is not importable by reference — a lambda, a <locals> closure, a __main__ def, or a callable instance with no __qualname__ of its own. Everything else ships by ref. transport="ref" insists on importability and raises TransportError otherwise; transport="pickle" always snapshots.

Only the callable is ever pickled. Args, kwargs, results, and journal frames stay on the msgpack codec.

Worker-side gates, in order:

  1. allow_pickled — a pickled fn is arbitrary code execution, so workers refuse by default. run_worker resolves None by trust domain: True on the local SQLite plane, False on any HTTP plane. lakeshore worker start --allow-pickled opts in.
  2. check_runtime() runs before any unpickle. Policy strict (exact python + cloudpickle) / minor (major.minor python, major cloudpickle — the default) / off.
  3. src_hash staleness check on the ref arm: if the worker's resolved source hashes differently from the submitter's, the invocation fails loud rather than running stale code. Skipped when the envelope carries no hash, the target's source is unavailable, or the policy is off.

INLINE_PICKLE_LIMIT is 32 KiB. A larger blob is a captured-state mistake, and the error says so — pass the data as arguments instead.

Submit

POST /v1/producer/submitpython
{
    "invocation_id": "01HX…",
    "queue": "gpu",
    "function": {"module": "_native_wire",
                 "qualname": "_native_wire.envelope",
                 "version": "v1"},
    "args_blob": b"<the Envelope>",
    "run_config": {"queue": "gpu", "priority": 0.0, "key": None},
    "code_context": {},
    "parent_invocation_id": None,       # omitted when absent
}

function is a stub on the native path — the real callable identity lives inside the Envelope, and the control plane only needs a Function row to hang the invocation off. args_blob is the Envelope, verbatim.

Server side:

  • The named queue is ensure-created (idempotent), defaulting to "default" when queue is absent.
  • run_config is stored with queue stamped onto it. Dispatch matches on runConfig.queue; there is no queue foreign key on the invocation row.
  • priority is denormalized off run_config.priority onto the row so the claim sort actually sees it. Non-numeric / non-finite values fall back to 0; fractional values are truncated.
  • parent_invocation_id resolves the row's rootInvocationId — the parent's root, or the parent itself.

Response: {"invocation_id": "…", "queue": "…"}.

`key` is carried, not enforced

The idempotency key rides Envelope.key and run_config.key, but no control-plane route writes Invocation.idempotencyKey, and the sparse unique index it would need has to be created by hand in MongoDB. Two submits with the same key produce two invocations today.

Dispatch

A daemon long-polls POST /v1/daemon/poll. When the claim succeeds it gets back one run command:

poll responsepython
{
    "commands": [{
        "id": "01HX…",                 # command id (ULID)
        "kind": "run",
        "body": {
            "invocation_id": "01HX…",
            "function": {"module": …, "qualname": …, "version": …},
            "args_blob": b"<the Envelope>",
            "run_config": {...},
            "prefix": "scenes/0007",   # only when run_config.prefix is set
        },
    }],
    "target_nymph_version": "0.1.4",   # only when the namespace pins one
}

The claim is a single MongoDB findAndModify on the Invocation collection: query: {state: "queued", "runConfig.queue": …}, sort: {priority: -1, submittedAt: 1}, flipping the row to running.

nymph's RunBody additionally accepts parent_invocation_id, root_invocation_id, script (run bash -c instead of the Python entry), and as_user (wrap the spawn in runuser -u). See Daemon protocol.

The workdir contract

nymph's process / docker / gvisor / slurm / kube runners do not speak the Envelope. They write args_blob verbatim to {workdir}/{invocation_id}/args.msgpack, spawn python3 -m dreamlake.lakeshore.daemon.entry (or python -m … inside a container) with that directory as cwd, and read back:

  • exit 0 → result.msgpack, the wire-packed return value;
  • exit non-zero → error.json, {"type", "message", "traceback"}, reported as "<type>: <message>".

The workdir is deleted on success and left on disk on failure. Full contract: Daemons → Shared workdir contract.

No journal on the nymph path

entry.py drains generators and async-generators into the concatenated manifest and returns it as one value. Per-yield streaming exists only on the queue-native worker (lakeshore worker start), which appends frames as the body yields.

Ack

POST /v1/daemon/ackpython
{
    "invocation_id": "01HX…",
    "state": "succeeded",        # or "failed" | "killed" | "timeout"
    "result_blob": b"…",         # wire-packed return value; empty otherwise
    "error": {                   # ErrorInfo, on failure
        "type": "ValueError",
        "message": "…",
        "traceback": "…",        # optional
    },
    "worker_id": "…",            # nymph sends this; omitted when empty
}
FieldTypeRequiredNotes
invocation_idstringyesEchoed from the run command body.
statestringyesThe handler writes it straight to the row. Terminal states the rest of the system recognizes: succeeded, failed, killed, timeout.
result_blobbytes | nullnoStored as-is; null when absent.
errorobject | nullno{type, message, traceback?}. On the Rust side the field is type on the wire.

The handler sets finishedAt = now() unconditionally. startedAt is set server-side by the claim, not reported by the daemon.

Outcomes a nymph runner can produce: succeeded, failed (a bare-string failure gets type: "DaemonError"), and killed (cancellation, with {type: "Cancelled", message: "daemon shutdown grace expired"}).

POST /v1/producer/await is the producer-side counterpart: a long poll (default timeout_s 60, 100 ms spin) returning {state, result_blob, error} once terminal, or {"state": "pending"} on timeout.

No metrics or log pointers on the ack

The ack envelope carries no started_at / finished_at, no stdout_uri / stderr_uri, and no resource metrics. Progress and log lines flow on the separate POST /v1/daemon/event channel, which the control plane stores as invocation.progress / invocation.log / invocation.event. nymph itself does not post to /v1/daemon/event today.

The frame journal

A streaming invocation writes frames. One frame is a 4-element msgpack array:

python
[kind, seq, body, meta]
kindNameMeaning
0VALUEOne yielded/returned value; body is the value.
1CHUNKA slice of one large value; body is raw bytes.
2BEGINOpens a CHUNK run; meta may carry {total?, content_type?}.
3ENDTerminal — the producer finished.
4ERRORTerminal — the producer raised; meta is {type, message}.

seq is a monotone per-invocation cursor starting at 0 — the reconnect handle. An explicit END/ERROR frame closes the stream, so "producer finished" stays distinguishable from "connection dropped".

Frames are msgpack-self-delimiting: N frames concatenated, or split at arbitrary byte boundaries by HTTP chunking or Redis, reassemble through one msgpack.Unpacker. That is what wire.FrameReader does.

Append — POST /v1/daemon/result-chunk:

python
{"invocation_id": "01HX…", "seq": 3, "bytes": b"<frame>", "final": False}

The row is written first (Mongo is truth), then the frame is XADDed to the Redis stream for live readers. A retried POST of an already-journaled (invocation_id, seq) returns {"ok": true, "seq": 3, "duplicate": true} and deliberately does not re-XADD. A Redis failure is logged and swallowed — live readers catch up from Mongo.

Read — GET /v1/producer/invocations/:id/result/stream?since=<seq>:

A chunked application/x-msgpack-stream body. Frames are written verbatim, with no SSE wrapper and no delimiter. The route replays Mongo rows with seq > since, then follows the Redis stream, ending on a final frame or once the invocation reaches a terminal state. The seq you resume from comes from inside the decoded frame, not from the transport.

Versioning

Envelope.v is the single version integer. Envelope.from_bytes rejects a version newer than the reading client with a message naming the fix (upgrade dreamlake-lakeshore). Unknown ext codes and unknown cfg keys pass through untouched, so additive changes do not need a version bump.

Not implemented over HTTP

HttpDispatch.cancel() and HttpDispatch.count() raise NotImplementedError. Killing an invocation goes through DELETE /v1/namespaces/:ns/invocations/:id, which only works on a queued row — a running invocation returns 409, because the wire protocol has no daemon-side cancel signal yet.

See also