DreamLake

Simple functions

One decorator, four body kinds. The SDK inspects your function once at decoration time and the call convention follows the body you wrote. Everything on this page runs in-process, so you can paste any of it into a REPL.

Sync body

A plain def returns a plain value.

patterns.pypython
import dreamlake.lakeshore as dls

@dls.udf
def stats(xs: list[float]) -> dict[str, float]:
    return {"mean": sum(xs) / len(xs), "n": len(xs)}

print(stats([1.0, 2.0, 3.0]))
# {'mean': 2.0, 'n': 3}

Async body

An async def returns a coroutine. Await it.

patterns.pypython
import asyncio
import dreamlake.lakeshore as dls

@dls.udf
async def embed(x: float) -> float:
    await asyncio.sleep(0.01)
    return x * 2.0

async def main():
    print(await embed(21.0))
    # 42.0

asyncio.run(main())

Generator body

A generator returns an iterator. Nothing runs until the consumer pulls, so a long sequence costs only what you actually consume.

patterns.pypython
import dreamlake.lakeshore as dls

@dls.udf
def frames(n: int):
    for i in range(n):
        yield {"frame": i, "px": i * i}

it = frames(1000)          # nothing has executed yet
print(next(it))            # runs the body up to the first yield
# {'frame': 0, 'px': 0}

for f in frames(3):
    print(f)
# {'frame': 0, 'px': 0}
# {'frame': 1, 'px': 1}
# {'frame': 2, 'px': 4}

Async-generator body

An async def with yield returns an async iterator. Pull it with async for.

patterns.pypython
import asyncio
import dreamlake.lakeshore as dls

@dls.udf
async def ego_views(traj: str, n: int):
    for i in range(n):
        await asyncio.sleep(0)
        yield {"traj": traj, "frame": i}

async def main():
    views = [v async for v in ego_views("scenes/0007", 3)]
    print(views)

asyncio.run(main())

Chaining

A UDF's output feeds the next one. With plain values, chain by awaiting in sequence.

patterns.pypython
import asyncio
import dreamlake.lakeshore as dls

@dls.udf
async def decode(path: str) -> list[int]:
    return [ord(c) for c in path]

@dls.udf
async def embed(codes: list[int]) -> float:
    return sum(codes) / len(codes)

async def main():
    print(await embed(await decode("scenes/0007")))

asyncio.run(main())

With streaming bodies, chain by passing the iterator itself into the next stage. Each stage stays lazy, so the whole pipeline pulls one item at a time.

patterns.pypython
import dreamlake.lakeshore as dls

@dls.udf
def decode(path: str):
    for i, c in enumerate(path):
        yield {"i": i, "code": ord(c)}

@dls.udf
def embed(rows):
    for row in rows:
        yield {**row, "vec": row["code"] * 0.5}

@dls.udf
def dedupe(rows):
    seen = set()
    for row in rows:
        if row["code"] in seen:
            continue
        seen.add(row["code"])
        yield row

for row in dedupe(embed(decode("scenes/0007"))):
    print(row)

Fan-out over async bodies

Async UDF calls are ordinary awaitables, so asyncio.gather fans them out with no extra machinery.

patterns.pypython
import asyncio
import dreamlake.lakeshore as dls

@dls.udf
async def score(seed: int) -> float:
    await asyncio.sleep(0.05)
    return seed * 1.5

async def main():
    results = await asyncio.gather(*(score(i) for i in range(8)))
    print(results)
    # [0.0, 1.5, 3.0, 4.5, 6.0, 7.5, 9.0, 10.5]

asyncio.run(main())

Scopes and key-based returns

dls.scope(prefix) sets the ambient prefix that dls.run.read and dls.run.write resolve against. The body names prefix-relative string keys and returns the keys it wrote — never the bytes.

patterns.pypython
from pathlib import Path
import dreamlake.lakeshore as dls

@dls.udf
def normalize(raw: str) -> str:
    text = dls.run.read(raw).read_text()
    out = dls.run.write("process/clean.txt")
    out.write_text(text.strip().lower())
    return "process/clean.txt"

root = Path("/tmp/lakeshore-demo")
with dls.scope("scenes/0007", root=root):
    key = normalize("source/raw.txt")
    print(key)
    # process/clean.txt

Every call runs in its own forked context, and the contract is checked on the way out: every key you wrote must appear in what you return. Return a single string, a list or tuple of strings, or a dict whose values are the keys.

patterns.pypython
@dls.udf
def split(raw: str) -> dict[str, str]:
    for name in ("train", "val"):
        dls.run.write(f"process/{name}.txt").write_text(name)
    return {"train": "process/train.txt", "val": "process/val.txt"}

Scopes nest, so a parent prefix composes with a child one:

patterns.pypython
with dls.scope("scenes/0007", root=root):
    with dls.scope("cameras/front"):
        print(dls.run.write("rgb/0000.png"))
        # /tmp/lakeshore-demo/scenes/0007/cameras/front/rgb/0000.png

