DreamLake

F7 — the run object (gap)

Row F7 of the execution matrix is empty. There is no run = pipeline.submit(source) handle you can poll, stream, or draw. The SDK's own docs/execution-patterns.md marks the entire row GAP.

Design note, not API

The "what exists today" sections below are verified, runnable code. The PipelineRun sketch further down is a surface proposal: no PipelineRun class, no run.status, no run.events, and no DAG object exists in dreamlake-lakeshore. Do not write code against it.

The gap, precisely

A pipeline today is ordinary Python driving submit() calls. When it runs, exactly two kinds of run record come into existence:

  1. Invocation ids you keep by hand. Every submit() returns a durable string id. If you want the graph the run took, you append (src_inv, dst_inv, verb) tuples to a list yourself.
  2. Journal frames. Every generator invocation writes VALUE frames and a terminal END or ERROR frame to its journal; Topic is the same journal spelled pub/sub. That is live, durable, cursor-addressable status — per invocation.

Nothing joins them. The structure (which node feeds which) lives in your hand-kept edge list or in an author-time trace; the status (queued, running, done) lives in the journal, keyed by invocation id. No object holds both, so there is nothing to ask "how far along is the run?" and nothing to hand a canvas to draw live.

There is one unwired seam worth naming: the plane already has a parentage column. Dispatch.submit and HttpDispatch.submit both accept parent=, and wire.Envelope carries a parent field documented as "parent invocation id (pipeline lineage)". Nothing in the SDK populates it — UDF.submit passes only id, key, and priority. The storage for a run tree exists; the writer does not.

What exists today: the hand-kept edge list

examples/06_dynamic_graph.py already produces a full run record — by hand. Every edge is a pair of invocation ids plus a verb. Compact version, runnable as-is:

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

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

@dls.udf(queue="run-demo")
def propose(seed: int) -> list[int]:
    return list(range(seed, seed + (seed % 3) + 1))

@dls.udf(queue="run-demo")
def score(x: int) -> float:
    return (x * 37 % 10) / 10

def run(inv: str):
    dls.run_worker(q, once=True)
    return q.result(inv, timeout=5)

edges: list[tuple[str, str, str]] = []      # (src_inv, dst_inv, verb)

p_inv = propose.submit(2, _dispatch=plane)
candidates = run(p_inv)                     # width discovered at run time

accepted = []
for x in candidates:
    s_inv = score.submit(x, _dispatch=plane)
    edges.append((p_inv, s_inv, "score"))   # the edge, kept BY HAND
    if run(s_inv) >= 0.5:
        accepted.append(x)

assert len(edges) == len(candidates)        # the run's actual graph

The full example adds a rework loop — a failed candidate goes through refine and around again, up to three attempts — and prints the edge tally. A three-seed run produces 18 edges (12 score, 6 refine), a shape the author could not have known. The edge list is the dynamic graph. It just lives in a local variable that dies with the process, while the ids inside it are durable.

What exists today: the journal as live status

The other half already streams. A generator UDF's yields become journal frames any consumer can pull with a cursor, and Topic is the same mechanism for ad-hoc events (see examples/05_streams_topics.py). Runnable as-is, continuing the setup above:

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

inv = render_frames.submit(4, _dispatch=plane)
dls.run_worker(q, once=True)

frames = list(q.stream(inv, timeout=5))          # live per-invocation feed
assert [f["frame"] for f in frames] == [0, 1, 2, 3]

tail = list(q.stream(inv, since=1, timeout=5))   # reconnect from a cursor
assert [f["frame"] for f in tail] == [2, 3]

t = q.topic("progress")                          # same journal, pub/sub
t.post({"node": "score", "state": "running"})
t.post({"node": "score", "state": "ok"})
t.post(None, final=True)
assert [m["state"] for m in t.stream(since=-1, timeout=5)] == ["running", "ok"]

Everything a run object needs for liveness is here: durable frames, a 0-based since= cursor, reconnect from any process. What is missing is the map from journal frames back to nodes in a structure.

What exists today: the traced graph (author time)

The dataflow-primitives draft supplies the third piece — structure without execution. It lives at design/primitives/ in fortyfive-ai/macrodata-subtask_yanbing, not in the SDK. There, a @udf call under an active trace does not run its body: it appends a node to the active Graph and returns a Col, a lazy handle. The result algebra extends the graph — x > 0.5, a & ~b, review.all(axis=0), and mask-indexing labels[mask] each become filter nodes; batch(...) is a source node, to_dataset(...) a sink, and requeue(...) a dashed loop-back edge. @udf(kind=...) sets the node's bucket, and each bucket owns one hue of the shared six-hue palette: source green, transform blue, model purple, filter amber, sink red, idle grey — the same colors the canvas renders.

So the traced Graph knows every node and edge of the pipeline before anything runs. It has no invocation ids and no status: it is the exact complement of the journal.

The PipelineRun sketch

The run object is the join. pipe.submit(source) would return a handle holding three things: the traced structure, a node-to-invocation-id map, and a journal cursor.

python
# SKETCH — not shipped. Surface proposal only; names and shapes will move.

run = image_object_annotation.submit(source)     # -> PipelineRun

run.id                    # durable string, like any invocation id
run.graph                 # the traced Graph: nodes, edges, kinds
run.ids                   # {node -> invocation id | None (not submitted yet)}

# Per-node status, colored by the same six-hue palette the graph uses:
run.status                # {node -> "idle" | "queued" | "running" | "ok" | "error"}
run.status["score"]       # one node

run.result("detect_objects")        # q.result on the node's invocation(s)

# Live events ride the journal, cursor and all:
for ev in run.events(since=cursor):          # sync
    ...                                      # (seq, node, frame)
async for ev in run.events(since=cursor):    # async, same shape
    ...

run.graph.to_dot()        # render: DOT out, canvas draws it live
                          # (nodes colored by kind, tinted by status)

# Ids are durable, so reattaching from another process is natural:
run = PipelineRun.attach(run_id)
for ev in run.events(since=run.cursor):      # resume where you left off
    ...

Nothing in the sketch invents new machinery. run.ids is the hand-kept dict from the dynamic-graph example, kept by the system instead; run.events is q.stream with node labels attached; run.graph is the trace the primitives draft already produces; attach is the same reconnect story q.result(inv) and since= cursors already have for single invocations.

Design questions

  • Trace-to-id binding. The traced graph exists at author time; ids exist at submit time. When a node fans out — one score node, N invocations — is run.ids["score"] a list, and who appends to it: the driver, or the plane?
  • Dynamic edges. Rework edges are discovered mid-run. Does the run object grow its graph live, with edges arriving as events, or does run.graph stay the static trace with rework shown only in status?
  • One journal or many. Does run.events() merge the per-invocation journals, and if so how are they ordered across journals? Or does the run own a single run-level Topic that workers post to alongside their own journals?
  • Is the run itself an invocation? If pipe.submit() returns a durable id and writes journal frames, the pipeline is just an invocation whose body submits others — attach, result, and stream come for free, parent gets its writer, and nested pipelines recurse. That economy is attractive; the open question is whether driver-side control flow can run on a worker without occupying one for the whole run.