@udf decorator
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.
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.
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:
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:
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). |
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
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:
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:
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.
@udf(queue=...) — remote by default
A plain call on a queue-bound UDF submits and waits. An async def or
async-generator body returns an awaitable instead of blocking.
.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). |
submit on a UDF with no queue binding raises:
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:
.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_LIMITis 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 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.
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:
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 — the queue the id lives on.
- Invocation ids — results, reconnect-by-id, fan-out + collect.
- Dispatch planes — where submits actually go, and how to run a worker.