# `@udf` decorator

```python
import dreamlake.lakeshore as dls
from dreamlake.lakeshore import udf
```

One decorator. The function's own shape — sync, async, generator, or
async generator — picks the execution shape; there is no mode flag to
forget and no second decorator.

```python
udf(fn=None, *, queue: str | None = None, transport: str = "auto")
```

Those are the only parameters. There is no `mode=`, `name=`, `retries=`,
`timeout=`, `image=`, or `resources=` on `@udf`.

| Form | Plain call does |
| --- | --- |
| `@udf` (bare) | runs the body in-process and returns the result |
| `@udf(queue="gpu")` | submits and waits — equivalent to `.remote(...)` |

## Bare `@udf` — local by default

A plain call runs in-process, exactly like the undecorated function.
No control plane, no daemon, no network — the local tier imports
nothing distributed.

```python
@dls.udf
def stats(a: list[float]) -> dict:
    return {"mean": sum(a) / len(a)}

stats([1.0, 2.0, 3.0])   # {'mean': 2.0}
```

Every plain call runs inside a **fork** of the ambient run context — its
own key ledger — so sequential calls in one `dls.scope` never
cross-contaminate the return-your-keys check.

## The four body kinds

Detected once at decoration with `inspect`; each dispatches naturally:

```python
@dls.udf
def f(x): ...                 # sync       -> the value

@dls.udf
async def g(x): ...           # async      -> coroutine, await it

@dls.udf
def h(x):                     # generator  -> iterator
    yield ...

@dls.udf
async def k(x):               # async-gen  -> async iterator
    yield ...
```

The `UDF` object exposes the detection as `.is_async`, `.is_gen`, and
`.is_asyncgen`, alongside `.queue` and `.transport`.

For a generator body, the run context is installed only around each
`next()` — never while control is back with the consumer — so a yield
cannot leak the call's scope into consumer code.

## Data UDFs — string keys and the `dls.run` seam

A data UDF never sees a bucket, a credential, or an absolute path. It
asks the ambient run context for local paths and returns plain string
**keys**:

```python
@dls.udf
async def splats_to_mesh(splats: str) -> str:
    mesh = extract(dls.run.read(splats))            # key -> local Path
    mesh.save(dls.run.write("process/mesh.ply"))    # key -> local Path
    return "process/mesh.ply"                       # a KEY, not a path
```

Keys are prefix-relative POSIX strings (`"process/mesh.ply"`). An
absolute key is rejected; `../` may climb to a parent scope but never
out of the run root. Keys resolve to `root / prefix / key` on the local
filesystem, and file bytes never ride the msgpack wire — only keys do.

### `dls.run` surface

| Member | Signature | Does |
| --- | --- | --- |
| `run.write` | `(key, *, dir=False) -> Path` | Register an output key; return the local path to write at. Parents created; `dir=True` makes the directory itself. |
| `run.read` | `(key) -> Path` | Resolve an input key to a local path. Raises `KeyError_` if it does not exist. |
| `run.read_all` / `run.write_all` | `(keys) -> same shape` | Shape-mirroring: `str -> Path`, `list -> list`, `dict -> dict`. |
| `run.prefix` / `run.root` | properties | The ambient prefix (`str`) and filesystem root (`Path`). |

> **Warning:** `RunContext` has a `storage` field and a fetch-on-read path, but nothing
> binds a `Storage` to it today — `run_worker` installs a context with
> `storage=None`. In practice a worker's `dls.run.read/write` resolve
> against the worker's own filesystem under its `--root`, so a remote data
> UDF needs a shared or mounted filesystem. The automatic
> upload-recorded-writes side channel is designed but not implemented.

### `dls.scope` — set the prefix

```python
dls.scope(sub: str = "", *, root: Path | str | None = None)
```

```python
with dls.scope("scenes/0007"):                     # prefix = scenes/0007
    key = splats_to_mesh("source/splats.ply")

with dls.scope(scene_root, root="/data"):          # nesting joins prefixes
    for i in range(num_episodes):
        with dls.scope(f"outputs/{i:04d}"):
            ...
```

`root=` re-anchors the filesystem root — use it on the outermost scope
of a local run. Each scope tracks its own writes, so the key contract
stays per-invocation.

## The return-your-keys contract

The contract is **one-directional**: every key written through
`dls.run.write` must appear in the returned manifest. A returned string
that was never written is fine — it may be plain data, or an upstream
key passed through.

Accepted manifest shapes: `str`, `list[str]`, `tuple[str, ...]`, or
`dict[str, str]` (named outputs; the values are the keys). Violations
raise `KeyError_`, a subclass of `ValueError`.

| You did | You get |
| --- | --- |
| wrote via `dls.run.write`, returned a non-manifest value | `KeyError_: UDF wrote N file(s) via dls.run.write (...) but did not return their keys` |
| wrote a key you left out of the manifest | `KeyError_: UDF wrote key(s) it did not return: [...]` |
| returned a key you never wrote | allowed — plain data or a pass-through key |
| returned plain data, wrote nothing | allowed — not every UDF is a data UDF |

A typo'd return is still caught: the correctly written key goes
unreturned and fails.

Multi-output shapes are stdlib primitives:

```python
@dls.udf
def poses_to_views(splats: str, poses: str) -> dict[str, str]:
    rgb = dls.run.write("process/images/rgb", dir=True)
    depth = dls.run.write("process/images/depth", dir=True)
    render(splats, poses, rgb=rgb, depth=depth)
    return {"rgb": "process/images/rgb", "depth": "process/images/depth"}
```

