DreamLake

F1 — sync blocking

F1 is the first row of the execution matrix: invoke, then stand still until the value comes back. It is the least clever form and the one to exhaust before reaching for ids, generators, or async. Most code is control flow around values, and a blocking call is the shortest path from a call site to a value.

LevelSpellingBlocks on
L0 one callf(x) · f.local(x)the body itself
L1 queue fan-outq.result(inv, timeout=) per idone invocation
L2 static pipelinelist-resolve between stagesthe slowest job in a stage
L3 dynamic graphwhile / if around blocking callseach step, by design

L0 — one call

A bare @dls.udf runs in-process on a plain call, inside a fresh child of the ambient run context. This is the whole local tier: no plane, no worker, no wire, and no import of httpx or msgpack.

python
import dreamlake.lakeshore as dls

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

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

On a queue-bound UDF the plain call keeps F1 semantics — it routes to .remote(), which is submit plus wait. Same call site, same shape of code, different tier:

python
import threading
import dreamlake.lakeshore as dls
from dreamlake.lakeshore.dispatch import Dispatch

plane = Dispatch(":memory:")
q = dls.SyncQueue("f1-demo", dispatch=plane)

@dls.udf(queue="f1-demo")
def embed(x: str) -> list[float]:
    return [float(len(x))]

done = threading.Event()
worker = threading.Thread(target=dls.run_worker, args=(q,), kwargs={"stop": done.is_set})
worker.start()

embed("hello", _dispatch=plane)   # [5.0] — blocks until the worker answers
done.set()
worker.join()
A plain remote call has a 30-second ceiling

.remote() resolves through Dispatch.result(inv_id) with no timeout argument, so it inherits that method's 30-second default and raises TimeoutError past it. A plain call is therefore only F1-shaped for jobs that finish inside 30 seconds. For anything longer, submit and resolve explicitly: q.result(f.submit(x), timeout=3600).

f.local(x) forces the local tier even when a queue is bound, so debugging and unit tests stay plain function calls:

python
embed.local("hello")   # [5.0] — no queue, no network

.local() builds a fresh unbound copy of the function, which means it does not carry the original transport= setting across. That is irrelevant for local execution — nothing is serialized — but do not read .local() as "the same UDF, minus the queue".

Data UDFs are F1 at L0 too: under with dls.scope(prefix): a plain call reads and writes through the dls.run seam and returns plain string keys. See F6 — scopes for the key contract, and @udf for the decorator reference.

L1 — fan-out, blocking resolve

At L1 the F1 spelling is q.result(inv, timeout=) — the blocking resolve for one durable invocation id. Fan out with submit, then block per id, in submission order:

python
import dreamlake.lakeshore as dls
from dreamlake.lakeshore.dispatch import Dispatch

plane = Dispatch(":memory:")
q = dls.SyncQueue("f1-demo", dispatch=plane)

@dls.udf(queue="f1-demo")
def square(x: int) -> int:
    return x * x

invs = [square.submit(i, _dispatch=plane) for i in range(4)]
for _ in invs:
    dls.run_worker(q, once=True)     # zero-infra drain: one call per job

results = [q.result(inv, timeout=5) for inv in invs]   # blocks, in order
assert results == [0, 1, 4, 9]

There is no gather. A list comprehension of submit calls plus one result per id is the fan-out, and it is deliberate: every element of the fan-out is a string that outlives the process holding it.

The cost is head-of-line waiting. Jobs still run concurrently on the workers, so total wall time is the slowest job rather than the sum — but you observe completions in submission order, and a fast finisher at the end of the list sits resolved and unread while you wait on an earlier slow one. Reacting to whichever finishes first is the F2 row's problem, not a reason to make this code async.

L2 — static pipeline, the barrier

Chain stages with ordinary calls; the list-resolve between stages is a barrier. Stage 2 sees plain values, never ids:

python
@dls.udf(queue="f1-demo")
def total(xs: list[int]) -> int:
    return sum(xs)

inv2 = total.submit(results, _dispatch=plane)   # stage-1 values feed stage 2
dls.run_worker(q, once=True)
q.result(inv2, timeout=5)                       # 14

Barrier semantics are the point, not a limitation: every stage boundary is a checkpoint where all upstream values exist and the whole DAG so far is debuggable with print. The price is that no stage-2 work starts until the slowest stage-1 job lands. For wide, uneven stages that idle time is real — pipeline the ids instead, as examples/02_pipeline.py does.

L3 — dynamic graph, blocking on purpose

When the next step depends on the last value, blocking is not a compromise — it is the semantics. A rework loop cannot proceed without the result it is looping on:

python
x = 3
while x < 100:
    inv = square.submit(x, _dispatch=plane)
    dls.run_worker(q, once=True)
    x = q.result(inv, timeout=5)     # each step waits — by design

x   # 6561

The graph this loop traces was not knowable at author time; its shape is the sequence of resolved values. Conditionals, retries, and data-dependent fan-out all read this way: ordinary Python control flow, with a blocking resolve wherever an edge needs a value. Width comes back only when a single dynamic wave is itself wide — then the wave fans out via F2 and F1 blocks once at the wave boundary.

When to reach for F1

Default to it. F1 is correct whenever the caller has nothing better to do than wait: scripts, tests, .local() debugging, stage barriers you actually want, and any dynamic step whose successor needs the value. Leave it only for a named reason — overlapping uneven fan-out (F2), streaming partials (F3/F5), or a caller that must stay responsive (F4).

Runnable end to end in the dreamlake-lakeshore repo: examples/02_pipeline.py (the two-stage barrier) and examples/04_ids_not_futures.py (the same levels, spelled with ids).