Skip to content

Go SDK: order tasks with Before, After and Label - #74119

Merged
henry3260 merged 6 commits into
apache:mainfrom
henry3260:go-sdk-node-before-after-label
Oct 3, 2026
Merged

henry3260 merged 6 commits into
apache:mainfrom
henry3260:go-sdk-node-before-after-label

Conversation

@henry3260

@henry3260 henry3260 commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

Why

A Dag authored in Go could declare its tasks and the data edges between them
(airflow.Inputs), but not an ordering between tasks that exchange no data — Python's
a >> b and b << a — and no way to label an edge for the graph view.

This implements decisions 6, 7 and 8 of
go-sdk/adr/0008-native-dag-interface.md:
airflow.Node is the Go counterpart of Python's DAGNode/DependencyMixin.

loaded.Before(notified, cleaned)  // loaded >> [notify, cleanup]
cleaned.After(extracted)          // cleanup << extracted

loaded.Before(airflow.Label(emptyNotice, "when empty"))

What

go-sdk/airflow/node.go (new):

  • Node, sealed by an unexported node(), satisfied by *TaskRef.
  • Before/After, variadic so one call fans out. Both return their argument set as one
    Node, not the receiver, so a.Before(b, c).Before(d) is a >> [b, c] >> d.
  • Label(node, text) wraps the endpoint, so each edge of a fan-out can carry its own label.
    A later label on an edge replaces it, as Python's DAG.set_edge_info does. A label on a
    verb's receiver has no edge to land on and is dropped, as Label("x") >> b is in Python.
  • Declaring an edge the Dag already has only applies the label, which is how an
    airflow.Inputs edge gets one: extracted.Before(Label(transformed, "rows")).

go-sdk/airflow/dag.go: DagRef.edgeLabels holds every edge and its label;
TaskRef.upstreams/downstreams hold both sides. DagRef.Task records the edges Inputs
declares into the same structure. markRegistered rejects a Dag whose dependencies contain a
cycle, next to the existing check for a condition without a Then: the Dag is whole by then,
so one walk of the graph answers for every edge, and the Dag stays unregistered so it can be
corrected. Declaring an edge never walks the graph, so building a Dag is linear in its edges.


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Code (Opus 5)

@henry3260
henry3260 requested a review from guan404ming October 2, 2026 18:13
@henry3260
henry3260 marked this pull request as draft October 2, 2026 19:06
@henry3260
henry3260 marked this pull request as ready for review October 2, 2026 19:54

@FrankYang0529 FrankYang0529 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Overall LGTM. Thanks for the PR.

Comment thread go-sdk/airflow/dag.go

@jason810496 jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice, thanks! There're the claude review findings.

Comment thread go-sdk/airflow/node.go
Comment thread go-sdk/airflow/node.go
Comment thread go-sdk/airflow/node.go Outdated
Comment thread go-sdk/airflow/node.go Outdated
Comment thread go-sdk/airflow/node.go Outdated

@jason810496 jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

My one design decision concern before addressing the claude review:

Would it be better to validate the graph at once "after the Dag be registered" so that we don't need to traverse whole graph again whenever a new edge be added?

Just like the TS side that we will only validate the cycle at the "finalize stage". I think we should do the same "last mile validation right before the serve" instead of having iterative validation (since we will compile the SDK into artifact anyway)

A Dag authored in Go can now order tasks that exchange no data, the way
Python's >> and << do, and label the edge that reaches a task.

    loaded.Before(notified, cleaned)  // loaded >> [notify, cleanup]
    cleaned.After(extracted)          // cleanup << extracted

    loaded.Before(airflow.Label(emptyNotice, "when empty"))

Both verbs are variadic, so one call fans out, and both return their
argument set as one Node, so a.Before(b, c).Before(d) is a >> [b, c] >> d.

airflow.Label wraps the endpoint rather than the call, so each edge of a
fan-out can carry a label of its own. Declaring an edge the Dag already
has only applies the label, which is how an airflow.Inputs edge gets one.
Before and After with no node declare no edge, but they now still reject a
Dag that has been registered and a *TaskRef that DagRef.Task did not
return, so a mistake in the receiver is not hidden by an empty fan-out.

