DreamLake

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.

python
import dreamlake.lakeshore as dls

@dls.udf                          # local tier — a plain in-process call
def stats(xs: list[float]) -> dict:
    return {"mean": sum(xs) / len(xs)}

@dls.udf(queue="gpu")             # remote tier — submit and wait
def embed(text: str) -> list[float]:
    ...

The full signature is short, and there is nothing else on it:

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

No mode=, no retries=, no timeout=, no image=. Runtime shape belongs to the queue and the daemon, not to the decorator.

The two tiers

Formfn(...) 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:

MethodReturnsBehaviour
fn.local(*a, **kw)the valueAlways in-process, even on a queue-bound UDF.
fn.submit(*a, **kw)strEnqueue and return the invocation id immediately.
fn.remote(*a, **kw)the valueSubmit, 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.

python
inv = embed.submit("hello")           # -> "01JZ…" immediately
...                                   # crash, restart, different machine
q = dls.SyncQueue("gpu")
value = q.result(inv, timeout=60.0)   # re-attach by id

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:

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

bash
lakeshore worker start --queue gpu --allow-pickled

Thunks — ad-hoc callables

A Thunk bundles a non-importable callable with concrete arguments so it can be submitted by value without a decorator:

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

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:

python
@dls.udf
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

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

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:

ConditionPlane
LAKESHORE_URL is setHttpDispatch against that control plane. Token from LAKESHORE_CLIENT_TOKEN, namespace from LAKESHORE_NAMESPACE (default "default").
otherwiseThe 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.

Queue-bound UDFs need the Python worker

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:

python
@dls.udf(queue="gpu")
def bad(x):
    return weights @ x        # `weights` comes from the enclosing scope

@dls.udf(queue="gpu")
def good(x, weights):
    return weights @ x

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

  • @udf decorator — every parameter, the four body kinds, and the dls.run seam 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.