DreamLake

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

FormPlain 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

MemberSignatureDoes
run.write(key, *, dir=False) -> PathRegister an output key; return the local path to write at. Parents created; dir=True makes the directory itself.
run.read(key) -> PathResolve an input key to a local path. Raises KeyError_ if it does not exist.
run.read_all / run.write_all(keys) -> same shapeShape-mirroring: str -> Path, list -> list, dict -> dict.
run.prefix / run.rootpropertiesThe ambient prefix (str) and filesystem root (Path).
Storage-backed remote I/O is not wired yet

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 didYou get
wrote via dls.run.write, returned a non-manifest valueKeyError_: UDF wrote N file(s) via dls.run.write (...) but did not return their keys
wrote a key you left out of the manifestKeyError_: UDF wrote key(s) it did not return: [...]
returned a key you never wroteallowed — plain data or a pass-through key
returned plain data, wrote nothingallowed — 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.

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

Plain remote calls inherit a 30 s timeout

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

KwargDoes
_keyIdempotency key — resubmitting the same key on the same queue returns the existing invocation id, no duplicate job.
_priorityFloat attached to the row. fifo and filo claim highest priority first; boltzmann samples proportional to exp(priority / T). Default 0.0.
_dispatchOverride the dispatch plane (default: the process plane — see Dispatch planes).

submit on a UDF with no queue binding raises:

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.

transportShipsFails when
"auto" (default)ref when the callable is importable, else pickle—
"ref"module:qualname plus a src_hash version stampthe callable is a lambda, a <locals> closure, a __main__ def, or a callable instance
"pickle"a cloudpickle snapshot plus a runtime stampthe 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 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:

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