DreamLake

Pipeline Functions

A DreamLake pipeline is a Python module. A static tracer reads its source and derives a typed DAG from it — it never imports the module and never executes a line of it. That is why every UDF body in every example on these pages is ... and traces perfectly well. The graph is the artifact; there is no run.

The tracer parses the file, treats each @ls.udf function as a node template, walks the body of the one @dl.pipeline function statement by statement, and tracks the provenance of every local variable — the set of nodes its data descends from. Calling a UDF creates a node and an edge from each argument's provenance. Method calls, subscripts and operators create no nodes at all; they pass provenance through. Everything the graph knows, it knows from syntax.

Proposal — not importable yet

The authoring API described in this section is not implemented in shipped Python. The dreamlake package exports no pipeline, no node, no batch, no stream, no to_dataset and no requeue. The real signature of lakeshore.udf is udf(fn=None, *, queue=None, transport="auto", config=None) — it takes no kind= argument. kind is a convention the tracer reads out of the decorator's keyword arguments, and it never reaches the runtime. lakeshore.types — Mask, String, Tensor, Tuple — does not exist, and the canonical import for the real package is import dreamlake.lakeshore as dls, not import lakeshore as ls.

The honest conclusion: these .py files are a specification format that happens to be Python syntax. They trace because the tracer never imports or executes anything. Read every page in this section as the shape the graph takes, not as code you can run today.

Here is the smallest complete pipeline — a source, a transform and a sink — exactly as the tracer lays it out. The figure below is the real <PipelineGraph> component from @dreamlake/uikit, not a picture of one: drag a card, scroll to pan, hold ⌘ or ctrl and scroll to zoom.

clips
rows
load_clips
source · 0→1
caption
transform · 1→1
save_captions
sink · 1→0
source → transform → sink, as the tracer derives it

The node and edge schema

The tracer emits JSON, the server stores it on a PipelineVersion, and the renderer consumes it. These are the server-side wire types, verbatim:

ts
interface NodeDef {
  id: string
  title: string
  kind: string
  inputs: string[]
  outputs: string[]
  columns: string[]
  code: string
  config: Record<string, unknown>
  pos: { x: number; y: number }
}

interface EdgeDef {
  from: string
  fromPort: string
  to: string
  toPort: string
  kind: 'data' | 'mask'
}

interface GraphObject {
  id: string
  title: string
  subtitle: string
  nodeCount: number
  nodes: Record<string, NodeDef>
  edges: EdgeDef[]
}

Read the fields against the Python they come from. inputs are the UDF's parameters, one port each. outputs is a single port, ['out'], because a UDF returns one table and passing that table downstream passes the whole thing — a sink, which returns nothing, has [] instead. columns is the schema of that one output, taken from the return annotation, and is emphatically not a list of ports. config is the decorator's keyword arguments verbatim, which is where kind lives. pos comes from a longest-path layering the tracer runs after the walk, at a pitch of 208 px per layer and 108 px per row.

Two differences separate GraphObject from the PipelineGraphData the renderer wants, and both are bridged client-side in toGraphData():

  • The server graph has no top-level code. It is filled in from the pipeline's sourceCode field, which is stored beside the graph rather than inside it.
  • Server NodeDef has no status. Runtime state lives in a separate NodeState record, so the static view is entirely idle until a runtime overlay is wired in. Every figure in this section is therefore idle, and that is not a stylistic choice — it is what a traced graph actually contains.

Node kinds

kind is cosmetic in the renderer — it drives the card's category dot — but it is load-bearing in the tracer, because two of the kinds change what the node is. Not every kind the renderer paints is a kind the tracer can produce:

KindHow it arisesWhat it changes
source@ls.udf(kind="source") — no parameters, has a return annotationinputs: [], outputs: ['out']. Its outgoing edges are excluded from merge inference
sink@ls.udf(kind="sink") — has parameters, no return annotationoutputs: []. The graph's terminal
review@ls.udf(kind="review") returning Mask[…]Nothing structurally, but the result is used as a selector, so its outgoing edge is a dashed mask edge
mergeInferred when a plain UDF has two or more non-source data inputs, or set explicitlyPurely cosmetic — fan-in is already visible in the edges
transformThe default for a plain UDF with fewer than two stage inputsNothing
modelRenderer-only. The tracer never emits it today—
filterRenderer-only. The tracer never emits it today—

