DreamLake

Pipeline Graph

PipelineGraph draws a DreamLake pipeline as an interactive flow-chart. It's purely presentational: you hand it one JSON object and it renders the DAG — dotted canvas, status-tinted node cards, orthogonal edges, live flow animation — with no data fetching of its own. PipelineSource is the paired read-only source inspector.

If you just want to render a graph, this page is the whole story. For the exact data shape see Pipeline Graph JSON; for a field-by-field visual map see Anatomy; for how the component is built and what's coming, see Architecture & Roadmap.

Background — what this renders and why

DreamLake runs an AI auto-labeling + human-review loop: models propose labels, reviewers gate them, accepted rows land in a dataset and the rest go back for rework. That loop is written as a pipeline — a plain Python @dl.pipeline function whose stages are @ls.udf functions.

You never draw this graph by hand. The dl_trace tracer reads the pipeline's Python statically (no import, no execution) and derives the node/edge graph from its dataflow. So:

  • The graph is a derived view of code. Change the .py, re-trace, the graph updates. There is no separate diagram to keep in sync.
  • Placeholder bodies (...) trace like real ones, so a graph renders long before any stage is implemented.
  • This component only draws. Tracing happens elsewhere (a Python service); the runtime status that animates the graph streams in separately.

Who this is for: anyone rendering a traced pipeline in a UI — a Studio view, a dashboard, a docs page. You supply the traced JSON (and optionally a live status overlay); the component owns the canvas.

One stage, three faces

The thing that makes the data model click: a node is one stage seen three ways. You write a Python UDF; the tracer emits a JSON node; the component draws a card. They are the same object.

① The Python you write — a @ls.udf stage. Its parameters become input ports; its return columns become the result schema:

python
@ls.udf
def detect_objects(images: Tensor["N", "H", "W", 3]) -> Tuple["boxes", "classes", "confidence"]:
    """Detect objects in each image; one row per box."""
    ...

② The node you see, ③ the JSON in between — flip between the rendered card (Preview), the render call (Source), and the tracer's node JSON (Data). The card is the JSON: the kind glyph ← kind, the title ← title, the 1→1 meta ← inputs/outputs lengths, and a single input / output dot on each edge (every parameter shares the one input dot).

detect_objects
transform · 1→1

The return columns (boxes, classes, confidence) live in columns — the result's schema, not extra ports. A UDF returns one table; passing it downstream passes the whole table. See Anatomy for the full field-to-pixel map.

Basic

A freshly-traced graph is entirely idle. Drag nodes to rearrange, scroll (or two-finger drag) to pan, ⌘/ctrl-scroll or pinch to zoom, click a node to select it. Once the canvas is focused, arrow keys walk the selection — ↑ / ↓ step through the pipeline in topological order, ← / → jump to the upstream / downstream neighbour, Esc clears — and the selected node pans into view. Selecting a node tints its card border and the edges touching it in that node's status colour — a running selection reads blue, an errored one red, an untouched idle one stays neutral — while the unrelated edges fade back; the other node cards are never dimmed. Edges come in two kinds — data (solid) and mask (dashed gate); the Anatomy page shows both.

videos
frames
keypointsdescriptors
poses
rows
rows
rows
rows
load_videos
source · 0→1
extract_frames
transform · 1→1
detect_features
transform · 1→1
estimate_poses
transform · 2→1
bundle_adjust
transform · 1→1
save_dataset
sink · 1→0
rework
sink · 1→0

Graph + source, linked

The design layout is a canvas with a source right rail. The two share one selection: click a node and the rail jumps to it; click the background to clear. Both components are controlled — you own the selectedNodeId state.

The seam between them is draggable — grab it and pull to trade canvas width for code width (the rail keeps a 180px floor, the canvas 200px). Neither component owns that behaviour: the layout is the caller's, and the specimen below builds it from ResizeDivider, the same drag primitive ResizableLayout uses. The rail holds a pixel width while the canvas is flex: 1 1 0, so the canvas absorbs viewport changes and the code column stays put.

