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:
| Hop | Route | Payload |
|---|---|---|
| Producer → control plane | POST /v1/producer/submit | the Envelope, opaque inside args_blob |
| Control plane → daemon | POST /v1/daemon/poll → a run command | the same bytes, echoed as args_blob |
| Daemon → control plane | POST /v1/daemon/ack | the 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.
| Code | Name | Ext body |
|---|---|---|
1 | numpy.ndarray | msgpack([dtype_str, shape_list, raw_bytes]) |
2 | torch.Tensor | same as ndarray (detached, moved to CPU) |
3 | PIL.Image | msgpack([format_str, encoded_bytes]) |
5 | tuple | wire.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(...)).
Three properties worth knowing:
strict_types=True— atupleround-trips as a tuple (code 5), not silently coerced to a list.- Unknown type at pack time raises
TypeErrorwith the fix in the message, instead of failing deep insidemsgpack. - 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.
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.
| Field | Type | Notes |
|---|---|---|
v | int | WIRE_VERSION, currently 1. from_bytes rejects v < 1 and any v newer than the reader. |
id | string | ULID generated by the producer. This is the handle — there is no Future type. |
fn | dict | How the callable crosses the wire. See Function transport. |
args | bytes | Pre-packed, so routers never decode user values. |
kwargs | bytes | Same. |
prefix | string | Item root the worker installs into its RunContext. "" when no dls.scope. |
key | string | null | Idempotency key. See the caveat under Submit. |
cfg | dict | Only code is written today — the code-archive stamp {remote, commit, archiveKey, storageName}. |
parent | string | null | Set 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.
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:
allow_pickled— a pickledfnis arbitrary code execution, so workers refuse by default.run_workerresolvesNoneby trust domain:Trueon the local SQLite plane,Falseon any HTTP plane.lakeshore worker start --allow-pickledopts in.check_runtime()runs before any unpickle. Policystrict(exact python + cloudpickle) /minor(major.minor python, major cloudpickle — the default) /off.src_hashstaleness check on therefarm: 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 isoff.
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
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"whenqueueis absent. run_configis stored withqueuestamped onto it. Dispatch matches onrunConfig.queue; there is no queue foreign key on the invocation row.priorityis denormalized offrun_config.priorityonto the row so the claim sort actually sees it. Non-numeric / non-finite values fall back to0; fractional values are truncated.parent_invocation_idresolves the row'srootInvocationId— the parent's root, or the parent itself.
Response: {"invocation_id": "…", "queue": "…"}.
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:
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.
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
| Field | Type | Required | Notes |
|---|---|---|---|
invocation_id | string | yes | Echoed from the run command body. |
state | string | yes | The handler writes it straight to the row. Terminal states the rest of the system recognizes: succeeded, failed, killed, timeout. |
result_blob | bytes | null | no | Stored as-is; null when absent. |
error | object | null | no | {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.
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:
kind | Name | Meaning |
|---|---|---|
0 | VALUE | One yielded/returned value; body is the value. |
1 | CHUNK | A slice of one large value; body is raw bytes. |
2 | BEGIN | Opens a CHUNK run; meta may carry {total?, content_type?}. |
3 | END | Terminal — the producer finished. |
4 | ERROR | Terminal — 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:
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.
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
- Daemon protocol — the full
/v1/daemon/*surface. - Dispatch planes —
Dispatch,HttpDispatch,run_worker. - Invocation ids — reconnect-by-id and streaming from Python.
- Payload store — the out-of-band path for large blobs.