henry3260 opened a new pull request, #74119:
URL: https://github.com/apache/airflow/pull/74119

   ## Why
   
   A Dag authored in Go could declare the tasks it runs 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`. There was also no way to label an edge, which the 
grid and graph
   views render.
   
   This implements decisions 6, 7 and 8 of
   
[`go-sdk/adr/0008-native-dag-interface.md`](https://github.com/apache/airflow/blob/main/go-sdk/adr/0008-native-dag-interface.md):
   `airflow.Node` is the Go counterpart of Python's 
`DAGNode`/`DependencyMixin`, where
   `set_upstream` and `set_downstream` live.
   
   ```go
   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) — the authoring surface and the edges 
behind it:
   
   - `Node`, sealed by an unexported `node()` method, satisfied by `*TaskRef`.
   - `TaskRef.Before` / `TaskRef.After`, variadic so one call fans out. Both 
return their
     argument set as one `Node` rather than the receiver, so `a.Before(b, 
c).Before(d)` is
     `a >> [b, c] >> d` and not a second fan-out from `a`.
   - `Label(node, text)`, which wraps the endpoint rather than the call, so 
each edge of a
     fan-out can carry a label of its own. The label applies to the one verb it 
is passed to
     and is not carried onward by the `Node` that verb returns.
   - Declaring an edge the Dag already has only applies the label. That is how 
an
     `airflow.Inputs` edge gets one: `extracted.Before(Label(transformed, 
"rows"))`.
   - Rejected at declaration time, each with a message naming the tasks: a task 
ordered
     against itself, an edge between two Dags, a `*TaskRef` that `DagRef.Task` 
did not return,
     an edge declared after `Register`, a second label on one edge, and a cycle 
(reported with
     the path that closes it). Every check but the cycle runs over the whole 
fan-out before any
     of it is recorded, the way `DagRef.Task` settles its checks before it 
writes a task.
   
   **`go-sdk/airflow/dag.go`** (+17) — `DagRef.edgeLabels` holds every edge of 
the Dag and its
   label; `TaskRef.upstreams`/`downstreams` hold both sides of each edge. 
`DagRef.Task` now
   records the edges `airflow.Inputs` declares into the same structure, which 
is what makes
   redeclaring one idempotent.
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes — Claude Code (Opus 5)


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to