DreamLake

Sources and Sinks

Every pipeline has at least one root and at least one terminal. Neither is a special kind of object — both are ordinary UDFs, and the tracer tells them apart by their shape alone.

A root is a node nothing flows into. A terminal is a node nothing flows out of. In a lot of dataflow frameworks those are separate classes with their own base types, their own registration, and their own lifecycle hooks. Here they are not. A root and a terminal are both @ls.udf functions, and what makes one a root and the other a terminal is the signature — the parameter list and the return annotation.

That is why every rule below is a rule about Python syntax rather than a rule about an API. There is no Source class to subclass, no sink() builder, and no registry to add yourself to. The tracer reads the def line, and the graph follows from it.

Sources

A source is written @ls.udf(kind="source"). It takes no parameters, and it carries a return annotation:

python
@ls.udf(kind="source")
def load_clips() -> Tuple["clips", "prompts"]:
    """Read the work queue; one row per clip."""
    ...

The tracer derives two things from that signature. Input ports are the function's parameters, and there are none, so inputs is []. The return annotation is present, so the node gets its one output port and outputs is ['out']. The names inside Tuple[…] are not ports — they are the result's column schema, and they land in columns as ['clips', 'prompts']. A UDF returns exactly one table; passing it downstream passes the whole table.

The kind="source" argument does less than it looks like it does. It sets the card's kind dot, and it exempts the node from _infer_kinds, the pass that relabels un-annotated nodes as transform or merge after the walk. What actually makes load_clips a root is the empty parameter list: with no arguments there is no provenance to draw an edge from, so nothing can point at it. Delete the kind= and the node is still a root; add a parameter and it stops being one, whatever the label says.

Sinks

A sink is written @ls.udf(kind="sink"). It takes parameters, and it has no return annotation:

python
@ls.udf(kind="sink")
def save_dataset(rows):
    """Write the accepted rows to the dataset."""
    to_dataset(rows)

The mechanical rule, stated once so it is not mistaken for a convention: a node has an output port if and only if the function carries a return annotation. The tracer's line for it is exactly that literal — "outputs": ["out"] if fn.returns is not None else []. So kind="sink" labels the node, but it is the missing annotation that makes it terminal.

The two are genuinely independent, and the consequence is worth being blunt about. Write @ls.udf(kind="sink") on a function that does return something and the node gets an output port anyway — downstream UDFs will connect to it and the graph will render those edges without complaint. The label paints the card; it does not enforce the shape. Drop the annotation on a function you never labelled and it becomes terminal regardless.

to_dataset and requeue

These are the two framework terminals: to_dataset(rows) writes rows into a dataset, and requeue(rows) sends them back around for another pass. Both must be wrapped inside a kind="sink" UDF. Calling either one directly in the pipeline body raises at trace time.

This is the rule that trips everyone, so here is why it exists — it is not a special case for these two names. Every call inside a @dl.pipeline body must resolve to a decorated UDF. When the walk reaches a call it cannot attribute to an @ls.udf / @dl.node definition, it fails loudly rather than passing an unknown call through:

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

to_dataset and requeue are imported from the framework, not defined as UDFs in your file, so they hit that check like any other bare helper would. Wrapping them is the whole fix, and it costs two lines each:

python
@ls.udf(kind="sink")
def save_dataset(rows):
    to_dataset(rows)          # writes the accepted rows


@ls.udf(kind="sink")
def rework(rows):
    requeue(rows)             # sends the rejected rows back for another pass
clips
rows
rows
load_clips
source · 0→1
annotate
transform · 2→1
save_dataset
sink · 1→0
rework
sink · 1→0
a source with no inputs, and two sinks with no outputs

Look at what is not in that graph. There is no to_dataset node and no requeue node. The tracer never descends into a UDF body — it reads the def line and stores the whole function's text verbatim in the node's code field. The sink is the node; the framework call inside it is body text. to_dataset(rows) is in save_dataset's code string, where a source viewer can show it — but it is never a node of its own, and no edge ever points at it.

The rest of the shape reads straight off the signatures. load_clips has no inbound edge and no input dots, because it has no parameters. save_dataset and rework have no outbound edge and no output dot, because neither carries a return annotation. Nothing about the picture required a special node type.

Two sinks is the normal shape

One terminal is the exception, not the rule. A pipeline that decides anything usually has two: one for the rows that passed and one for the rows that did not.

python
@dl.pipeline
def image_object_annotation():
    src = load_clips()
    for items in batch(src, n=64):
        labels = annotate(items.clips, items.prompts)
        consensus = review_boxes(labels)
        save_dataset(labels[consensus])   # accepted
        rework(labels[~consensus])        # rejected → another pass

The figure above is that pipeline with the review step elided, and its edges drawn solid. In the full pipeline they would be dashed. labels[consensus] puts a mask in a subscript's slice, which taints that argument's whole provenance to mask, and a mask edge renders dashed — the visual claim being that review_boxes does not contribute data to the sink, it only gates which rows reach it. That gating is the subject of Review and masks, and the mask subscript as a way of tapping rows off the main stream is the subject of Sampling. This page keeps its edges solid on purpose: the terminals themselves are the point here, ungated.

Note also that rework closes a loop conceptually — rejected rows go back to the queue load_clips reads — but it does not close one in the graph. The tracer is static and the queue is outside it, so the round trip is invisible. A traced pipeline graph is always acyclic.

Proposal — not importable yet

to_dataset, requeue, batch, stream, pipeline, and node are not exported by the shipped dreamlake package, and lakeshore.udf accepts no kind= argument — its real signature is udf(fn=None, *, queue=None, transport="auto", config=None). kind is a convention the tracer reads and the runtime never sees. These files are a specification format that happens to be Python syntax: they trace because the tracer parses them and never imports or executes anything.