Merge inference is narrower than it first looks. Only data edges count toward the fan-in, and inputs drawn from a source node are excluded — a side input such as a prompt table should not turn an ordinary stage into a merge. An explicit kind= on the decorator always wins over inference.

A known inconsistency

dreamlake-server's own copy of the node type narrows kind to 'source' | 'transform' | 'sink', dropping review and merge — both of which the tracer emits and the renderer paints. The values survive the round trip because the field is serialised as JSON, but the server's type is wrong about its own data. Widening it is recorded as a hand-off.

Two edge kinds

An edge carries no runtime style — its flow animation is derived at render time from the status of its two endpoint nodes — but it does carry one static tag, and the tag means something precise:

  • data (solid) — the value flows through. The target receives the source's table.
  • mask (dashed) — the source only gates or filters the target. It is a review veto, a confidence threshold, a boolean selector. No table crosses this edge; a decision does.
rows
rows
detect_objects
transform · 1→1
save_boxes
sink · 1→0
review_boxes
review · 1→1
keep_agreed
sink · 1→0
above: a solid data edge. below: a dashed mask edge

A mask tag is produced by position, not by declaration. When the tracer evaluates a subscript such as labels[consensus], it evaluates the slice with the mask role, so everything consensus descends from reaches the target tagged mask while labels reaches it tagged data. That single rule is what makes review and consensus edges render as dashed gates without anyone annotating them.

The promotion rule matters when both tags reach the same edge, which happens whenever one node feeds another twice — once as a value and once as a selector. data dominates mask. The edge is drawn solid, because a table really is crossing it, and the gating relationship is already implied. There is exactly one edge per (from, fromPort, to, toPort) key; the second contribution merges into the first rather than drawing a parallel line.

Trust the pipeline; fail loudly

Any undecorated function call inside a @dl.pipeline body raises at trace time. Not a warning, not a passthrough node, not a silent skip — the trace fails and the CLI exits non-zero:

pipeline calls undecorated function 'normalize' — every function called in a
@dl.pipeline must be a @ls.udf / @dl.node (wrap framework helpers like
to_dataset/requeue in a udf)

This is a deliberate design choice and it is the reason the graph can be trusted. If unknown calls passed through quietly, the rendered DAG would be a partial picture of the program, and the one thing worse than no graph is a graph that silently omits a stage. So there are no hidden framework helpers: to_dataset and requeue do not get a special case, they get wrapped in a kind="sink" UDF like everything else.

There is exactly one exemption. A for-loop's iterator — batch(src, n=64) or stream(...) — is evaluated leniently: it adds no node, and it passes the iterated source's provenance straight to the loop variable. Structural iteration is not a stage, so it is not required to be a UDF. The loop body is strict again, and UDF calls nested inside it become nodes exactly as they would outside. See Batching for what the exemption does and does not cover.

How to trace one

Locally, against the source file, with no server involved:

bash
python -m dl_trace pipelines/image_object_annotation.py --pretty

To push it, so the server traces it and stores the graph on a version:

bash
dreamlake pipeline create <name> --file pipelines/image_object_annotation.py
The server-side parser is a stub by default

dreamlake pipeline create does not run dl_trace. The server dispatches to an external parser at PIPELINE_PARSER_URL. If that variable is unset it falls back to a regex mock that matches only a bare @ls.udf line — not @ls.udf(kind="source"), not @dl.node — and returns inputs: [] and edges: []. The result is an empty graph that looks like a successful push. Real graphs require the external dl_trace service to be configured.

The nine constructs

One page per construct, each with its own live node view.

@dl.pipeline →

The one entry point. What the decorator marks, why there is exactly one per module, and what the walk does with the body.

@ls.udf →

Ports from parameters, columns from the return annotation, and the decorator kwargs that end up in config.

Sources and Sinks →

Where data enters and leaves — kind="source", kind="sink", and the wrapped to_dataset / requeue helpers.

Transform and Merge →

The default kind, and the fan-in rule that infers merge at two or more non-source data inputs.

Review and Masks →

Gates, Mask[…] returns, and the mask algebra ~ & | ^ as edge tags rather than nodes.

Sampling →

Tapping data off the main stream with a mask subscript — and where the real statistical samplers actually live.

Batching →

batch() and stream() — the single exemption from the decoration rule, and its exact boundary.

Control Flow →

if / else, for, comprehensions — what the tracer keeps, and the two named fidelity losses.

Model Functions →

The canonical model API — a contract the pipelines are written against, not working code.