DreamLake

API design decisions

TL;DR. @udf(mode=...) is canonical; modes are runtime expressions. Sync vs async is decided at def, not at the call site. Decorated calls produce thunks; spawn / call / apply_async dispatch them. Pipelines and run-scoping live at the top level, not nested.

The user-facing API went through several iterations before it settled. This page records the decisions and the alternatives that lost, so the design isn't re-litigated.

1. Decorator shape

Decision. Canonical form is @udf(mode=...). udf.compute was added briefly and removed.

examples/hello.pypython
@udf(mode="local" if DEBUG else "mit-supercloud")
def train(seed: int) -> Model: ...

mode accepts a string tag or a ComputeMode instance. The shorthand is load-bearing: mode selection is a runtime expression, not a deploy-time config flip.

Why no udf.compute alias

One entry point is easier to reason about than two synonyms. The alias made the design feel like there were two execution models when there is only one.

2. Sync vs async

Decision. The function declaration decides sync/async. There are no .aio / .remote shims on the function itself.

RemovedReplaced by
fn.aioasync def — the function is async natively
fn.remoteRemote dispatch is configured at the decorator

apply_async requires sync function inputs:

python
@udf
async def fetch(...): ...

udf.apply_async(fetch(...))   # raises TypeError — await directly instead

3. Thunks and dispatch

Decision. Calling a decorated function invokes it — train(...) runs sync, await fetch(...) runs async. To get a deferred handle (a thunk) you go through a separate API.

python
# Direct call: invokes.
result = train(seed=42)               # sync — runs and returns
coro   = fetch(seed=42)               # async — returns a coroutine

# Thunk: deferred handle.
thunk = udf.bind(train, seed=42)       # build a deferred call (no dispatch)
udf.spawn(thunk)                       # fire and forget
udf.call(thunk)                        # block; sync/async per fn
udf.apply_async(thunk)                 # returns Future
udf.apply_spawn(thunk)                 # spawn-many variant
PrimitiveBehavior
udf.spawn(thunk)Fire and forget
udf.call(thunk)Block; sync or async per the underlying fn
udf.apply_async(t)Always returns a Future
udf.apply_spawn(t)Spawn-many variant of spawn
Decided: thunk constructor is udf.bind

The constructor is udf.bind(fn, *args, **kwargs). It reads as English — "bind these args to this function" — and the verb/noun split is clean: bind() is the verb that constructs, Thunk is the noun it returns. bind collides conceptually with functools.partial / method binding / DI containers, but that's a guide, not a trap: all of those mean roughly the same thing (associate args with a callable), so a Python reader guesses the right shape on first sight. Considered: udf.thunk (literal but overloaded with the type name); udf.defer (verb, but reads like an action you're taking now, not a handle); udf.task (collides with asyncio Tasks / Celery / Airflow, all of which already run — the opposite of cheap-and-inert).

Why split call from thunk

Overloading train(...) to mean both "invoke" and "build a thunk" forces every reader to know the context to predict what the line does. A dedicated constructor keeps the call site honest: train(...) always invokes, udf.bind(train, ...) always defers.

4. Pipelines at top level

Decision. @dls.pipeline lives parallel to @udf and @agent, not nested under @udf.

python
import dreamlake.lakeshore as dls

@dls.pipeline
def my_workflow(...): ...

Rejected: nesting under @dls.udf.pipeline. Pipelines compose UDFs; they are not a flavor of UDF.

5. Run scoping

Decision. dls.run.emit / log / progress stay under the run namespace.

ConsideredDecided
dls.emit(…)dls.run.emit(…) — scope kept

The .run. scope clarifies that these messages belong to a specific run, not the module-level namespace. Flat naming would have made it read like a static log sink.

6. Iteration shapes

The first-class shape for stream processing is the inline for loop:

python
for batch in DataSource().chunks:
    results = video_segmentation(batch)
    scoring = vlm_agent @ f"""
        Review these {results} per @review_criteria.md and return
        {{ quality: int, alt_label: str, start: int, end: int }}
    """
  • DataSource().chunks yields batches.
  • vlm_agent @ "<prompt>" runs the VLM agent against the prompt and its referenced data.

A second motivating example, for batch jobs that don't produce a result:

python
for batch in DataSource().chunks:
    run_colmap(batch, save_to="./poses")

These two shapes drove the dispatch / iteration design — chunked sources composed with single-call UDFs, no special pipeline DSL required.

Open items

  • Zero-copy data path. UDF-to-UDF data on the same host currently round-trips through serialization. Planned: shared memory / Arrow IPC.