Generator bodies may yield partial manifests; they are concatenated and checked once the generator is exhausted.

patterns.pypython
@dls.udf
def shard(n: int):
    for i in range(n):
        dls.run.write(f"process/shard-{i}.txt").write_text(str(i))
        yield f"process/shard-{i}.txt"

with dls.scope("scenes/0007", root=root):
    print(list(shard(3)))
    # ['process/shard-0.txt', 'process/shard-1.txt', 'process/shard-2.txt']

Declaring what it may touch

A scope sets where keys resolve. It does not say what the body is allowed to reach. Declaring the readable source, the writable target, and S3 mounts on the decorator — or once for the whole repo — is its own page: Declaring access.

The registry payload

A simple function does not travel as source. It is registered — the control plane keeps its identity and signature, and the invocation carries only a reference to that record plus the arguments.

registered.pypython
import dreamlake.lakeshore as dls

@dls.udf(queue="gpu")
def encode(
    src: str,                    # Input video key, resolved against the run scope
    dest: str = "frames/",       # Output key prefix for the extracted frames
    fps: int = 4,                # Frames sampled per second of source video
    overwrite: bool = False,     # Re-encode even when dest already holds frames
) -> dict[str, str]:
    """Sample frames from a video and write them into the run scope.

    Returns the written keys so the caller does not have to re-list the prefix.
    """
    ...

inv = encode.submit(src="raw/run01.mp4", fps=8)

The inline comment on each parameter is its description — the same convention a command line program uses to generate --help. One rule covers both UDF types.

That decoration produces one Function row, keyed by (namespaceId, module, qualname, version):

Functionjson
{
  "namespaceId": "ns_7Yc…",
  "module":      "pipelines.video",
  "qualname":    "encode",
  "version":     "9f2c1e…",
  "source":      "@dls.udf(queue=\"gpu\")\ndef encode(src: str, fps: int = 4, …",
  "signature": {
    "params": [
      { "name": "src",       "annotation": "str",  "default": null,
        "doc": "Input video key, resolved against the run scope" },
      { "name": "dest",      "annotation": "str",  "default": "frames/",
        "doc": "Output key prefix for the extracted frames" },
      { "name": "fps",       "annotation": "int",  "default": 4,
        "doc": "Frames sampled per second of source video" },
      { "name": "overwrite", "annotation": "bool", "default": false,
        "doc": "Re-encode even when dest already holds frames" }
    ],
    "returns": "dict[str, str]",
    "doc": "Sample frames from a video and write them into the run scope. …"
  },
  "createdAt": "2026-08-06T17:02:11.480Z"
}

Three things about that record are worth knowing:

Why the comment, and not a docstring param list

A Google- or NumPy-style Args: block is a second copy of the signature. Rename a parameter and the docstring keeps the old name, silently — nothing checks them against each other. An inline comment sits on the line it describes, so it moves with the parameter or disappears with it.

Extraction costs nothing new: the SDK already captures source at decoration time for the source column, so the comments are in hand before anything is parsed. typing.Annotated[str, "…"] is the runtime-visible alternative and stays available for cases that need it — it is just noisier for the common one.

  • version is the sha256 of the source, not a number you bump. Edit the body and you get a new row; the old one stays, so an invocation recorded last week still names the code that actually ran.
  • source is captured at decoration time and is what the dashboard shows when you click into an invocation.
  • signature is where the prose lands. signature.doc is __doc__, and each entry in signature.params carries its own doc taken from the parameter's inline comment. That is what makes the registry self-describing — a reader asking "what is dest?" gets an answer without opening the source.
Best-effort, and one place it degrades

Comments are discarded by the interpreter, so this is a source parse. Where source cannot be retrieved — a function defined in a REPL, built by exec, or shipped by value because it was not importable — the parameters still register with names, annotations and defaults, and simply carry no doc. Missing prose is never an error, and it is the same condition that already pushes transport from ref to pickle.

The envelope on the wire is much smaller — it names the function rather than carrying it:

envelope (abridged)json
{
  "kind":     "ref",
  "module":   "pipelines.video",
  "qualname": "encode",
  "queue":    "gpu",
  "args":     { "src": "raw/run01.mp4", "fps": 8 }
}

kind: "ref" is the importable case — the worker imports module:qualname, fetching and extracting the code snapshot only on ImportError, cached by commit. kind: "pickle" is the fallback for lambdas, closures, and functions defined in __main__, where there is nothing importable to name.

The native path does not register yet

The Function model above is real and the dashboard reads it, but the native HTTP dispatch currently sends a fixed stub — {"module": "_native_wire", "qualname": "_native_wire.envelope", "version": "v1"} — rather than the function's own identity. So a registry keyed by real module and qualname is the design, not yet the behaviour on that path. Worker-side resolution already works from module:qualname or an explicit registry={qualname: fn}.