PipelineSource is more than a code viewer — its tabs are contextual:

  • Nothing selected → a single PIPELINE tab: the pipeline's status (its title, a status-count grid over all six node states, and a nodes · edges line) above the whole pipeline.py in a read-only, line-numbered, Python- highlighted editor.
  • A node selected → NODE (its live execution status, the default) and CODE (just that stage's source), plus PIPELINE to step back out (which clears the selection). The NODE tab is the inspector: a status pill (colour + label, and — once a run supplies them — progress, duration, rows), the node's resolved i/o (each input port traced to its upstream node, each output to its downstream consumers), the result schema (the columns), the decorator config, and an output preview — a sampled result table once the node is ok, or a status-keyed one-liner before then.
videos
frames
keypointsdescriptors
poses
rows
rows
rows
rows
load_videos
source · 0→1
extract_frames
transform · 1→1
detect_features
transform · 1→1
estimate_poses
transform · 2→1
bundle_adjust
transform · 1→1
save_dataset
sink · 1→0
rework
sink · 1→0
camera_pose_trajectory
7idle0running0waiting0ok0error0stale
7 nodes · 9 edges
"""Recover the camera pose trajectory from a video.
frames → features → pose estimation → bundle adjustment; a confidence
mask (σ-algebra: intersection of two boolean columns) decides what goes
to the dataset vs. rework. EVERY function the pipeline calls is a
`@ls.udf` — the source and the two sinks included. UDF bodies are
placeholders.
"""
import dreamlake as dl
import lakeshore as ls
from dreamlake import batch, requeue, to_dataset
from lakeshore.types import Tensor, Tuple
@ls.udf(kind="source")
def load_videos() -> Tuple["videos"]:
"""Pull the batch of videos to process."""
...
@ls.udf
def extract_frames(videos: Tensor["N"]) -> Tuple["frames", "timestamps"]:
"""Decode each video into a frame column plus per-frame timestamps."""
...
@ls.udf
def detect_features(frames: Tensor["N", "H", "W", 3]) -> Tuple["keypoints", "descriptors"]:
"""Per-frame keypoints and descriptors (e.g. SuperPoint)."""
...
@ls.udf
def estimate_poses(keypoints, descriptors) -> Tuple["poses", "confidence"]:
"""Visual odometry / SfM placeholder; poses are 4x4 world-from-camera."""
...
@ls.udf
def bundle_adjust(poses: Tensor["N", 4, 4]) -> Tuple["poses", "residual"]:
"""Global refinement; residual is the reprojection error per frame."""
...
@ls.udf(kind="sink")
def save_dataset(rows):
"""Write the trajectory to the dataset (wraps dreamlake.to_dataset)."""
to_dataset(rows)
@ls.udf(kind="sink")
def rework(rows):
"""Send low-confidence frames back (wraps dreamlake.requeue)."""
requeue(rows)
@dl.pipeline
def camera_pose_trajectory():
src = load_videos()
for items in batch(src, n=8): # batch = chunked/streamed run
frames = extract_frames(items.videos)
feats = detect_features(frames.frames)
poses = estimate_poses(feats.keypoints, feats.descriptors)
traj = bundle_adjust(poses.poses)
ok = (poses.confidence > 0.5) & (traj.residual < 1.0) # A ∩ B
save_dataset(traj[ok]) # sink: the trajectory
rework(traj[~ok]) # low-confidence frames go back

Running from the rail

The rail is also where a run is driven and gated — all UI-only, so the actual execution is the caller's. Wire onRun and a ▶ RUN button appears on the PIPELINE tab (targeting the whole pipeline) and on a review node's NODE tab (targeting that node, disabled until every upstream is ok). The other props switch what that button shows, so the rail mirrors the run's lifecycle:

  • running → the button becomes a disabled running… with a spinner.
  • reviewNodeId → it becomes a purple ⏸ REVIEW button that selects the paused node so a human can inspect it (takes precedence over running).
  • done → a green ✓ DONE marker in its place (click to run again).
  • onContinue → on a review node that is waiting, a green ▶ CONTINUE button releases the pipeline to run downstream (also gated on upstreams).

None of these props run anything themselves — the host owns the run driver and feeds progress back through statusById. The live-status demo below is one such driver (a tiny in-browser simulation) animating the graph through statusById.

Live status — a runnable pipeline

Edges carry no stored style. Each edge's visual flow is derived from the status of its two endpoint nodes, so animating a running pipeline is just a matter of feeding node statuses in via the statusById overlay. The six flow states, and how each is derived, are catalogued in Anatomy → Connector states.

The demo below is a tiny in-browser "runner". Each node, when it finishes, settles to a random outcome — mostly ok, occasionally stale or error. A node only starts once all its upstreams have settled and none errored, so a failure blocks everything downstream (those nodes never run and stay idle), exactly like a real pipeline. Nothing real executes — it only drives statusById. Press Run (each run re-rolls): running edges flow blue, ok edges go green, a stale node's edges turn amber, an error node's edges turn red, and nodes blocked by an upstream failure stay idle.

random outcomes — a failure blocks everything downstream0 ok · 0 failed · 7 blocked
videos
frames
keypointsdescriptors
poses
rows
rows
rows
rows
load_videos
source · 0→1
extract_frames
transform · 1→1
detect_features
transform · 1→1
estimate_poses
transform · 2→1
bundle_adjust
transform · 1→1
save_dataset
sink · 1→0
rework
sink · 1→0

Ports & param tags

A stage's signature still sets its port counts — the card's N→1 meta reads straight off the UDF, from a source with zero inputs to a merge gathering four. But every parameter shares one input dot on the left edge (and the result shares one output dot on the right); the card never fans a dot out per parameter. The parameter names live in a floating param tag instead: one tag per node-pair, listing every param that pair transfers (each with a small leading dot, stacked), tinted to match the edge's flow.

Each tag rests on its edge (at the bend), tracking it as nodes move. They're also draggable: drag one along its edge to rebend the edge (the orthogonal jog follows), or across it to lift the tag onto a dashed leader (release near the line to snap back). Press-and-hold highlights the tag and its edge — the same treatment a node selection gives, so you can trace a crowded corner.

rows
left
a
w
rows
right
ingest
source · 0→1
map1
transform · 1→1
zip
transform · 2→1
blend
merge · 3→1
collate
merge · 4→1
sink
sink · 1→0

The same graph paired with its source rail — click a node to read the UDF whose parameters and return columns produced those ports:

rows
left
a
w
rows
right
ingest
source · 0→1
map1
transform · 1→1
zip
transform · 2→1
blend
merge · 3→1
collate
merge · 4→1
sink
sink · 1→0
port_showcase
6idle0running0waiting0ok0error0stale
6 nodes · 6 edges
"""Nodes with 0–4 ports.
A hand-built showcase graph: each stage's parameters become input ports
and its result is the single output port, spanning 0–4 inputs. One edge
loops backward to exercise the centered S-bend.
"""

Props

PipelineGraph

PropTypeDefaultDescription
graphPipelineGraphData—The traced graph JSON.
statusByIdStatusOverlay—Live per-node { status, progress, ... }, merged onto the graph.
selectedNodeIdstring | null—Controlled selection. Omit for uncontrolled.
onSelectNode(id: string | null) => void—Selection change (also fires on background click).
showControlsbooleantrueShow the overlay chrome — the legend (top-right: node-kind glyphs + edge-flow swatches) and the keyboard-hint strip (bottom). Pass false for tiny embeds.
classNamestring—Extra classes on the canvas.

PipelineSource

PropTypeDefaultDescription
graphPipelineGraphData—The traced graph JSON.
selectedNodeIdstring | null—Which node to inspect (else the PIPELINE tab: status + the whole .py).
onSelectNode(id: string | null) => void—Fired by the PIPELINE tab (to clear) and the REVIEW button (to select the paused node).
statusByIdStatusOverlay—Live per-node status, merged over each node's static status; drives the status panels, i/o, and the output preview.
onRun(target: RunTarget) => void—Enables the ▶ RUN button. RunTarget is { kind: 'pipeline' } (PIPELINE tab) or { kind: 'node'; id } (a review node's NODE tab).
onContinue(nodeId: string) => void—Enables ▶ CONTINUE on a waiting review node, releasing the pipeline downstream.
runningboolean—A run is executing → the RUN button goes to a disabled running… state.
reviewNodeIdstring | null—A run paused on this review node → RUN becomes a ⏸ REVIEW button that selects it (takes precedence over running).
doneboolean—The run finished → a ✓ DONE marker shows in place of RUN (click to run again).
classNamestring—Extra classes.

Both are theme-aware (uikit tone tokens) and load @dreamlake/uikit/styles.css for their colours. RunTarget and StatusOverlay are exported from @dreamlake/uikit.


Next: Anatomy maps every field to a pixel · Pipeline Graph JSON is the data-model reference · Architecture & Roadmap covers the internals and what's next.