# 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.

> **Warning:** 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
`` component from `@dreamlake/uikit`, not a picture of one:
drag a card, scroll to pan, hold ⌘ or ctrl and scroll to zoom.

## 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:

| Kind | How it arises | What it changes |
| --- | --- | --- |
| `source` | `@ls.udf(kind="source")` — no parameters, has a return annotation | `inputs: []`, `outputs: ['out']`. Its outgoing edges are excluded from merge inference |
| `sink` | `@ls.udf(kind="sink")` — has parameters, **no** return annotation | `outputs: []`. 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 |
| `merge` | **Inferred** when a plain UDF has two or more non-source `data` inputs, or set explicitly | Purely cosmetic — fan-in is already visible in the edges |
| `transform` | The default for a plain UDF with fewer than two stage inputs | Nothing |
| `model` | **Renderer-only.** The tracer never emits it today | — |
| `filter` | **Renderer-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.

> **Note:** `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.

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:

```text
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](/pipelines/batching.md) 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
```

> **Warning:** `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.

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

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

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

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

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

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

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

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

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