# F5 — async generators

The streaming row of the async column: values arrive one at a time, the
consumer pulls with `async for`, and the event loop stays free between
pulls. F5 is [F3](/python-sdk/architecture/f3-generators.md) with the loop underneath —
same journal at system level, same cursor reconnect — but the consumer can
hold many streams open at once without a thread per stream *in its own
code*. Where the SDK itself still burns a thread per pull is this page's
named gap.

| Level | Spelling |
| --- | --- |
| L0 one call | `async for y in f(x)` — the body is an async iterator |
| L1 systems | `async for f in Queue.stream(inv)` — offloaded sync, not native async |
| L2 static pipeline | async-gen stages composed; `async for` all the way down |
| L3 dynamic graph | `TaskGroup` plus `asyncio.Queue` fan-in — works, hand-rolled, **GAP** |

## L0 — an async-gen body is an async iterator

`@dls.udf` wraps async-generator bodies natively —
`inspect.isasyncgenfunction` detects the kind at decoration, and a plain
call returns an async iterator. Nothing runs until the first `__anext__`;
each step may await real I/O.

```python
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)          # real async I/O goes here
        yield {"traj": traj, "frame": i}

async def main():
    frames = [f async for f in ego_views("scenes/0007", 3)]
    assert [f["frame"] for f in frames] == [0, 1, 2]

asyncio.run(main())
```

Laziness and backpressure work exactly as in F3; the only difference is that
the pull is an `await`, so the consumer can interleave other work between
frames. Yields of string keys are still partial manifests — the
return-your-keys contract concatenates them into one manifest and checks it
when the body finishes, same as the sync row.

## L1 — `Queue.stream`: what it actually is today

`dls.Queue` is the async sibling of `dls.SyncQueue` — a facade holding the
blocking queue as `._q`, with every verb offloaded. `Queue.stream` is a real
async generator with correct semantics, but its engine is the blocking
`SyncQueue.stream` iterator, pulled one item at a time:

```python file="dreamlake/lakeshore/queues.py"
async def stream(self, inv_id, *, since=-1, timeout=30.0):
    it = self._q.stream(inv_id, since=since, timeout=timeout)
    while True:
        item = await asyncio.to_thread(next, it, _SENTINEL)
        if item is _SENTINEL:
            return
        yield item
```

So: **async iterator, yes; native async I/O, no.** Each pull parks one
executor thread inside the blocking journal poll until a frame lands. The
event loop is never blocked — that part of the contract holds — but the
concurrency ceiling is the thread pool, not the loop. The `Queue` class
docstring reserves the fix: an httpx-native async transport slots in behind
this same surface when the remote control plane wire lands. It is not built.

It works today, end to end, including *live* consumption — worker and
consumer overlapping in time:

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

@dls.udf(queue="render")
async def render(n: int):
    for i in range(n):
        await asyncio.sleep(0)
        yield {"frame": i}

async def main():
    plane = Dispatch(":memory:")
    sync_q = dls.SyncQueue("render", dispatch=plane)
    q = dls.Queue("render", dispatch=plane)

    inv = render.submit(4, _dispatch=plane)      # durable string id

    async with asyncio.TaskGroup() as tg:
        tg.create_task(asyncio.to_thread(dls.run_worker, sync_q, once=True))
        frames = [f async for f in q.stream(inv, timeout=10)]

    assert [f["frame"] for f in frames] == [0, 1, 2, 3]

asyncio.run(main())
```

On the worker side the async-gen body is driven under `asyncio.run` and each
yield becomes a VALUE frame on the invocation's journal, with END on return
— identical frames to the sync-generator row. The cursor works unchanged
through the async surface (`since` is the last 0-based seq you saw):

```python
tail = [f async for f in q.stream(inv, since=1, timeout=5)]
assert [f["frame"] for f in tail] == [2, 3]
```

## L2 — async-gen stages chained

Stages compose the way generators always have; the outermost `async for`
pulls through the whole chain, one item at a time, awaiting at every hop. No
stage runs ahead of demand.

```python
import asyncio
from typing import AsyncIterator
import dreamlake.lakeshore as dls

@dls.udf
async def decode(path: str):
    for i in range(4):
        await asyncio.sleep(0)
        yield {"frame": i, "src": path}

@dls.udf
async def embed(frames: AsyncIterator[dict]):
    async for f in frames:
        await asyncio.sleep(0)
        yield {**f, "emb": [f["frame"] * 0.5]}

@dls.udf
async def dedupe(embs: AsyncIterator[dict]):
    seen = set()
    async for e in embs:
        key = tuple(e["emb"])
        if key not in seen:
            seen.add(key)
            yield e

async def main():
    pipeline = dedupe(embed(decode("clips/0001.mp4")))
    out = [x async for x in pipeline]
    assert len(out) == 4

asyncio.run(main())
```

This chaining is **local-tier only**: an async iterator is in-process state
and cannot ride the msgpack wire as an argument. To chain across queues you
put a stream at each stage boundary — stage A yields onto its invocation
journal, stage B consumes `q.stream(inv_a)` and yields onto its own. The
frames compose; the iterators do not.

## L3 — merging a set of streams (GAP)

The dynamic-graph cell: `async for` over several invocation streams at once,
taking frames in whatever order they land. There is no merge primitive in
the SDK — the legacy `as_completed` helper is gone and nothing stream-shaped
replaced it — so today you hand-roll fan-in with a `TaskGroup`, an
`asyncio.Queue`, and a sentinel per stream. Clunky, but it runs:

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

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

async def main():
    plane = Dispatch(":memory:")
    sync_q = dls.SyncQueue("render", dispatch=plane)
    q = dls.Queue("render", dispatch=plane)

    invs = [
        render.submit("scenes/0007", 3, _dispatch=plane),
        render.submit("scenes/0008", 2, _dispatch=plane),
    ]

    DONE = object()
    merged: asyncio.Queue = asyncio.Queue()

    async def pump(inv: str) -> None:
        async for frame in q.stream(inv, timeout=10):
            await merged.put(frame)
        await merged.put(DONE)

    async with asyncio.TaskGroup() as tg:
        for _ in invs:                                   # one drain per job
            tg.create_task(asyncio.to_thread(dls.run_worker, sync_q, once=True))
        for inv in invs:
            tg.create_task(pump(inv))

        open_streams = len(invs)
        got = []
        while open_streams:
            item = await merged.get()
            if item is DONE:
                open_streams -= 1
                continue
            got.append(item)

    assert len(got) == 5

asyncio.run(main())
```

Naming the gap precisely:

- **No fan-in verb.** There is no `merge(*invs)` or `as_completed(invs)`
  over streams, so every caller re-invents the pump-and-sentinel dance above
  — and re-invents error propagation with it. A failed stream raises inside
  its pump task and `TaskGroup` surfaces it, but partial frames already sit
  in the merge queue.
- **A thread per live stream.** Each open `Queue.stream` occupies one
  executor thread inside the blocking journal poll, per the L1 shape above.
  Two streams is fine; hundreds saturate asyncio's default thread pool. The
  honest ceiling today is tens of concurrent streams, and the fix is the
  native-async transport, not user code.
- **No growing-set spelling.** Adding a stream to the merge after it started
  is `tg.create_task(pump(new_inv))` plus manual bookkeeping of
  `open_streams` — workable, but the counting is on you.

What holds even in the clunky version: the merge inputs are invocation ids,
so a crashed merger restarts from persisted `(inv, cursor)` pairs and loses
nothing. The gap is convenience and thread economics, not durability.

## Read next

- [F3 — generators](/python-sdk/architecture/f3-generators.md) — the sync row this column
  mirrors; journal frames and cursors in full.
- [Queue API](/python-sdk/queue.md) — `stream`, `Topic`, and the worker verbs.
- `examples/03_body_kinds.py` in the dreamlake-lakeshore repo — all four
  body kinds, including this row's L0, as one runnable script.
- `examples/05_streams_topics.py` — the journal and cursor mechanics this
  page's L1 rides, zero infra.
