DreamLake

The execution matrix

Python gives you an unusually large number of ways to express "run this" — sync calls, futures, generators, coroutines, async generators, context managers, run objects — and the pre-rewrite API had accreted most of them. This page is the matrix that framed the question before any opinion: two axes, the form (the Python execution pattern) and the level (what you are running).

The question as originally posed: generators, async generators, sync, async, pipeline run object — what are the different ways we have, for running the basics, the systems, and the system-level imperative dynamic graph?

Axis 1 — the form

The full vocabulary Python offers, and how you get values out of each:

FormInvokeGet resultsDriverBarrier?
F1 sync blockingy = f(x)return valuecaller—
F2 deferred handleh = f.submit(x)resolve h latercaller pollsdeferred
F3 sync generatorfor x in g()lazy pull, one at a timeconsumer (backpressure)streaming
F3′ gen with .send()x = yield ytwo-way (trampoline / scheduler)schedulerstreaming
F4 async coroutineawait f(x)awaited valueevent loopdeferred
F5 async generatorasync for x in g()lazy pull plus awaitconsumer plus loopstreaming
F6 context managerwith X() as h: / async withscoped lifecycle——
F7 run/handle objectr = pipe.submit(...)poll status, stream events, awaitcaller trackseither

Each form F1–F7 has its own page in the Patterns section, expanded into runnable examples against the shipped SDK.

Axis 2 — the level

  • L0 basic — one UDF, one invocation.
  • L1 systems — N calls across queues and workers; the scheduling substrate.
  • L2 static pipeline — a fixed multi-stage DAG, knowable at author time (examples/02_pipeline.py).
  • L3 dynamic imperative graph — structure discovered while running: data-dependent fan-out width, conditional branches, rework loops, recursion.

The original grid — drawn against the pre-rewrite API

Historical. None of this API exists.

Every symbol in the grid below (apply, apply_spawn, apply_async, gather, map, as_completed, Future, QueueFuture, Pool, rpc_future, @pipeline) was deleted in the native rewrite. The grid is preserved because the shape of the analysis — where the mass sat, where the holes were — is what drove the design. For what to write today, skip to the re-grounded grid below.

L0 basicL1 systems / fan-outL2 static pipelineL3 dynamic graph
F1 sync blockingapply(fn, …), fut.get(), @udf localQueue.map(fn, xs), gather([…])gather → gather (barriers)while loop plus blocking calls plus requeue (serial)
F2 sync plus futureapply_spawn → Future; @udf(queue=q) → QueueFuture[fn(i) for i …] → list[Future]; as_completeda stage-1 Future passed into a stage-2 call (implicit edge)spawn futures off resolved values; parent_invocation_id builds the tree
F3 sync generatorfor chunk in q.rpc_future(…) streamas_completed yields by completion; draft batch(src, n)draft @pipeline (a generator that yields to_dataset(…)).send()-driven dynamic pull → GAP
F4 async coroutineawait apply_async, await futawait gather_async, map_asyncan await chain of stagesasyncio.create_task dynamic spawn
F5 async generatorstreaming one call via async for → GAPas_completed_async(futs)async-gen stages → GAPasync for over a growing stream → GAP
F6 context managerwith Queue(…) as q:Pool(size=4)——
F7 run object——GAP (@pipeline returned the function unchanged)GAP — no run you can poll, stream, or await

What the grid showed

  • The legacy API lived almost entirely in F1, F2, and F4 — blocking, futures, async. That was roughly fifteen symbols spent on one region: apply / spawn / call in sync and async spellings, gather / map / as_completed in sync and async spellings, and two Future classes.
  • The dataflow-primitives draft lived in F3 — batch, @pipeline, all generator-shaped. A different form, which is why it read so differently from the client.
  • The whole F7 row was empty. @pipeline was a stub that returned the function unchanged. The draft's Graph is an author-time structure, not a live handle, and callee-side run telemetry is not a caller-side tracker.
  • Column L3 was only reachable imperatively. It was real — the parent_invocation_id tree captured it at runtime — but never a first-class object you could hold, stream, or visualize.

Two cross-cutting tensions fall straight out of the grid: barrier versus streaming (F1/F2 gather against F3/F5 generators), and static versus dynamic (L2 knowable at author time against L3 discovered at run time). The empty F7 row is exactly where those two reconcile — a run object is the thing that can be both streamed and dynamic.

Re-grounded: the shipped grid (0.3.7)

The rewrite deleted the legacy surface and collapsed most of the form axis into the function body itself:

FormShipped spellingStatus
F1 syncf(x) — local tier, or submit-and-wait when a queue is boundshipped
F2 deferredinv = f.submit(x) then q.result(inv)shipped (ids, not Futures)
F3 generatorgenerator body; q.stream(inv, since=cursor)shipped
F4 asyncasync def body; dls.Queue, the async siblingshipped
F5 async-genasync def generator body; async for … in q.stream(…)shipped
F6 scopewith dls.scope("scenes/0007"):shipped
F7 run handlenone — there is no pipeline run objectGAP — the open seam

Three of the original F5 gaps closed because the async generator became a first-class body kind, detected by inspect.isasyncgenfunction at decoration. Two cells remain open rather than shipped:

  • F3 × L3, the two-way .send() trampoline. The journal is one-way — workers append frames, consumers pull them — so there is no way to send a value into a paused remote body. Locally, plain Python generators do it fine.
  • F5 × L3, merging a growing set of streams. No fan-in verb exists; you hand-roll it with a TaskGroup and an asyncio.Queue. See F5 for the working, clunky version.

The SDK's own matrix marks both of those cells "explore" rather than GAP — they are reachable by hand today. Only the F7 row is marked GAP outright. See F7 — the run object.

The canonical, runnable version of this matrix lives in the SDK repo at docs/execution-patterns.md, with explorations examples/03_body_kinds.py through examples/06_dynamic_graph.py (PR #18).