## Streaming bodies

A generator or async generator yields **partial manifests** — each yield
is shipped as one journal frame, and the yields concatenate into THE
manifest for the key contract:

```python
@dls.udf(queue="gpu")
async def traj_to_ego_views(traj: str):
    for i, frame in enumerate(render(dls.run.read(traj))):
        frame.save(dls.run.write(f"images/{i:04d}.png"))
        yield [f"images/{i:04d}.png"]        # one VALUE frame, live
```

If the yields are not key manifests (a plain data stream), the contract
is skipped for that invocation. Consumers iterate the journal while the
job runs — see [streaming](/python-sdk/invocations.md#streaming).

## `@udf(queue=...)` — remote by default

```python
@dls.udf(queue="gpu")
def embed(x: str) -> list[float]: ...

embed("hello")               # plain call = submit + wait for the result
inv = embed.submit("hello")  # non-blocking -> durable invocation id (str)
```

A plain call on a queue-bound UDF submits and waits. An `async def` or
async-generator body returns an awaitable instead of blocking.

> **Warning:** `.remote()` (and therefore a plain call on a queue-bound UDF) fetches
> the result with the dispatch plane's default `timeout=30.0`. For jobs
> that run longer, use `inv = f.submit(...)` and then
> `q.result(inv, timeout=...)`.

### `.submit(*args, **kwargs) -> str`

Returns the durable invocation id immediately. Control kwargs are
**underscore-prefixed** so they cannot collide with your function's own
keyword arguments:

| Kwarg | Does |
| --- | --- |
| `_key` | Idempotency key — resubmitting the same key on the same queue returns the existing invocation id, no duplicate job. |
| `_priority` | Float attached to the row. `fifo` and `filo` claim highest priority first; `boltzmann` samples proportional to `exp(priority / T)`. Default `0.0`. |
| `_dispatch` | Override the dispatch plane (default: the process plane — see [Dispatch planes](/python-sdk/dispatch.md)). |

`submit` on a UDF with no queue binding raises:

```text
TypeError: embed has no queue binding; decorate with @udf(queue=...) or
pass one at submit time
```

Arguments ride `wire.pack` inside an `Envelope` — ndarray, tensor,
image, and tuple all survive the round trip. The ambient
`dls.run.prefix` travels in the envelope, so the worker resolves keys
exactly as you would have locally.

### `.remote(*args, _dispatch=None, **kwargs)`

Submit and wait. Sync and generator bodies block and return the value;
async and async-generator bodies return an awaitable.

### `.local(*args, **kwargs)`

Force the local tier even on a queue-bound UDF — debugging and unit
tests stay plain function calls:

```python
@dls.udf(queue="gpu")
def add(a: int, b: int) -> int:
    return a + b

add.local(2, 3)   # 5 — no queue, no network
```

`.local()` builds a fresh unbound `UDF` around the same function, which
also drops the original `transport=` setting. That is irrelevant for
local execution but worth knowing.

## How the function crosses the wire

`transport=` on the decorator picks the encoding of the callable itself.
Arguments, keyword arguments, results, and stream frames always use the
msgpack/ExtType codec — only the callable is ever pickled.

| `transport` | Ships | Fails when |
| --- | --- | --- |
| `"auto"` (default) | `ref` when the callable is importable, else `pickle` | — |
| `"ref"` | `module:qualname` plus a `src_hash` version stamp | the callable is a lambda, a `<locals>` closure, a `__main__` def, or a callable instance |
| `"pickle"` | a cloudpickle snapshot plus a runtime stamp | the blob exceeds 32 KiB |

Anything else raises
`TransportError: unknown transport '...'; use "ref" | "pickle" | "auto"`.

Two guardrails on the `pickle` path:

- **Size.** `INLINE_PICKLE_LIMIT` is 32 KiB. A bigger blob means the
  callable captured large state; pass the data as arguments instead.
- **Trust.** Workers refuse pickled callables unless started with
  `allow_pickled=True`. `run_worker(allow_pickled=None)` — the default —
  resolves by trust domain: allowed on the local SQLite plane, refused
  on an HTTP plane.

On the `ref` path the worker compares the submitter's `src_hash` with
the source it resolved and fails loud on a mismatch, so a stale worker
checkout cannot silently run different code. The
[`runtime_policy`](/python-sdk/dispatch.md#the-worker-loop) flag
(`"strict" | "minor" | "off"`, default `"minor"`) governs that check.

When the plane exposes a code channel (`HttpDispatch`), a `ref` submit
also pushes the current git tree so a fresh worker can import your
modules — see [code mount](/python-sdk/happy-paths/10-code-mount.md).

## `dls.thunk` — an ad-hoc callable

For a lambda or a closure you did not decorate, `dls.thunk` bundles the
callable with its arguments for by-value submission:

```python
t = dls.thunk(lambda: score(model, val_ds))
inv = q.submit(t)                 # -> durable string id
val = q.result(inv, timeout=60)
```

`thunk(fn, /, *args, **kwargs)` returns a `Thunk`. Only the callable is
pickled; the arguments ride the codec. Passing a `UDF` unwraps to its
inner function. The receiving worker must allow pickled functions.

## Read next

- [Queue API](/python-sdk/queue.md) — the queue the id lives on.
- [Invocation ids](/python-sdk/invocations.md) — results, reconnect-by-id,
  fan-out + collect.
- [Dispatch planes](/python-sdk/dispatch.md) — where submits actually go,
  and how to run a worker.
