DreamLake

Dataflow primitives

A draft in a different repo. None of this ships in the Python SDK.

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

design/primitives/examples/image_object_annotation.pypython
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.

PrimitiveMeaningBucket / hue
@udf / @udf(kind=…)compute node; kind sets its color everywheretransform #23aaff · model #7c5bd9
udf.local(...)bypass the trace, run the body in-process—
batch(src, n=)source yielding fixed-size batchessource #1f8f4a
stack([…])fan-in on a new leading axistransform #23aaff
x > t · & | ~ · .all/.mean(axis=) · x[mask]result algebra: predicates, set algebra, consensus, selectfilter #c0922e
to_dataset(x)data leaves the graphsink #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 · Stringshape- 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:

PipelineNodesEdges
image_object_annotation911
camera_pose_trajectory1317
video_task_annotation1729

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.

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.