The label a verb puts on an edge is now settled in one place, mergeLabel,
which both the check over the call's pairs and the recording step use. An
edge verb enumerates its pairs once and records what it checked, so the
checks and the recording cannot drift apart, and addEdgeLocked keeps only
the cycle check, which cannot run up front.

An edge given two labels in one call now says so rather than reporting the
first as already on the edge, which no earlier declaration had put there.
Declaring a label on an edge that already carries one now replaces it,
which is what Python's DAG.set_edge_info does: "this will overwrite,
rather than merge with, existing info". It used to panic, a rule the Go
SDK had of its own.

A label on the receiver of an edge verb is still dropped, since Label
marks the edge that reaches a node and the receiver is the node an edge
leaves from. Python is silent there too: Label("x") >> b, with nothing
upstream of the label, sets no label either. Tests now pin both, and the
nesting that Label(Label(x, "inner"), "outer") reads as.
DagRef.Task records the edges Inputs declares so that the cycle check sees
them, which is the case ADR-0008 names: b := dag.Task(B, Inputs(a)) then
b.Before(a) is a genuine cycle in accepted syntax. Only the upstream and
downstream lists were asserted, so deleting that recording failed no test
over the cycle itself.
Before, After and Inputs each walked the graph to see whether the edge they
were recording closed a cycle, which made building a Dag cost one walk per
edge: a 8,000-task chain declared in reverse took seconds. Register now
walks the graph once, after the Dag is whole, as the Java SDK's
Bundle.register does and the TypeScript SDK plans to. The same chain builds
in tens of milliseconds.

Recording an edge no longer panics, so a fan-out that is rejected leaves the
Dag exactly as it was, and the caveat about a cycle keeping the earlier
edges of its call is gone with it. A Dag whose dependencies contain a cycle
stays unregistered, so it can still be corrected.

A task ordered against itself is still rejected where it is declared: that
needs no walk.
@henry3260
henry3260 force-pushed the go-sdk-node-before-after-label branch from 653028a to 2effac3 Compare October 3, 2026 14:09
@henry3260

Copy link
Copy Markdown
Contributor Author

My one design decision concern before addressing the claude review:

Would it be better to validate the graph at once "after the Dag be registered" so that we don't need to traverse whole graph again whenever a new edge be added?

Just like the TS side that we will only validate the cycle at the "finalize stage". I think we should do the same "last mile validation right before the serve" instead of having iterative validation (since we will compile the SDK into artifact anyway)

Agreed, done in the latest push. The cycle check moved into markRegistered, next to the existing Then check, so it runs once over the whole graph instead of once per edge: O(V+E) rather than O(E*(V+E)).

I rebuilt your reverse-declared chain locally, 2,000 tasks in 9ms and 8,000 in 24ms end to end (building the Dag plus the Register walk), against the 0.19s/2.1s you measured. Different machine, but it is linear now.

This also removes the need for the rollback in your other comment: recording an edge no longer panics, so a rejected fan-out leaves the Dag exactly as it was, and the caveat in the Before doc is gone with it. A task ordered against itself is still rejected where it is declared, since that needs no walk.

@jason810496 jason810496 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, LGTM now.
After solving the final comment from claude review, it's good to merge.

Comment thread go-sdk/airflow/dag.go
IfRef.Then and IfRef.Else stored the task on the condition without recording
an edge, so the ordering their doc states, that the two run after the
condition, was in no Dag the SDK builds: the task had no upstream, and a
cycle running through a condition passed Register.

	read := dag.Task(readRows)
	loaded := dag.Task(load)
	dag.If(hasRows, airflow.Inputs(read)).Then(loaded)
	loaded.Before(read) // read -> hasRows -> load -> read

Naming a task now records the edge from the condition to it, which is also
the upstream the serialized Dag needs. Declaring that edge again, as an
author who writes it out does, changes nothing.
@henry3260
henry3260 merged commit f070997 into apache:main Oct 3, 2026
90 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants