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 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.
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:
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:
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):
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.
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:
Naming the gap precisely:
- No fan-in verb. There is no
merge(*invs)oras_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 andTaskGroupsurfaces it, but partial frames already sit in the merge queue. - A thread per live stream. Each open
Queue.streamoccupies 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 ofopen_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 — the sync row this column mirrors; journal frames and cursors in full.
- Queue API —
stream,Topic, and the worker verbs. examples/03_body_kinds.pyin 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.