# Dataflow primitives

> **Warning:** `import lakeshore`, `import dreamlake as dl`, `@dl.pipeline`, `batch`,
>   `stack`, `to_dataset`, `requeue`, `trace`, and `@udf(kind=…)` exist only
>   in the runnable draft at `design/primitives/` in
>   [`macrodata-subtask_yanbing`](https://github.com/fortyfive-ai/macrodata-subtask_yanbing/pull/12).
>   The shipped package is `dreamlake-lakeshore`, imported as
>   `import dreamlake.lakeshore as dls`, and its `udf()` accepts exactly
>   `queue=` and `transport=`. Nothing on this page is installable from PyPI.

Three contract pipelines (`image_object_annotation`,
`camera_pose_trajectory`, `video_task_annotation`) already *import* a
primitive API that did not exist. The draft makes it real: tracing a
pipeline yields a `Graph` whose nodes are colored by one shared six-hue
palette. The noun the author types is the node the tracer graphs and the hue
the canvas paints.

This is the pipeline the draft traces, close to verbatim:

```python file="design/primitives/examples/image_object_annotation.py"
import dreamlake as dl
import lakeshore as ls
from dreamlake import batch, requeue, to_dataset
from lakeshore.types import Mask, Tensor, Tuple

@ls.udf
def detect_objects(images: Tensor["N", "H", "W", 3]) -> Tuple["boxes", "classes", "confidence"]:
    ...

@ls.udf(kind="review")                      # review folds onto the filter bucket (amber)
def review_boxes(labels) -> Mask["R", "N"]:
    ...

@dl.pipeline
def image_object_annotation(source):
    for items in batch(source, n=64):
        labels = detect_objects(items.images)
        review = review_boxes(labels)

        consensus = review.all(axis=0)      # ∩ over reviewers
        good = labels[consensus]
        needs_attn = labels[~consensus]     # complement ¬

        yield to_dataset(good)              # sink
        requeue(needs_attn)                 # rework loop = plain Python
```

## The set

The draft splits into two namespaces: `lakeshore` is compute, `dreamlake` is
dataflow built on top of it.

| Primitive | Meaning | Bucket / hue |
| --- | --- | --- |
| `@udf` / `@udf(kind=…)` | compute node; `kind` sets its color everywhere | transform `#23aaff` · model `#7c5bd9` |
| `udf.local(...)` | bypass the trace, run the body in-process | — |
| `batch(src, n=)` | source yielding fixed-size batches | source `#1f8f4a` |
| `stack([…])` | fan-in on a new leading axis | transform `#23aaff` |
| `x > t` · `&` `\|` `~` · `.all/.mean(axis=)` · `x[mask]` | result algebra: predicates, set algebra, consensus, select | filter `#c0922e` |
| `to_dataset(x)` | data leaves the graph | sink `#c8513b` |
| `requeue(x)` | rework loop-back (dashed edge to the source) | idle `#9c907a` |
| `@pipeline` / `trace()` | driver; plain Python control flow *is* the DAG | — |
| `Tuple["a","b"]` · `Tensor` · `Mask` · `String` | shape- and column-carrying annotations | — |

Six hues, no seventh — that is the palette contract in
`uikit-workspace/staging/design.md`, which forbids a seventh hue outright.
Aliases fold by *function*, not by vibe: `review` → filter, `merge` → model,
`select` → filter, `requeue` → idle. The mapping lives in one table,
`KIND_COLOR`, with `KIND_ALIAS` beside it.

## What tracing yields

Under a trace, a `@udf` call does not run its body: it appends a node to the
active `Graph` and returns a `Col`, a lazy handle. Outside a trace the same
call runs the real body. All three contract pipelines trace to non-empty
graphs:

| Pipeline | Nodes | Edges |
| --- | --- | --- |
| `image_object_annotation` | 9 | 11 |
| `camera_pose_trajectory` | 13 | 17 |
| `video_task_annotation` | 17 | 29 |

The set algebra shows up as amber filter nodes, `requeue` as a dashed
loop-back into the source, and the renderer emits both 24-bit-color ASCII
(for the REPL) and Graphviz DOT (the canvas handoff). Both read color from
the same `KIND_COLOR` table, so Python, terminal, and canvas cannot drift.

```bash
python design/primitives/run_trace.py       # trace the 3 contract pipelines
python -m pytest design/primitives/tests    # the contract, locked (6 tests)
```

The six tests lock exactly the properties the design depends on: every
pipeline traces to a non-empty graph, every node kind resolves to a palette
bucket, source and sink are present and `requeue` loops back, `kind="review"`
folds onto `filter`, named columns flow on edges, and a `@udf` outside a
trace runs its body locally.

## Relationship to the shipped SDK

The primitives are the **authoring front-end**; the shipped SDK is the
**execution back-end**. A traced node would dispatch by submitting through a
queue (`inv = f.submit(...)`, returning a durable string id), and because
the shipped `udf()` signature takes only `queue=` and `transport=`, the
`kind=` keyword is free to carry the semantic bucket without colliding with
anything.

Free, but not wired: no shipped decorator accepts `kind=` today, and no
object joins a traced `Graph` to its live invocations. That object is the
missing [run object](/python-sdk/architecture/f7-run-object.md).

## Open questions

These come from the draft's own README, surfaced by actually running the
trace:

- **Do ML UDFs default to `model` (purple) or stay `transform` (blue)?**
  Today `detect_objects` and the caption models trace as blue, because the
  contract examples do not tag them.
- **Is `merge` a model?** `merge_captions` is described as the quotient
  operation but traces blue; `kind="merge"` would fold it onto purple.
- **Is mask-indexing a node or an edge?** `labels[mask]` currently makes an
  explicit amber `select` node. Nodes are visible and colorable; folding the
  operation into a styled edge keeps the graph leaner. The same call applies
  to `>`, `&`, and `~`.
- **Namespace split.** The draft assumes `lakeshore` is compute and
  `dreamlake` is dataflow, matching the org repos. The shipped client is
  `dreamlake.lakeshore` — promote `lakeshore` to top level, or re-export?
- **`stack` axis semantics.** `stack([...]).mean(axis=0)` assumes `stack`
  adds a leading model axis; the shape types (`Mask[3, N]`) should make that
  explicit.
