Functions
A Function is a Python callable decorated with @dls.udf. The
decorator has exactly one job: decide whether calling the function runs
it here or submits it somewhere else.
The full signature is short, and there is nothing else on it:
No mode=, no retries=, no timeout=, no image=. Runtime shape
belongs to the queue and the daemon, not to the decorator.
The two tiers
| Form | fn(...) does |
|---|---|
@udf (bare) | Runs the body in-process and returns the value. Imports nothing distributed. |
@udf(queue="…") | Submits, then blocks for the result — equivalent to fn.remote(...). |
Three entry points sit on the decorated object:
| Method | Returns | Behaviour |
|---|---|---|
fn.local(*a, **kw) | the value | Always in-process, even on a queue-bound UDF. |
fn.submit(*a, **kw) | str | Enqueue and return the invocation id immediately. |
fn.remote(*a, **kw) | the value | Submit, then fetch the result. An async def body returns an awaitable instead of blocking. |
submit() accepts three private keyword arguments alongside your own:
_dispatch= to target a specific plane, _key= for an idempotency key,
and _priority=.
There is no Future type
A pending job is its invocation id — a plain string (a ULID). Persist it, log it, hand it to another process; anything that can reach the same dispatch plane can re-attach with it.
Fan-out is therefore a list comprehension and the join is a loop (or
asyncio.gather over the async Queue.result) — there is no gather()
helper and no as_completed() in the SDK.
The four body kinds
The function's own shape picks the execution shape — sync returns the
value, async returns a coroutine, a generator body returns an iterator,
and an async-generator body returns an async iterator. inspect detects
which one once at decoration, and the UDF object exposes the result as
.is_async, .is_gen and .is_asyncgen.
See the four body kinds on the
@udf reference for the worked examples of all four.
How a callable crosses the wire
Envelope.fn is a tagged union with two arms, chosen by transport:
| Transport | Wire shape | Meaning |
|---|---|---|
ref | {kind, module, qualname, src_hash} | The worker imports the module and walks the qualname, so it runs its own checked-out code. src_hash is a version stamp — a mismatch fails loud rather than running stale code. |
pickle | {kind, blob, hash, source, qualname, runtime} | A by-value cloudpickle snapshot, for lambdas, closures, and __main__ defs. |
transport="auto" — the default — picks pickle exactly when the
callable is not importable by reference, and ref otherwise.
transport="ref" insists on importability and fails loud if it cannot.
Only the callable is ever pickled. Arguments, keyword arguments,
results, and journal frames all ride the msgpack + ExtType codec.
Workers refuse pickled functions unless started with allow_pickled —
unpickling is code execution, so accepting it is a deliberate act of
trust. The inline pickle ceiling is 32 KiB.
Thunks — ad-hoc callables
A Thunk bundles a non-importable callable with concrete arguments so
it can be submitted by value without a decorator:
The thunk's callable ships under the same auto rule (so in practice
by pickle); its args and kwargs ride the normal codec.
Data UDFs return keys, not bytes
A function that touches files never sees a bucket, a credential, or an absolute path. It asks the ambient run context for local paths and returns plain string keys:
The return-your-keys contract is enforced after every call: every key
written through dls.run.write must appear in the returned manifest. The
reverse is deliberately not enforced — a returned string that was never
written may be plain data or an upstream key passed through.
File bytes never ride the msgpack wire; only keys do. See
@udf decorator for the full dls.run surface and
the current limits on storage-backed remote I/O.
Where the call goes
_default_dispatch() resolves the process-wide plane once:
| Condition | Plane |
|---|---|
LAKESHORE_URL is set | HttpDispatch against that control plane. Token from LAKESHORE_CLIENT_TOKEN, namespace from LAKESHORE_NAMESPACE (default "default"). |
| otherwise | The local SQLite plane at $LAKESHORE_HOME/dispatch.sqlite (default ~/.lakeshore). |
Pass _dispatch= on submit() — or dispatch= on a Queue /
SyncQueue — to override it per call.
nymph, the Rust daemon, does not execute queue-bound @udf bodies.
Run lakeshore worker start --queue <name> (which spawns
python -m dreamlake.lakeshore.daemon) on a host that can import your
code. The worker deliberately ignores saved lakeshore auth login
credentials — pass --url or set LAKESHORE_URL.
What the decorator does not capture
The decorator does not capture closures under the ref transport. If
your function depends on outer state, pass it as an argument:
Under pickle a closure is snapshotted, which is why auto falls back
to it — but that snapshot freezes at submit time, which is rarely what
you want for production code.
Read next
@udfdecorator — every parameter, the four body kinds, and thedls.runseam in full.- Queue API —
Queue/SyncQueue, streaming, and the worker verbs. - Invocation ids — the no-Future model and reconnect-by-id.
- Queues — the operator surface for the queue the decorator binds to.
- Function protocol — the wire-level shape.