# 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](/python-sdk/udf.md) and [Queue API](/python-sdk/queue.md).

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

| Hop | Route | Payload |
| --- | --- | --- |
| Producer → control plane | `POST /v1/producer/submit` | the [Envelope](#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(...))`.

```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.

> **Note:** 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](/python-sdk/udf.md) 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.

```python file="wire.Envelope.to_bytes()"
{
    "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)
}
```

| 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](#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](#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.

```python file="fn (kind=ref)"
{
    "kind": "ref",
    "module": "user_pkg.search",
    "qualname": "topk",
    "src_hash": "sha256:5a4f0c…",   # 16 hex chars, or None if source unavailable
}
```

```python file="fn (kind=pickle)"
{
    "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

```python file="POST /v1/producer/submit"
{
    "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": "…"}`.

> **Warning:** 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:

```python file="poll response"
{
    "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](/nymph/protocol.md).

### 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](/nymph/daemons.md#shared-workdir-contract).

> **Note:** `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](#the-frame-journal) as the body yields.

## Ack

```python file="POST /v1/daemon/ack"
{
    "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
}
```

| 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.

> **Warning:** 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]
```

| `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`:

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

The row is written first (Mongo is truth), then the frame is `XADD`ed 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.

> **Warning:** `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](/nymph/protocol.md) — the full `/v1/daemon/*` surface.
- [Dispatch planes](/python-sdk/dispatch.md) — `Dispatch`, `HttpDispatch`, `run_worker`.
- [Invocation ids](/python-sdk/invocations.md) — reconnect-by-id and streaming from Python.
- [Payload store](/get-started/payloads.md) — the out-of-band path for large blobs.
