This is an automated email from the ASF dual-hosted git repository.

jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new a2cf3289367 Go SDK: answer the Dag-parsing request from Serve (#74391)
a2cf3289367 is described below

commit a2cf328936783aff1ced3b70da920ccb7895a9d2
Author: PoAn Yang <[email protected]>
AuthorDate: Wed Oct 7 17:15:33 2026 +0800

    Go SDK: answer the Dag-parsing request from Serve (#74391)
    
    * Go SDK: answer the Dag-parsing request from Serve
    
    Signed-off-by: PoAn Yang <[email protected]>
    
    * Drop the stderr excerpt from the Lang-SDK Dag processor
    
    The stderr excerpt is unrelated to Go Dag parsing, so this PR leaves it out.
    
    * Go SDK: serve the Dag parse request from bundle.Bundle
    
    Serve takes a bundle.Bundle again, and the Dag serializer is an optional 
interface the bundle can implement.
    
    * Go SDK: reply to the Dag parse request on its request id
    
    The runtime replies on the request id and exits without waiting for an 
acknowledgement, as the TS runtime and ADR 0003 do.
    
    * Go SDK: log the stack when serializing a Dag panics
    
    The panic becomes the Dag's error, and the stack now reaches the Dag 
processor log through the default slog logger.
    
    * Go SDK: trim the comments of the Dag parse path
    
    The comments now state the behavior in one or two lines.
    
    ---------
    
    Signed-off-by: PoAn Yang <[email protected]>
    Co-authored-by: ZHE YOU LIU <[email protected]>
---
 go-sdk/README.md                         |   4 +
 go-sdk/airflow/bundle.go                 |  55 ++++++++-
 go-sdk/airflow/bundle_test.go            |  41 +++++++
 go-sdk/airflow/dag.go                    |   4 +-
 go-sdk/airflow/serve.go                  |   7 +-
 go-sdk/airflow/serve_test.go             | 102 +++++++++++++---
 go-sdk/internal/bundle/doc.go            |   8 +-
 go-sdk/internal/bundle/task.go           |  14 +++
 go-sdk/pkg/execution/dag_parse.go        |  79 ++++++++++++
 go-sdk/pkg/execution/dag_parse_test.go   | 202 +++++++++++++++++++++++++++++++
 go-sdk/pkg/execution/integration_test.go |  65 +++++++++-
 go-sdk/pkg/execution/server.go           |  22 ++--
 12 files changed, 569 insertions(+), 34 deletions(-)

diff --git a/go-sdk/README.md b/go-sdk/README.md
index ad344112abe..f099bee045e 100644
--- a/go-sdk/README.md
+++ b/go-sdk/README.md
@@ -389,6 +389,10 @@ Python supervisor / task runner
   protocol on the comm socket, with structured JSON-line logs on the logs 
socket.
 - The Python runtime is the worker. It proxies every `GetConnection` / 
`GetVariable` / `GetXCom` /
   `SetXCom` call through to the Execution API. The Go binary just runs the 
task function.
+- When the first frame on the comm socket is a `DagFileParseRequest` from the 
Dag processor, the
+  binary runs no task. It answers with one `DagFileParsingResult` that holds 
the Dags from
+  `airflow.Dag` that the binary registered. The Dags are serialized as
+  [Serializing a native Dag](#serializing-a-native-dag) describes. The binary 
then exits.
 
 The Go side of the protocol is implemented in `pkg/execution/`. On the Python 
side it is the
 `ExecutableCoordinator` in 
`task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py`.
diff --git a/go-sdk/airflow/bundle.go b/go-sdk/airflow/bundle.go
index 9abb02717a0..df10c5d5255 100644
--- a/go-sdk/airflow/bundle.go
+++ b/go-sdk/airflow/bundle.go
@@ -19,6 +19,8 @@ package airflow
 
 import (
        "fmt"
+       "log/slog"
+       "runtime/debug"
        "slices"
        "sync"
        "sync/atomic"
@@ -193,10 +195,11 @@ func (m *taskHandlerMap) ListTaskHandlers() 
[]bundle.TaskHandlerInfo {
        return slices.Clone(m.order)
 }
 
-// dagMap holds the registered Dags by dag_id.
+// dagMap holds the registered Dags by dag_id, in registration order.
 type dagMap struct {
-       mu   sync.Mutex
-       dags map[string]*DagRef
+       mu    sync.Mutex
+       dags  map[string]*DagRef
+       order []*DagRef
 }
 
 func (m *dagMap) add(dag *DagRef) {
@@ -211,6 +214,7 @@ func (m *dagMap) add(dag *DagRef) {
                m.dags = make(map[string]*DagRef)
        }
        m.dags[dag.dagID] = dag
+       m.order = append(m.order, dag)
 }
 
 func (m *dagMap) has(dagID string) bool {
@@ -220,3 +224,48 @@ func (m *dagMap) has(dagID string) bool {
        _, exists := m.dags[dagID]
        return exists
 }
+
+// serialize serializes the Dags in registration order. If a Dag panics, the 
panic becomes that
+// Dag's Err.
+func (m *dagMap) serialize(fileloc, relativeFileloc string) 
[]bundle.SerializedDag {
+       m.mu.Lock()
+       dags := slices.Clone(m.order)
+       m.mu.Unlock()
+
+       serialized := make([]bundle.SerializedDag, len(dags))
+       for i, dag := range dags {
+               serialized[i] = serializeRecovering(dag, fileloc, 
relativeFileloc)
+       }
+       return serialized
+}
+
+func serializeRecovering(dag *DagRef, fileloc, relativeFileloc string) (s 
bundle.SerializedDag) {
+       s.DagID = dag.dagID
+       defer func() {
+               if r := recover(); r != nil {
+                       slog.Error(
+                               "Dag serialization panicked",
+                               "dag_id", dag.dagID, "panic", r, "stack", 
string(debug.Stack()),
+                       )
+                       s.Data, s.Err = nil, fmt.Errorf("%v", r)
+               }
+       }()
+       s.Data = dag.serialize(fileloc, relativeFileloc)
+       return s
+}
+
+// coordinatorSource is what Serve hands to execution.Serve: the task handlers 
and the Dags of one
+// bundle.
+type coordinatorSource struct {
+       *taskHandlerMap
+       dags *dagMap
+}
+
+var (
+       _ bundle.Bundle        = coordinatorSource{}
+       _ bundle.DagSerializer = coordinatorSource{}
+)
+
+func (s coordinatorSource) SerializeDags(fileloc, relativeFileloc string) 
[]bundle.SerializedDag {
+       return s.dags.serialize(fileloc, relativeFileloc)
+}
diff --git a/go-sdk/airflow/bundle_test.go b/go-sdk/airflow/bundle_test.go
index 7a5874e29b4..d6c2da67f51 100644
--- a/go-sdk/airflow/bundle_test.go
+++ b/go-sdk/airflow/bundle_test.go
@@ -412,3 +412,44 @@ func TestRegisterableRejectsForeignTypes(t *testing.T) {
        assert.Contains(t, string(out), "foreignItem does not implement 
airflow.Registerable")
        assert.Contains(t, string(out), "unexported method registerable")
 }
+
+func TestSerializeDagsKeepsTheOtherDagsWhenADagCannotBeSerialized(t 
*testing.T) {
+       b := Bundle()
+       etl := Dag("etl")
+       etl.Task(noop)
+       b.Register(etl)
+       // A Dag that skipped Register stands in for one the serializer fails 
on.
+       broken := Dag("broken")
+       b.dags.dags["broken"] = broken
+       b.dags.order = append(b.dags.order, broken)
+       reports := Dag("reports")
+       reports.Task(noop)
+       b.Register(reports)
+
+       serialized := b.dags.serialize("/bundles/go/etl", "etl")
+
+       require.Len(t, serialized, 3)
+       assert.Equal(t, bundle.SerializedDag{
+               DagID: "etl",
+               Data:  etl.serialize("/bundles/go/etl", "etl"),
+       }, serialized[0])
+       assert.Equal(t, "broken", serialized[1].DagID)
+       assert.Nil(t, serialized[1].Data)
+       assert.ErrorContains(t, serialized[1].Err, `Dag "broken" is not 
registered`)
+       assert.Equal(t, bundle.SerializedDag{
+               DagID: "reports",
+               Data:  reports.serialize("/bundles/go/etl", "etl"),
+       }, serialized[2])
+}
+
+func TestSerializeDagsLeavesOutADagThatRegisterRejected(t *testing.T) {
+       b := Bundle()
+       cyclic := Dag("cyclic")
+       extracted := orderedTask(t, cyclic, "extract")
+       loaded := orderedTask(t, cyclic, "load")
+       extracted.Before(loaded)
+       loaded.Before(extracted)
+       require.Panics(t, func() { b.Register(cyclic) })
+
+       assert.Empty(t, b.dags.serialize("/bundles/go/etl", "etl"))
+}
diff --git a/go-sdk/airflow/dag.go b/go-sdk/airflow/dag.go
index e95505de0ca..9415e514caf 100644
--- a/go-sdk/airflow/dag.go
+++ b/go-sdk/airflow/dag.go
@@ -71,8 +71,8 @@ type DagRef struct {
 // [IfRef.Then], [IfRef.Else], [SwitchRef.Case] and the methods of 
[TaskGroupRef] panic once the Dag
 // is registered.
 //
-// [BundleRef.Serve] does not yet serve the Dags that Dag returns. It leaves 
them out of the
-// --airflow-metadata manifest and cannot run their tasks.
+// [BundleRef.Serve] sends the registered Dags to the Dag processor, but does 
not yet list them in
+// the --airflow-metadata manifest or run their tasks.
 //
 // Dag panics if it gets more than one DagSpec, or if the DagSpec has a value 
that Python rejects
 // when it builds or validates a Dag:
diff --git a/go-sdk/airflow/serve.go b/go-sdk/airflow/serve.go
index bcfb33ab934..7fafd795965 100644
--- a/go-sdk/airflow/serve.go
+++ b/go-sdk/airflow/serve.go
@@ -57,8 +57,9 @@ const (
 // The command-line flags of the executable decide what Serve does.
 // With --airflow-metadata it prints the bundle's manifest and returns, which 
is how
 // airflow-go-pack reads the Dag and task ids of the registered task handlers.
-// With --comm and --logs, which the Airflow supervisor passes, it runs one 
task over the
-// coordinator protocol.
+// Airflow starts the executable with --comm and --logs. Serve then speaks the 
coordinator
+// protocol and either runs one task or answers the Dag processor's parse 
request with the Dags
+// from [Dag].
 //
 // main must exit with a non-zero status when Serve returns an error, because 
the exit status
 // is how the supervisor learns that the task failed:
@@ -141,7 +142,7 @@ func (b *BundleRef) serve(args []string, stdout io.Writer) 
error {
                }
                return execution.DumpAirflowMetadata(stdout, &b.taskHandlers, 
format)
        case modeCoordinator:
-               return execution.Serve(&b.taskHandlers, *commAddr, *logsAddr)
+               return execution.Serve(coordinatorSource{&b.taskHandlers, 
&b.dags}, *commAddr, *logsAddr)
        case modeCoordinatorUsageError:
                return errCoordinatorFlagsRequired
        }
diff --git a/go-sdk/airflow/serve_test.go b/go-sdk/airflow/serve_test.go
index 93de70af023..bca506bc14c 100644
--- a/go-sdk/airflow/serve_test.go
+++ b/go-sdk/airflow/serve_test.go
@@ -33,6 +33,7 @@ import (
        "gopkg.in/yaml.v3"
 
        "github.com/apache/airflow/go-sdk/pkg/execution"
+       "github.com/apache/airflow/go-sdk/pkg/execution/genmodels"
 )
 
 func TestDecideMode(t *testing.T) {
@@ -219,28 +220,22 @@ func TestServeHelpIsNotAnError(t *testing.T) {
        assert.Contains(t, stdout.String(), "--airflow-metadata")
 }
 
-// A fake supervisor sends StartupDetails over the comm socket, as the Python
-// ExecutableCoordinator does after it starts the bundle with --comm and 
--logs.
-func TestServeRunsTaskForSupervisor(t *testing.T) {
+// serveForSupervisor runs b.serve with --comm and --logs and returns the 
supervisor's end of the
+// comm socket and the channel that gets what serve returns.
+func serveForSupervisor(t *testing.T, b *BundleRef) 
(*execution.CoordinatorComm, <-chan error) {
+       t.Helper()
        commLn, err := net.Listen("tcp", "127.0.0.1:0")
        require.NoError(t, err)
-       defer commLn.Close()
+       t.Cleanup(func() { commLn.Close() })
        logsLn, err := net.Listen("tcp", "127.0.0.1:0")
        require.NoError(t, err)
-       defer logsLn.Close()
+       t.Cleanup(func() { logsLn.Close() })
        // Without a deadline, a Serve that never dials would leave Accept 
blocked until the test
        // binary times out.
        deadline := time.Now().Add(10 * time.Second)
        require.NoError(t, commLn.(*net.TCPListener).SetDeadline(deadline))
        require.NoError(t, logsLn.(*net.TCPListener).SetDeadline(deadline))
 
-       ran := false
-       b := Bundle()
-       b.Register(TaskHandler("py_etl", "transform", func(Context) error {
-               ran = true
-               return nil
-       }))
-
        done := make(chan error, 1)
        go func() {
                done <- b.serve(
@@ -251,13 +246,26 @@ func TestServeRunsTaskForSupervisor(t *testing.T) {
 
        commConn, err := commLn.Accept()
        require.NoError(t, err)
-       defer commConn.Close()
+       t.Cleanup(func() { commConn.Close() })
        logsConn, err := logsLn.Accept()
        require.NoError(t, err)
-       defer logsConn.Close()
+       t.Cleanup(func() { logsConn.Close() })
        require.NoError(t, commConn.SetDeadline(deadline))
 
-       supervisor := execution.NewCoordinatorComm(commConn, commConn, 
discardLogger())
+       return execution.NewCoordinatorComm(commConn, commConn, 
discardLogger()), done
+}
+
+// A fake supervisor sends StartupDetails over the comm socket, as the Python
+// ExecutableCoordinator does after it starts the bundle with --comm and 
--logs.
+func TestServeRunsTaskForSupervisor(t *testing.T) {
+       ran := false
+       b := Bundle()
+       b.Register(TaskHandler("py_etl", "transform", func(Context) error {
+               ran = true
+               return nil
+       }))
+
+       supervisor, done := serveForSupervisor(t, b)
        require.NoError(t, supervisor.SendRequest(0, map[string]any{
                "type": "StartupDetails",
                "ti": map[string]any{
@@ -284,3 +292,67 @@ func TestServeRunsTaskForSupervisor(t *testing.T) {
        }
        assert.True(t, ran)
 }
+
+// The Dag processor sends a DagFileParseRequest instead of StartupDetails.
+func TestServeAnswersTheDagParseRequestWithoutRunningATask(t *testing.T) {
+       var ran []string
+       record := func(name string) func(Context) error {
+               return func(Context) error {
+                       ran = append(ran, name)
+                       return nil
+               }
+       }
+       etl := Dag("etl")
+       extracted := etl.Task(func(Context) (int, error) {
+               ran = append(ran, "extract")
+               return 3, nil
+       }, TaskSpec{TaskID: "extract"})
+       loaded := etl.Task(record("load"), TaskSpec{TaskID: "load"})
+       etl.If(func(_ Context, rows int) (bool, error) {
+               ran = append(ran, "has_rows")
+               return rows > 0, nil
+       }, Inputs(extracted), TaskSpec{TaskID: "has_rows"}).Then(loaded)
+       reports := Dag("reports")
+       reports.Task(record("publish"), TaskSpec{TaskID: "publish"})
+
+       b := Bundle()
+       b.Register(etl, TaskHandler("py_etl", "transform", 
record("transform")), reports)
+
+       supervisor, done := serveForSupervisor(t, b)
+       const requestID = 7
+       require.NoError(t, supervisor.SendRequest(requestID, map[string]any{
+               "type":        "DagFileParseRequest",
+               "file":        "/bundles/go/etl",
+               "bundle_path": "/bundles/go",
+               "bundle_name": "go",
+       }))
+
+       frame, err := supervisor.ReadMessage()
+       require.NoError(t, err)
+       assert.EqualValues(t, requestID, frame.ID)
+       var result genmodels.DagFileParsingResult
+       require.NoError(t, msgpack.Unmarshal(frame.Body, &result))
+       assert.Equal(t, "DagFileParsingResult", result.Type)
+       assert.Equal(t, "/bundles/go/etl", result.Fileloc)
+       assert.Nil(t, result.ImportErrors)
+       require.Len(t, result.SerializedDags, 2)
+       for i, dag := range []*DagRef{etl, reports} {
+               want := jsonOf(t, dag.serialize("/bundles/go/etl", "etl"))
+               assertJSON(t, want, result.SerializedDags[i].Data)
+       }
+
+       select {
+       case err := <-done:
+               require.NoError(t, err)
+       case <-time.After(2 * time.Second):
+               t.Fatal("Serve did not return after it answered the Dag parse 
request")
+       }
+       assert.Empty(t, ran)
+}
+
+func jsonOf(t *testing.T, value any) string {
+       t.Helper()
+       raw, err := json.Marshal(value)
+       require.NoError(t, err)
+       return string(raw)
+}
diff --git a/go-sdk/internal/bundle/doc.go b/go-sdk/internal/bundle/doc.go
index dda8b6ed6d4..8a77698fe70 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/internal/bundle/doc.go
@@ -15,8 +15,10 @@
 // specific language governing permissions and limitations
 // under the License.
 
-// Package bundle defines what the coordinator runtime needs from a bundle:
-// the tasks it looks up and runs, and the Dag and task ids it lists in the 
manifest.
+// Package bundle defines what the coordinator runtime needs from a bundle: the
+// tasks it looks up and runs, the Dag and task ids it lists in the manifest, 
and
+// the serialized Dags it sends to the Dag processor.
 //
-// Package airflow builds both from the task handlers a bundle registers.
+// Package airflow builds the tasks and ids from the task handlers a bundle
+// registers, and the serialized Dags from its Dags.
 package bundle
diff --git a/go-sdk/internal/bundle/task.go b/go-sdk/internal/bundle/task.go
index c0cccb95697..c3cb8cc6557 100644
--- a/go-sdk/internal/bundle/task.go
+++ b/go-sdk/internal/bundle/task.go
@@ -57,6 +57,20 @@ type EnumerableBundle interface {
        ListTaskHandlers() []TaskHandlerInfo
 }
 
+// SerializedDag is one Dag from airflow.Dag, serialized for the Dag 
processor. Data is nil when
+// Err is set.
+type SerializedDag struct {
+       DagID string
+       Data  map[string]any
+       Err   error
+}
+
+// DagSerializer serializes the Dags from airflow.Dag that a bundle 
registered, in registration
+// order.
+type DagSerializer interface {
+       SerializeDags(fileloc, relativeFileloc string) []SerializedDag
+}
+
 type taskFunction struct {
        fn       reflect.Value
        fullName string
diff --git a/go-sdk/pkg/execution/dag_parse.go 
b/go-sdk/pkg/execution/dag_parse.go
new file mode 100644
index 00000000000..96d33416ea7
--- /dev/null
+++ b/go-sdk/pkg/execution/dag_parse.go
@@ -0,0 +1,79 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package execution
+
+import (
+       "fmt"
+       "log/slog"
+       "path/filepath"
+       "strings"
+
+       "github.com/apache/airflow/go-sdk/internal/bundle"
+       "github.com/apache/airflow/go-sdk/pkg/execution/genmodels"
+)
+
+// parseDags answers a DagFileParseRequest with the bundle's Dags from 
airflow.Dag. A Dag that
+// cannot be serialized adds a line to the import error of the file, and the 
other Dags are still
+// sent.
+func parseDags(
+       s bundle.DagSerializer,
+       req *genmodels.DagFileParseRequest,
+       logger *slog.Logger,
+) genmodels.DagFileParsingResult {
+       relative := computeRelativeFileloc(req.File, req.BundlePath)
+       // A nil slice would be sent as null, which the Dag processor rejects.
+       result := genmodels.DagFileParsingResult{
+               Fileloc:        req.File,
+               SerializedDags: []genmodels.LazyDeserializedDAG{},
+       }
+       var dagIDs, failures []string
+       if s != nil {
+               for _, dag := range s.SerializeDags(req.File, relative) {
+                       dagIDs = append(dagIDs, dag.DagID)
+                       if dag.Err != nil {
+                               logger.Error("Dag could not be serialized", 
"dag_id", dag.DagID, "error", dag.Err)
+                               failures = append(failures, fmt.Sprintf("Dag 
%q: %v", dag.DagID, dag.Err))
+                               continue
+                       }
+                       result.SerializedDags = append(
+                               result.SerializedDags, 
genmodels.LazyDeserializedDAG{Data: dag.Data},
+                       )
+               }
+       }
+       if len(failures) > 0 {
+               result.ImportErrors = &genmodels.ImportErrors{relative: 
strings.Join(failures, "\n")}
+       }
+       logger.Info("Parse-mode response",
+               "fileloc", req.File,
+               "dag_ids", dagIDs,
+               "serialized", len(result.SerializedDags),
+               "import_errors", len(failures),
+       )
+       return result
+}
+
+// computeRelativeFileloc returns file relative to bundlePath, in the same way 
that Python's DagBag
+// computes the relative_fileloc of a Dag file. Like DagBag, it returns file 
unchanged when file is
+// not inside bundlePath.
+func computeRelativeFileloc(file, bundlePath string) string {
+       rel, err := filepath.Rel(bundlePath, file)
+       if err != nil || rel == ".." || strings.HasPrefix(rel, 
".."+string(filepath.Separator)) {
+               return file
+       }
+       return rel
+}
diff --git a/go-sdk/pkg/execution/dag_parse_test.go 
b/go-sdk/pkg/execution/dag_parse_test.go
new file mode 100644
index 00000000000..27f4d8a95d4
--- /dev/null
+++ b/go-sdk/pkg/execution/dag_parse_test.go
@@ -0,0 +1,202 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package execution
+
+import (
+       "bytes"
+       "encoding/json"
+       "errors"
+       "io"
+       "log/slog"
+       "testing"
+
+       "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
+       "github.com/vmihailenco/msgpack/v5"
+
+       "github.com/apache/airflow/go-sdk/internal/bundle"
+       "github.com/apache/airflow/go-sdk/pkg/execution/genmodels"
+)
+
+// serializedDags is a DagSerializer that records the paths it was called with.
+type serializedDags struct {
+       dags              []bundle.SerializedDag
+       fileloc, relative string
+}
+
+func (s *serializedDags) SerializeDags(fileloc, relativeFileloc string) 
[]bundle.SerializedDag {
+       s.fileloc, s.relative = fileloc, relativeFileloc
+       return s.dags
+}
+
+func serializedDag(dagID string) bundle.SerializedDag {
+       return bundle.SerializedDag{
+               DagID: dagID,
+               Data:  map[string]any{"__version": 3, "dag": 
map[string]any{"dag_id": dagID}},
+       }
+}
+
+var etlParseRequest = &genmodels.DagFileParseRequest{
+       File:       "/bundles/go/etl",
+       BundlePath: "/bundles/go",
+       BundleName: "go",
+}
+
+func discardLogger() *slog.Logger { return 
slog.New(slog.NewTextHandler(io.Discard, nil)) }
+
+// wireJSON returns body as the Dag processor receives it, after the frame 
encoding of SendRequest.
+func wireJSON(t *testing.T, body any) string {
+       t.Helper()
+       payload, err := encodeRequest(0, body)
+       require.NoError(t, err)
+       frame, err := decodeFrame(payload)
+       require.NoError(t, err)
+       dec := msgpack.NewDecoder(bytes.NewReader(frame.Body))
+       dec.UseLooseInterfaceDecoding(true)
+       var decoded any
+       require.NoError(t, dec.Decode(&decoded))
+       raw, err := json.Marshal(decoded)
+       require.NoError(t, err)
+       return string(raw)
+}
+
+func TestParseDagsAnswersWithTheSerializedDags(t *testing.T) {
+       dags := &serializedDags{dags: []bundle.SerializedDag{
+               serializedDag("etl"),
+               serializedDag("reports"),
+       }}
+
+       result := parseDags(dags, etlParseRequest, discardLogger())
+
+       assert.Equal(t, "/bundles/go/etl", dags.fileloc)
+       assert.Equal(t, "etl", dags.relative)
+       assert.JSONEq(t, `{
+               "type": "DagFileParsingResult",
+               "fileloc": "/bundles/go/etl",
+               "serialized_dags": [
+                       {"data": {"__version": 3, "dag": {"dag_id": "etl"}}},
+                       {"data": {"__version": 3, "dag": {"dag_id": "reports"}}}
+               ]
+       }`, wireJSON(t, result))
+}
+
+func TestParseDagsAnswersWithNoDagsWhenTheBundleHasNoSerializer(t *testing.T) {
+       result := parseDags(nil, etlParseRequest, discardLogger())
+
+       assert.NotNil(t, result.SerializedDags)
+       assert.Empty(t, result.SerializedDags)
+       assert.Nil(t, result.ImportErrors)
+}
+
+func TestParseDagsSendsAnEmptyListForABundleWithoutDags(t *testing.T) {
+       result := parseDags(&serializedDags{}, etlParseRequest, discardLogger())
+
+       assert.JSONEq(t, `{
+               "type": "DagFileParsingResult",
+               "fileloc": "/bundles/go/etl",
+               "serialized_dags": []
+       }`, wireJSON(t, result))
+}
+
+func TestParseDagsReportsEachDagThatCannotBeSerialized(t *testing.T) {
+       dags := &serializedDags{dags: []bundle.SerializedDag{
+               serializedDag("etl"),
+               {DagID: "reports", Err: errors.New("no schema field")},
+               serializedDag("cleanup"),
+               {DagID: "audit", Err: errors.New("bad field type")},
+       }}
+       var logs bytes.Buffer
+       logger := slog.New(slog.NewJSONHandler(&logs, nil))
+
+       result := parseDags(dags, etlParseRequest, logger)
+
+       assert.JSONEq(t, `{
+               "type": "DagFileParsingResult",
+               "fileloc": "/bundles/go/etl",
+               "serialized_dags": [
+                       {"data": {"__version": 3, "dag": {"dag_id": "etl"}}},
+                       {"data": {"__version": 3, "dag": {"dag_id": "cleanup"}}}
+               ],
+               "import_errors": {
+                       "etl": "Dag \"reports\": no schema field\nDag 
\"audit\": bad field type"
+               }
+       }`, wireJSON(t, result))
+
+       var logged []map[string]any
+       for line := range bytes.Lines(logs.Bytes()) {
+               var record map[string]any
+               require.NoError(t, json.Unmarshal(line, &record))
+               logged = append(logged, record)
+       }
+       require.Len(t, logged, 3)
+       for i, dagID := range []string{"reports", "audit"} {
+               assert.Equal(t, "ERROR", logged[i]["level"])
+               assert.Equal(t, dagID, logged[i]["dag_id"])
+       }
+       assert.Equal(t, "INFO", logged[2]["level"])
+       assert.Equal(t, []any{"etl", "reports", "cleanup", "audit"}, 
logged[2]["dag_ids"])
+       assert.EqualValues(t, 2, logged[2]["serialized"])
+       assert.EqualValues(t, 2, logged[2]["import_errors"])
+}
+
+func TestComputeRelativeFileloc(t *testing.T) {
+       tests := []struct {
+               name       string
+               file       string
+               bundlePath string
+               want       string
+       }{
+               {name: "in the bundle", file: "/bundles/go/etl", bundlePath: 
"/bundles/go", want: "etl"},
+               {
+                       name:       "in a directory of the bundle",
+                       file:       "/bundles/go/dags/etl",
+                       bundlePath: "/bundles/go",
+                       want:       "dags/etl",
+               },
+               {
+                       name:       "the bundle itself",
+                       file:       "/bundles/go/etl",
+                       bundlePath: "/bundles/go/etl",
+                       want:       ".",
+               },
+               {
+                       name:       "a name that starts with two dots",
+                       file:       "/bundles/go/..etl",
+                       bundlePath: "/bundles/go",
+                       want:       "..etl",
+               },
+               {
+                       name:       "outside the bundle",
+                       file:       "/bundles/gopher/etl",
+                       bundlePath: "/bundles/go",
+                       want:       "/bundles/gopher/etl",
+               },
+               {
+                       name:       "the directory that holds the bundle",
+                       file:       "/bundles",
+                       bundlePath: "/bundles/go",
+                       want:       "/bundles",
+               },
+               {name: "a relative file", file: "etl", bundlePath: 
"/bundles/go", want: "etl"},
+       }
+       for _, tt := range tests {
+               t.Run(tt.name, func(t *testing.T) {
+                       assert.Equal(t, tt.want, 
computeRelativeFileloc(tt.file, tt.bundlePath))
+               })
+       }
+}
diff --git a/go-sdk/pkg/execution/integration_test.go 
b/go-sdk/pkg/execution/integration_test.go
index 664563511d4..8e4b24358c8 100644
--- a/go-sdk/pkg/execution/integration_test.go
+++ b/go-sdk/pkg/execution/integration_test.go
@@ -1361,7 +1361,8 @@ func TestServeFailureAfterConnectClosesComm(t *testing.T) 
{
        logsConn := <-logsCh
        defer logsConn.Close()
 
-       // Serve expects StartupDetails as the first frame, so it fails to 
decode a VariableResult.
+       // Serve expects StartupDetails or DagFileParseRequest first, so it 
fails to decode a
+       // VariableResult.
        payload, err := encodeRequest(
                0,
                map[string]any{"type": "VariableResult", "key": "k", "value": 
"v"},
@@ -1382,3 +1383,65 @@ func TestServeFailureAfterConnectClosesComm(t 
*testing.T) {
        _, err = readFrame(commConn)
        require.Error(t, err)
 }
+
+// parseBundle is a bundle that serializes Dags and has no task to run.
+type parseBundle struct {
+       testBundle
+       *serializedDags
+}
+
+const dagParseRequestID = 7
+
+// startDagParse runs Serve for dags and sends it a DagFileParseRequest. It 
returns the frame that
+// Serve answers with and the channel that gets what Serve returns.
+func startDagParse(t *testing.T, dags *serializedDags) (frame IncomingFrame, 
done <-chan error) {
+       t.Helper()
+       commAddr, logsAddr, commCh, logsCh, cleanup := startSupervisor(t)
+       t.Cleanup(cleanup)
+
+       served := make(chan error, 1)
+       go func() { served <- Serve(parseBundle{testBundle{}, dags}, commAddr, 
logsAddr) }()
+
+       commConn := <-commCh
+       t.Cleanup(func() { commConn.Close() })
+       logsConn := <-logsCh
+       t.Cleanup(func() { logsConn.Close() })
+       deadline := time.Now().Add(10 * time.Second)
+       require.NoError(t, commConn.SetDeadline(deadline))
+       require.NoError(t, logsConn.SetDeadline(deadline))
+
+       payload, err := encodeRequest(dagParseRequestID, map[string]any{
+               "type":        "DagFileParseRequest",
+               "file":        "/bundles/go/etl",
+               "bundle_path": "/bundles/go",
+               "bundle_name": "go",
+       })
+       require.NoError(t, err)
+       require.NoError(t, writeFrame(commConn, payload))
+
+       frame, err = readFrame(commConn)
+       require.NoError(t, err)
+       require.True(t, isNilRaw(frame.Err))
+       return frame, served
+}
+
+func TestServeDagFileParseRequestEndToEnd(t *testing.T) {
+       dags := &serializedDags{dags: 
[]bundle.SerializedDag{serializedDag("etl")}}
+       frame, done := startDagParse(t, dags)
+
+       assert.EqualValues(t, dagParseRequestID, frame.ID)
+       var result genmodels.DagFileParsingResult
+       require.NoError(t, decodeBody(frame.Body, &result))
+       assert.Equal(t, "DagFileParsingResult", result.Type)
+       assert.Equal(t, "/bundles/go/etl", result.Fileloc)
+       require.Len(t, result.SerializedDags, 1)
+       assert.Equal(t, "etl", 
result.SerializedDags[0].Data["dag"].(map[string]any)["dag_id"])
+       assert.Equal(t, "etl", dags.relative)
+
+       select {
+       case err := <-done:
+               require.NoError(t, err)
+       case <-time.After(2 * time.Second):
+               t.Fatal("Serve did not return after it sent the Dag parsing 
result")
+       }
+}
diff --git a/go-sdk/pkg/execution/server.go b/go-sdk/pkg/execution/server.go
index 0f56e005315..062658930d0 100644
--- a/go-sdk/pkg/execution/server.go
+++ b/go-sdk/pkg/execution/server.go
@@ -20,8 +20,8 @@
 // the Airflow supervisor (Python ExecutableCoordinator), the Serve method of
 // airflow.BundleRef dispatches here.
 //
-// The first inbound frame on the comm socket is a StartupDetails message
-// that drives multi-round task execution.
+// The first frame on the comm socket picks the mode: StartupDetails runs one
+// task, and DagFileParseRequest is answered with the bundle's serialized Dags.
 //
 // See go-sdk/adr/0003-coordinator-protocol-msgpack-ipc.md.
 package execution
@@ -58,11 +58,9 @@ const terminalSendTimeout = 30 * time.Second
 // comm and logs sockets, installs an slog handler that writes JSON-line
 // records to the logs connection, and dispatches on the first frame.
 //
-// Serve returns nil on a clean shutdown: the task ran and its terminal
-// TaskState/SucceedTask frame was delivered, and the caller should exit 0. A
-// non-nil error indicates a protocol-level failure (connection loss,
-// malformed frames, unknown first message type) that happens before or
-// instead of delivering a terminal frame.
+// Serve returns nil once it delivers the terminal frame of a task run or the 
DagFileParsingResult
+// of a Dag parse, and the caller should then exit 0. A non-nil error 
indicates a protocol-level
+// failure (connection loss, malformed frames, unknown first message type) 
before that.
 //
 // Failure-signaling contract: the caller (main) must turn a non-nil error
 // into a non-zero process exit. The supervisor derives the task's final state
@@ -167,6 +165,16 @@ func Serve(b bundle.Bundle, commAddr, logsAddr string) 
error {
                }
                logger.Debug("Task execution complete")
 
+       case *genmodels.DagFileParseRequest:
+               logger.Info("Received Dag parse request", "file", msg.File, 
"bundle_path", msg.BundlePath)
+               serializer, _ := b.(bundle.DagSerializer)
+               result := parseDags(serializer, msg, logger)
+               // Bound the write so a wedged socket cannot hang shutdown.
+               _ = 
commConn.SetWriteDeadline(time.Now().Add(terminalSendTimeout))
+               if err := comm.SendRequest(frame.ID, result); err != nil {
+                       return fmt.Errorf("sending Dag parsing result: %w", err)
+               }
+
        default:
                logger.Error("Unexpected initial message type", "type", 
fmt.Sprintf("%T", body))
                return fmt.Errorf("unexpected initial message type: %T", body)

Reply via email to