DreamLake

F3 — generators

The generator row of the execution matrix: results come back one at a time, and the consumer drives. Locally that is Python's own laziness. Remotely, each yield becomes a VALUE frame on the invocation's journal, and q.stream(inv) is the pull. Backpressure is the consumer's pace; durability is the journal; the since cursor makes every stream reconnectable — the property no Future ever had.

LevelSpelling
L0 one callfor y in f(x) — lazy pull, body runs on demand
L1 systemsq.stream(inv, timeout=) over journal frames; Topic pub/sub
L2 static pipelinestages as generators; the consumer pulls through the chain
L3 dynamic grapha .send() trampoline — local plain Python only, GAP at system level

L0 — a generator body is a generator

@dls.udf wraps generator bodies natively — no mode flag, no second decorator; inspect.isgeneratorfunction decides once at decoration. A plain call returns an iterator, and nothing runs until you pull.

python
import dreamlake.lakeshore as dls

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

it = frames(3)          # nothing has run yet
first = next(it)        # runs the body up to the first yield
rest = list(it)         # pulls the remainder

frames(3) prints nothing. The first print fires inside next(it) — laziness is the whole point of the row, and the local tier preserves it exactly. The wrapper installs a per-call run context only around each next(), never while control is back with the consumer, so a yield cannot leak the callee's scope into consumer code.

L1 — yields become journal frames

Bind the same body to a queue and each yield becomes a VALUE frame on the invocation's journal, followed by an END frame when the body returns. The consumer is q.stream(inv) — live while the job runs, from any process holding the id.

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

plane = Dispatch(":memory:")
q = dls.SyncQueue("render", dispatch=plane)

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

inv = render.submit(4, _dispatch=plane)     # durable string id
dls.run_worker(q, once=True)                # yields -> VALUE frames -> END

for frame in q.stream(inv, timeout=5.0):
    print(frame["frame"], frame["px"])

run_worker(q, once=True) executes one job; a generator job streams all its frames and ENDs during that call. In production the worker runs elsewhere and stream follows the journal while the body is still yielding — the frames are rows, not socket traffic, so producer and consumer never need to overlap in time.

If the body raises, the worker appends an ERROR frame instead of END, and q.stream re-raises it on the consumer side as a DispatchError naming the producer's exception type and message. A stream ends for exactly one of three reasons: an END frame, an ERROR frame, or a terminal invocation state with no further frames.

Cursor reconnect

Frame seqs are 0-based. since is the last seq you saw; since=-1 replays from the start. A new consumer — a different process, after a crash, an hour later — resumes without loss or duplication:

python
cursor = 1                                       # last seq this consumer saw
tail = list(q.stream(inv, since=cursor, timeout=5.0))
assert [f["frame"] for f in tail] == [2, 3]      # frames 0, 1 skipped

replay = list(q.stream(inv, since=-1, timeout=5.0))
assert len(replay) == 4                          # full history, any time

This is what F3 buys over a Future. A Future is process-bound state: once the process that held it dies, the stream is gone. The journal is a fact about the world — the invocation id plus a cursor reconstructs the consumer.

timeout here is patience for silence, not a deadline for the whole stream: every frame that arrives resets it. q.stream defaults to 30 seconds of silence before raising TimeoutError.

Topic — the same journal, spelled pub/sub

A topic is a named journal channel not tied to any invocation. Same frames, same cursor semantics. q.topic(name) returns one; an omitted name gets a generated topic-<ULID>.

python
t = q.topic("progress")
t.post({"done": 10, "of": 30})
t.post({"done": 20, "of": 30})
t.post({"done": 30, "of": 30})
t.post(None, final=True)         # END frame closes the stream

for msg in t.stream(since=-1):   # replay from the start
    print(msg)

late = t.listen(since=0, timeout=5.0)      # next value after seq 0
assert late == {"done": 20, "of": 30}

A late subscriber passes the last seq it saw and picks up mid-stream. listen returns one value, or None on timeout; stream iterates to END. Both default to timeout=5.0 — shorter than the invocation stream's 30, so do not assume one number covers both. One mechanism covers invocation streaming and pub/sub, so the cursor is learned once.

L2 — stages as generators

A static pipeline in this row is just composed generators. The consumer's pull propagates back through the chain; no stage runs ahead of demand.

python
import dreamlake.lakeshore as dls

@dls.udf
def load(paths):
    for p in paths:
        yield {"path": p}

@dls.udf
def score(rows):
    for r in rows:
        yield {**r, "score": len(r["path"])}

pipeline = score(load(["a/b.png", "c.png", "d/e/f.png"]))
best = max(pipeline, key=lambda r: r["score"])

first = next(score(load(["x.png", "never/loaded.png"])))

The second chain pulls exactly one item: "never/loaded.png" is never loaded. This composition is local-tier only — an iterator is in-process state and cannot ride the wire as an argument. Across queues, each stage's frames land on its own invocation journal, so every stage boundary is also a reconnect point.

L3 — the .send() trampoline (GAP)

Python generators are two-way: .send(value) delivers a value into the paused body. With plain generator functions that gives a driver/worker trampoline in a few lines:

python
def tuner():                     # a plain generator, NOT a @dls.udf
    lr = 0.1
    while True:
        loss = yield lr          # value OUT, measurement IN
        if loss is None:
            return
        lr = lr * 0.5 if loss > 1.0 else lr * 1.1

g = tuner()
lr = next(g)                     # prime: first suggestion
for loss in [2.0, 0.8, 0.5]:
    lr = g.send(loss)            # send the measurement, get the next lr
g.close()
A decorated generator does not forward `.send()`

@dls.udf returns a wrapper generator that advances the body with plain next() calls. A value sent into the wrapper is delivered to the wrapper's own yield and discarded — the body never sees it, and no error is raised. Two-way generators must stay undecorated today.

The system-level version — a remote generator whose journal carries frames in both directions, so a consumer can send a value into a paused remote body — does not exist. The journal is one-way: workers append, consumers pull. This cell is a named gap, not a hidden one. Today you spell two-way flow as two one-way channels — an invocation stream out, a Topic back in — with the loop in ordinary Python.

Read next

  • Queue API — stream, Topic, and the worker verbs in full.
  • F5 — async generators — the same row with an event loop underneath.
  • examples/05_streams_topics.py in the dreamlake-lakeshore repo — this page's L1 material as one runnable, zero-infra script.