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 1921597eeff Go SDK: embed each native Dag's source file in the bundle
(#74342)
1921597eeff is described below
commit 1921597eeffba2666946b90f76f8894f3fa0a40d
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Thu Oct 8 10:38:37 2026 +0800
Go SDK: embed each native Dag's source file in the bundle (#74342)
* Record the source file of each Go SDK native Dag
airflow.Dag records the file that called it, and --airflow-metadata
lists it per native Dag under dag_source_files so the packer can embed
the file. Task handlers add nothing.
* Embed each native Dag's source file in Go bundles
The packer reads dag_source_files from the binary's manifest and embeds
the entrypoint plus the file of every native Dag once, in a stable
order. The manifest drops source and gains entrypoint_path,
dag_source_paths and a sources index. A Dag whose file is outside the
module maps to the entrypoint with a warning.
* Print every embedded source file from airflow-go-pack inspect
inspect --source cuts the source region by the manifest's sources index
and prints each file under its own header. An index that does not fit
the region is an error.
* Document per-Dag source files in Go bundles
The Go SDK README, the Go language SDK page and the example Justfile now
say a bundle embeds the entrypoint plus the source file of each native
Dag.
* Go SDK: run the tasks of Dags built with airflow.Dag
Serve looked tasks up among the task handlers only, so every task of a
native Go Dag
would be reported as not registered once Airflow runs native Go Dags.
* Go SDK: pack an entrypoint at the module root
The module root itself now counts as inside the module, so a main.go at the
top of a plain go mod init project packs.
* Go SDK: quote dag_id keys in the bundle manifest
Dag IDs such as 2024, on or null were written unquoted and read back as
int, bool or None by a YAML parser.
* Go SDK: count native Dags in the pack summary
The summary showed dags=0 for a bundle with only airflow.Dag Dags, so it
now reports task handler and native Dags separately.
* Go SDK: tell the user to repack a bundle without a sources index
The error now says to repack the bundle with this airflow-go-pack instead
of only describing the mismatch.
* Go SDK: test the per-Dag source alias check through runPack
The old test called rejectOutputAlias directly, so removing the call in
runPack for per-Dag files kept the suite green.
* Document how Go maps a factory-built Dag to its source
Go records the file that calls airflow.Dag, so a Dag built in a factory
maps to the factory's file.
---
.../0015-per-dag-source-in-bundle-artifact.md | 13 +-
.../authoring-and-scheduling/language-sdks/go.rst | 5 +-
go-sdk/README.md | 4 +-
go-sdk/airflow/bundle.go | 53 ++-
go-sdk/airflow/bundle_test.go | 27 ++
go-sdk/airflow/dag.go | 11 +-
.../bundle/doc.go => airflow/dag_factory_test.go} | 14 +-
go-sdk/airflow/dag_source_test.go | 65 ++++
go-sdk/airflow/serve.go | 6 +-
go-sdk/cmd/airflow-go-pack/inspect.go | 59 ++-
go-sdk/cmd/airflow-go-pack/inspect_test.go | 89 ++++-
go-sdk/cmd/airflow-go-pack/main.go | 9 +-
go-sdk/cmd/airflow-go-pack/pack.go | 123 ++++--
.../cmd/airflow-go-pack/pack_integration_test.go | 21 +-
go-sdk/cmd/airflow-go-pack/pack_test.go | 102 ++++-
go-sdk/cmd/airflow-go-pack/sources.go | 256 ++++++++++++
go-sdk/cmd/airflow-go-pack/sources_test.go | 430 +++++++++++++++++++++
.../airflow-go-pack/testdata/handlersonly/main.go} | 26 +-
.../testdata/multidag/factory/factory.go} | 19 +-
.../airflow-go-pack/testdata/multidag/main.go} | 39 +-
.../testdata/multidag/reports/reports.go} | 19 +-
go-sdk/example/bundle/Justfile | 4 +-
go-sdk/internal/airflowmetadata/airflowmetadata.go | 3 +
go-sdk/internal/bundle/doc.go | 9 +-
go-sdk/internal/bundle/task.go | 5 +
go-sdk/pkg/execution/metadata.go | 5 +
go-sdk/pkg/execution/metadata_test.go | 36 ++
27 files changed, 1323 insertions(+), 129 deletions(-)
diff --git
a/airflow-core/adr/lang-sdk/0015-per-dag-source-in-bundle-artifact.md
b/airflow-core/adr/lang-sdk/0015-per-dag-source-in-bundle-artifact.md
index 90c48b593f7..602066d958f 100644
--- a/airflow-core/adr/lang-sdk/0015-per-dag-source-in-bundle-artifact.md
+++ b/airflow-core/adr/lang-sdk/0015-per-dag-source-in-bundle-artifact.md
@@ -95,6 +95,9 @@ packers adopt the same shape — per-Dag source resolution,
de-duplication by pa
always-embedded entrypoint fallback — so a bundle reader treats every language
the same way and the
Code tab behaves identically regardless of which SDK produced the artifact.
+Go has no module evaluation, so the Go SDK records the file that calls
`airflow.Dag`. A Dag built in
+a factory function maps to the factory's file, not to the file that calls the
factory.
+
## What a reader returns
| the Task is… | reader returns |
@@ -113,11 +116,11 @@ Code tab behaves identically regardless of which SDK
produced the artifact.
`DagCode` → the Code tab, deferred to
[ADR-0010](0010-native-dag-processing.md)'s open question and
future work) sees one contract.
- Generated Dags (built in a loop or a factory during module evaluation) are
resolved to their
- generator file like any other Dag — they are fully supported, not a fallback
case. Best-effort
- resolution only bites a Dag built in a detached callback after its module
finished, which then
- shows the entrypoint — a real file in the bundle rather than a wrong one.
This mirrors the source
- view already being best effort for Python factory-function Dags
- ([ADR-0006](0006-no-lang-sdk-source-display.md), "Why Not" #3).
+ generator file (for Go, the file that calls `airflow.Dag`) like any other
Dag — they are fully
+ supported, not a fallback case. Best-effort resolution only bites a Dag
built in a detached
+ callback after its module finished, which then shows the entrypoint — a real
file in the bundle
+ rather than a wrong one. This mirrors the source view already being best
effort for Python
+ factory-function Dags ([ADR-0006](0006-no-lang-sdk-source-display.md), "Why
Not" #3).
## References
diff --git a/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
b/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
index 7648c141531..f2d7613e356 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
@@ -28,7 +28,8 @@ to a compiled Go *bundle* that is launched by
:class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator` for each
task instance.
Because Go is a compiled language, every task must be compiled ahead of time
and registered inside a single,
-self-contained native executable called a **bundle**. The bundle also embeds
its Dag source and a metadata
+self-contained native executable called a **bundle**. The bundle also embeds
the source file of each Dag built with ``airflow.Dag``
+plus the entrypoint file (the one with ``func main``), and a metadata
manifest (the ``dag_id`` and ``task_id`` map) in a footer appended to the
executable, so the executable *is*
the bundle: one runnable file to ship, with no separate manifest or archive.
The
:ref:`airflow-go-pack <go-sdk/build>` tool builds and packs that bundle.
@@ -529,7 +530,7 @@ surface as Go values when read back via ``GetXCom``.
Building and packaging
----------------------
-A plain ``go build`` produces a runnable binary, but a *deployable* bundle
(binary + embedded source +
+A plain ``go build`` produces a runnable binary, but a *deployable* bundle
(binary + embedded source files +
manifest) must be produced with ``airflow-go-pack``. The packer compiles the
bundle and appends the embedded
metadata footer, so the coordinator can read its ``dag_id``\ s without
executing the binary, producing a
single runnable file. The on-disk format the packer emits (the ``AFBNDL01``
footer and the
diff --git a/go-sdk/README.md b/go-sdk/README.md
index f099bee045e..a865e52f24b 100644
--- a/go-sdk/README.md
+++ b/go-sdk/README.md
@@ -38,7 +38,9 @@ Python tasks are imported and run in-process. Go is compiled,
so the model is di
A single binary that bundles one or more Dags' task functions is called a
**bundle**. You build one with
the SDK's packer, `airflow-go-pack`, which compiles your code and appends a
metadata footer (the manifest
-of `dag_id`s and `task_id`s, plus the Dag source) to the executable. The
result is a **self-contained
+of `dag_id`s and `task_id`s, plus the source files) to the executable. The
source files are the
+entrypoint (the file with `func main`) and the file that declares each Dag
built with `airflow.Dag`.
+A Dag declared outside your module, such as in a dependency, uses the
entrypoint as its source. The result is a **self-contained
executable bundle**: a single runnable file that *is* the bundle, with no
separate manifest or archive to
ship alongside it.
diff --git a/go-sdk/airflow/bundle.go b/go-sdk/airflow/bundle.go
index df10c5d5255..61bcd77aa8f 100644
--- a/go-sdk/airflow/bundle.go
+++ b/go-sdk/airflow/bundle.go
@@ -217,6 +217,37 @@ func (m *dagMap) add(dag *DagRef) {
m.order = append(m.order, dag)
}
+// ListDagSourceFiles returns the file that declared each registered Dag, by
dag_id.
+func (m *dagMap) ListDagSourceFiles() map[string]string {
+ m.mu.Lock()
+ defer m.mu.Unlock()
+
+ files := make(map[string]string, len(m.dags))
+ for id, dag := range m.dags {
+ files[id] = dag.file
+ }
+ return files
+}
+
+// lookupTask returns the Go function of a task of a registered Dag. A task
from TriggerDagRun has
+// none, because a Python worker runs it.
+func (m *dagMap) lookupTask(dagID, taskID string) (bundle.Task, bool) {
+ m.mu.Lock()
+ dag, ok := m.dags[dagID]
+ m.mu.Unlock()
+ if !ok {
+ return nil, false
+ }
+
+ dag.mu.Lock()
+ defer dag.mu.Unlock()
+ ref, ok := dag.tasksByID[taskID]
+ if !ok || ref.task == nil {
+ return nil, false
+ }
+ return ref.task, true
+}
+
func (m *dagMap) has(dagID string) bool {
m.mu.Lock()
defer m.mu.Unlock()
@@ -254,18 +285,32 @@ func serializeRecovering(dag *DagRef, fileloc,
relativeFileloc string) (s bundle
return s
}
-// coordinatorSource is what Serve hands to execution.Serve: the task handlers
and the Dags of one
-// bundle.
+// coordinatorSource is what Serve hands to execution.Serve and
execution.DumpAirflowMetadata: the
+// task handlers and the Dags of one bundle.
type coordinatorSource struct {
*taskHandlerMap
dags *dagMap
}
var (
- _ bundle.Bundle = coordinatorSource{}
- _ bundle.DagSerializer = coordinatorSource{}
+ _ bundle.Bundle = coordinatorSource{}
+ _ bundle.DagSerializer = coordinatorSource{}
+ _ bundle.DagSourceLister = coordinatorSource{}
)
func (s coordinatorSource) SerializeDags(fileloc, relativeFileloc string)
[]bundle.SerializedDag {
return s.dags.serialize(fileloc, relativeFileloc)
}
+
+func (s coordinatorSource) ListDagSourceFiles() map[string]string {
+ return s.dags.ListDagSourceFiles()
+}
+
+// LookupTask finds a task of a Dag from airflow.Dag, then a task handler. A
dag_id belongs to only
+// one of the two, as Register checks.
+func (s coordinatorSource) LookupTask(dagID, taskID string) (bundle.Task,
bool) {
+ if task, ok := s.dags.lookupTask(dagID, taskID); ok {
+ return task, true
+ }
+ return s.taskHandlerMap.LookupTask(dagID, taskID)
+}
diff --git a/go-sdk/airflow/bundle_test.go b/go-sdk/airflow/bundle_test.go
index d6c2da67f51..1ddc3b2b109 100644
--- a/go-sdk/airflow/bundle_test.go
+++ b/go-sdk/airflow/bundle_test.go
@@ -453,3 +453,30 @@ func TestSerializeDagsLeavesOutADagThatRegisterRejected(t
*testing.T) {
assert.Empty(t, b.dags.serialize("/bundles/go/etl", "etl"))
}
+
+func TestServeLooksUpTheTasksOfDagsAndTaskHandlers(t *testing.T) {
+ dag := Dag("native_etl")
+ dag.Task(ping, TaskSpec{TaskID: "extract"})
+ dag.Task(
+ TriggerDagRun(TriggerDagRunSpec{DagID: "downstream_etl"}),
+ TaskSpec{TaskID: "trigger"},
+ )
+ b := Bundle()
+ b.Register(dag, TaskHandler("py_etl", "load", noop))
+ source := coordinatorSource{&b.taskHandlers, &b.dags}
+
+ for _, id := range [][2]string{{"native_etl", "extract"}, {"py_etl",
"load"}} {
+ task, ok := source.LookupTask(id[0], id[1])
+ assert.True(t, ok, "%s.%s", id[0], id[1])
+ assert.NotNil(t, task)
+ }
+ for _, id := range [][2]string{
+ {"native_etl", "trigger"},
+ {"native_etl", "load"},
+ {"py_etl", "extract"},
+ {"unknown", "extract"},
+ } {
+ _, ok := source.LookupTask(id[0], id[1])
+ assert.False(t, ok, "%s.%s", id[0], id[1])
+ }
+}
diff --git a/go-sdk/airflow/dag.go b/go-sdk/airflow/dag.go
index 9415e514caf..31950d51a0c 100644
--- a/go-sdk/airflow/dag.go
+++ b/go-sdk/airflow/dag.go
@@ -33,6 +33,9 @@ import (
// DagRef is a Dag authored in Go. [Dag] returns a new one.
type DagRef struct {
dagID string
+ // file is the source file that called Dag, as the compiler recorded
it. It is empty when the
+ // runtime cannot report a caller.
+ file string
// Dag and Task copy the specs they are given with copySpec, so a
caller cannot change a
// registered Dag through a spec it still holds.
spec DagSpec
@@ -71,8 +74,11 @@ type DagRef struct {
// [IfRef.Then], [IfRef.Else], [SwitchRef.Case] and the methods of
[TaskGroupRef] panic once the Dag
// is registered.
//
-// [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.
+// The bundle embeds the source file that calls Dag, so call it from the file
that declares the Dag.
+// A Dag built in a factory function belongs to the file of the factory.
+//
+// [BundleRef.Serve] sends the registered Dags to the Dag processor and runs
their tasks, but does
+// not list them in the dags of the --airflow-metadata manifest.
//
// 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:
@@ -91,6 +97,7 @@ func Dag(dagID string, spec ...DagSpec) *DagRef {
))
}
d := &DagRef{dagID: dagID}
+ _, d.file, _, _ = runtime.Caller(1)
if len(spec) == 1 {
if err := checkDagSpec(spec[0]); err != nil {
panic(fmt.Sprintf("airflow.Dag: Dag %q: %v", dagID,
err))
diff --git a/go-sdk/internal/bundle/doc.go b/go-sdk/airflow/dag_factory_test.go
similarity index 68%
copy from go-sdk/internal/bundle/doc.go
copy to go-sdk/airflow/dag_factory_test.go
index 8a77698fe70..ca794c5e670 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/airflow/dag_factory_test.go
@@ -15,10 +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, the Dag and task ids it lists in the manifest,
and
-// the serialized Dags it sends to the Dag processor.
-//
-// Package airflow builds the tasks and ids from the task handlers a bundle
-// registers, and the serialized Dags from its Dags.
-package bundle
+package airflow
+
+// newFactoryDag builds a Dag in a file other than the one that calls it, so a
test can check
+// that the Dag belongs to this file.
+func newFactoryDag(dagID string) *DagRef {
+ return Dag(dagID)
+}
diff --git a/go-sdk/airflow/dag_source_test.go
b/go-sdk/airflow/dag_source_test.go
new file mode 100644
index 00000000000..662a1f9b86a
--- /dev/null
+++ b/go-sdk/airflow/dag_source_test.go
@@ -0,0 +1,65 @@
+// 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 airflow
+
+import (
+ "bytes"
+ "path/filepath"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+ "gopkg.in/yaml.v3"
+
+ "github.com/apache/airflow/go-sdk/internal/airflowmetadata"
+)
+
+func TestDagRecordsCallerFile(t *testing.T) {
+ assert.Equal(t, "dag_source_test.go", filepath.Base(Dag("direct").file))
+}
+
+func TestDagFromFactoryRecordsFactoryFile(t *testing.T) {
+ assert.Equal(t, "dag_factory_test.go",
filepath.Base(newFactoryDag("built").file))
+}
+
+func TestServeListsDagSourceFiles(t *testing.T) {
+ b := Bundle()
+ b.Register(
+ TaskHandler("py_etl", "extract", noop),
+ Dag("direct"),
+ newFactoryDag("built"),
+ )
+
+ var stdout bytes.Buffer
+ require.NoError(t, b.serve([]string{"--airflow-metadata"}, &stdout))
+
+ var got airflowmetadata.Manifest
+ require.NoError(t, yaml.Unmarshal(stdout.Bytes(), &got))
+ require.Len(t, got.DagSourceFiles, 2)
+ assert.Equal(t, "dag_source_test.go",
filepath.Base(got.DagSourceFiles["direct"]))
+ assert.Equal(t, "dag_factory_test.go",
filepath.Base(got.DagSourceFiles["built"]))
+ assert.NotContains(t, got.DagSourceFiles, "py_etl")
+ assert.NotContains(t, got.Dags, "direct", "native Dags stay out of
dags")
+}
+
+func TestServeOmitsDagSourceFilesWithoutNativeDags(t *testing.T) {
+ var stdout bytes.Buffer
+ require.NoError(t, etlBundle().serve([]string{"--airflow-metadata"},
&stdout))
+
+ assert.NotContains(t, stdout.String(), "dag_source_files")
+}
diff --git a/go-sdk/airflow/serve.go b/go-sdk/airflow/serve.go
index 7fafd795965..a10fa43a910 100644
--- a/go-sdk/airflow/serve.go
+++ b/go-sdk/airflow/serve.go
@@ -140,7 +140,11 @@ func (b *BundleRef) serve(args []string, stdout io.Writer)
error {
if err != nil {
return err
}
- return execution.DumpAirflowMetadata(stdout, &b.taskHandlers,
format)
+ return execution.DumpAirflowMetadata(
+ stdout,
+ coordinatorSource{&b.taskHandlers, &b.dags},
+ format,
+ )
case modeCoordinator:
return execution.Serve(coordinatorSource{&b.taskHandlers,
&b.dags}, *commAddr, *logsAddr)
case modeCoordinatorUsageError:
diff --git a/go-sdk/cmd/airflow-go-pack/inspect.go
b/go-sdk/cmd/airflow-go-pack/inspect.go
index f9cca6d488b..9b0eb192ae7 100644
--- a/go-sdk/cmd/airflow-go-pack/inspect.go
+++ b/go-sdk/cmd/airflow-go-pack/inspect.go
@@ -21,6 +21,7 @@ import (
"fmt"
"github.com/spf13/cobra"
+ "gopkg.in/yaml.v3"
"github.com/apache/airflow/go-sdk/internal/bundlefooter"
)
@@ -29,7 +30,7 @@ func newInspectCmd() *cobra.Command {
var showSource bool
cmd := &cobra.Command{
Use: "inspect <bundle>",
- Short: "Print the manifest (and optionally source) embedded in
a bundle",
+ Short: "Print the manifest (and optionally source files)
embedded in a bundle",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
source, manifest, err := bundlefooter.Read(args[0])
@@ -38,10 +39,16 @@ func newInspectCmd() *cobra.Command {
}
out := cmd.OutOrStdout()
if showSource {
- fmt.Fprintln(out, "# --- source ---")
- out.Write(source)
- if len(source) > 0 && source[len(source)-1] !=
'\n' {
- fmt.Fprintln(out)
+ files, err := embeddedSources(source, manifest)
+ if err != nil {
+ return err
+ }
+ for _, f := range files {
+ fmt.Fprintf(out, "# --- source: %s
---\n", f.path)
+ out.Write(f.data)
+ if len(f.data) > 0 &&
f.data[len(f.data)-1] != '\n' {
+ fmt.Fprintln(out)
+ }
}
fmt.Fprintln(out, "# --- manifest ---")
}
@@ -52,6 +59,46 @@ func newInspectCmd() *cobra.Command {
return nil
},
}
- cmd.Flags().BoolVar(&showSource, "source", false, "also print the
embedded source file")
+ cmd.Flags().BoolVar(&showSource, "source", false, "also print each
embedded source file")
return cmd
}
+
+type embeddedSource struct {
+ path string
+ data []byte
+}
+
+// embeddedSources cuts the source region into the files the manifest's
sources index lists.
+func embeddedSources(region, manifest []byte) ([]embeddedSource, error) {
+ var index struct {
+ Sources []struct {
+ Path string `yaml:"path"`
+ Offset int `yaml:"offset"`
+ Length int `yaml:"length"`
+ } `yaml:"sources"`
+ }
+ if err := yaml.Unmarshal(manifest, &index); err != nil {
+ return nil, fmt.Errorf("decoding manifest: %w", err)
+ }
+ if len(index.Sources) == 0 && len(region) > 0 {
+ return nil, fmt.Errorf(
+ "the manifest lists no sources for the %d-byte source
region; "+
+ "repack the bundle with this airflow-go-pack",
+ len(region),
+ )
+ }
+ files := make([]embeddedSource, 0, len(index.Sources))
+ for _, src := range index.Sources {
+ if src.Offset < 0 || src.Length < 0 || src.Offset >
len(region)-src.Length {
+ return nil, fmt.Errorf(
+ "manifest source %q (offset %d, length %d) does
not fit the %d-byte source region",
+ src.Path, src.Offset, src.Length, len(region),
+ )
+ }
+ files = append(
+ files,
+ embeddedSource{path: src.Path, data: region[src.Offset
: src.Offset+src.Length]},
+ )
+ }
+ return files, nil
+}
diff --git a/go-sdk/cmd/airflow-go-pack/inspect_test.go
b/go-sdk/cmd/airflow-go-pack/inspect_test.go
index 8208136e801..696698950c6 100644
--- a/go-sdk/cmd/airflow-go-pack/inspect_test.go
+++ b/go-sdk/cmd/airflow-go-pack/inspect_test.go
@@ -20,7 +20,9 @@ package main
import (
"bytes"
"os"
+ "os/exec"
"path/filepath"
+ "strconv"
"testing"
"github.com/stretchr/testify/assert"
@@ -28,21 +30,30 @@ import (
)
// inspect reads a bundle through bundlefooter.Read and prints the embedded
-// manifest, prefixing the source too under --source.
+// manifest, prefixing each embedded source file too under --source.
func TestInspectCmd(t *testing.T) {
dir := t.TempDir()
exe := filepath.Join(dir, "input-bin")
require.NoError(t, os.WriteFile(exe, []byte("binary-bytes"), 0o755))
- source := []byte("package main\n\nfunc main() {}\n")
+ entry := "package main\n\nfunc main() {}\n"
+ dag := "package dags"
manifest := []byte(
"airflow_bundle_metadata_version: \"1.0\"\n" +
+ "entrypoint_path: \"cmd/main.go\"\n" +
+ "sources:\n" +
+ " - path: \"cmd/main.go\"\n" +
+ " offset: 0\n" +
+ " length: " + strconv.Itoa(len(entry)) + "\n" +
+ " - path: \"dags/etl.go\"\n" +
+ " offset: " + strconv.Itoa(len(entry)) + "\n" +
+ " length: " + strconv.Itoa(len(dag)) + "\n" +
"dags:\n" +
" my_dag:\n" +
" tasks:\n" +
" - \"t1\"\n",
)
bundle := filepath.Join(dir, "bundle")
- require.NoError(t, writeBundle(exe, bundle, source, manifest))
+ require.NoError(t, writeBundle(exe, bundle, []byte(entry+dag),
manifest))
for _, tc := range []struct {
name string
@@ -57,7 +68,8 @@ func TestInspectCmd(t *testing.T) {
{
name: "with source",
args: []string{"--source", bundle},
- expect: "# --- source ---\n" + string(source) +
+ expect: "# --- source: cmd/main.go ---\n" + entry +
+ "# --- source: dags/etl.go ---\n" + dag + "\n" +
"# --- manifest ---\n" + string(manifest),
},
} {
@@ -72,3 +84,72 @@ func TestInspectCmd(t *testing.T) {
})
}
}
+
+func TestEmbeddedSources_RejectsIndexOutsideRegion(t *testing.T) {
+ for _, tc := range []struct {
+ name string
+ region string
+ manifest string
+ wantErr string
+ }{
+ {
+ name: "past the end",
+ region: "abc",
+ manifest: "sources:\n - {path: a.go, offset: 2,
length: 5}\n",
+ wantErr: `source "a.go" (offset 2, length 5) does not
fit the 3-byte source region`,
+ },
+ {
+ name: "negative offset",
+ region: "abc",
+ manifest: "sources:\n - {path: a.go, offset: -1,
length: 2}\n",
+ wantErr: `source "a.go"`,
+ },
+ {
+ name: "region without an index",
+ region: "abc",
+ manifest: "dags: {}\n",
+ wantErr: "repack the bundle",
+ },
+ } {
+ t.Run(tc.name, func(t *testing.T) {
+ _, err := embeddedSources([]byte(tc.region),
[]byte(tc.manifest))
+ require.Error(t, err)
+ assert.Contains(t, err.Error(), tc.wantErr)
+ })
+ }
+
+ files, err := embeddedSources(nil, []byte("dags: {}\n"))
+ require.NoError(t, err)
+ assert.Empty(t, files)
+}
+
+// Inspect prints every file that packing the multidag fixture embedded.
+func TestInspectCmd_PrintsEachPackedFile(t *testing.T) {
+ if testing.Short() {
+ t.Skip("shells out to `go build`")
+ }
+ if _, err := exec.LookPath("go"); err != nil {
+ t.Skip("go toolchain not on PATH")
+ }
+
+ out := filepath.Join(t.TempDir(), "bundle")
+ pack := newRootCmd()
+ pack.SetArgs([]string{"./testdata/multidag", "--output", out})
+ pack.SetOut(&bytes.Buffer{})
+ pack.SetErr(&bytes.Buffer{})
+ require.NoError(t, pack.Execute())
+
+ inspect := newInspectCmd()
+ var printed bytes.Buffer
+ inspect.SetOut(&printed)
+ inspect.SetArgs([]string{"--source", out})
+ require.NoError(t, inspect.Execute())
+ for _, name := range []string{"main.go", "factory/factory.go",
"reports/reports.go"} {
+ assert.Contains(
+ t,
+ printed.String(),
+ "# --- source:
cmd/airflow-go-pack/testdata/multidag/"+name+" ---\n",
+ )
+ }
+ assert.Contains(t, printed.String(), "# --- manifest ---\n")
+}
diff --git a/go-sdk/cmd/airflow-go-pack/main.go
b/go-sdk/cmd/airflow-go-pack/main.go
index cc04b315ac4..2f30ea5acab 100644
--- a/go-sdk/cmd/airflow-go-pack/main.go
+++ b/go-sdk/cmd/airflow-go-pack/main.go
@@ -51,10 +51,11 @@ func newRootCmd() *cobra.Command {
Use: "airflow-go-pack [package]",
Short: "Build a self-contained Airflow bundle from a Go
package",
Long: `airflow-go-pack builds a Go bundle binary, queries it
for its DAG/task
-identity via --airflow-metadata, and appends the source plus an
+identity via --airflow-metadata, and appends the source files plus an
airflow-metadata.yaml manifest plus an AFBNDL01 trailer to the
-executable. The result is a single self-contained file that drops into
-[executable] bundles_folder.
+executable. The source files are the entrypoint (the file with func main)
+and the file that declares each Dag built with airflow.Dag. The result is a
+single self-contained file that drops into [executable] bundles_folder.
By default the packer builds the package in the current directory. Pass
a different package as the positional argument; pass extra go build
@@ -121,7 +122,7 @@ Examples:
root.Flags().StringVar(&opts.source, "source",
"",
- "path to the DAG source file (defaults to the file in the
target package containing func main)")
+ "path to the entrypoint source file (defaults to the file in
the target package containing func main)")
root.Flags().StringVar(&opts.executable, "executable",
"",
"pack a pre-built executable instead of running go build")
diff --git a/go-sdk/cmd/airflow-go-pack/pack.go
b/go-sdk/cmd/airflow-go-pack/pack.go
index f4be6fa8c45..c09cbe87e1e 100644
--- a/go-sdk/cmd/airflow-go-pack/pack.go
+++ b/go-sdk/cmd/airflow-go-pack/pack.go
@@ -30,6 +30,7 @@ import (
"path/filepath"
"runtime"
"sort"
+ "strconv"
"strings"
"gopkg.in/yaml.v3"
@@ -41,7 +42,7 @@ import (
// packOptions are the flags accepted by the root pack command.
type packOptions struct {
pkg string // target package (default ".")
- source string // override the auto-detected DAG source file
+ source string // override the auto-detected entrypoint
source file
executable string // pack a pre-built binary instead of building
output string // override the default <bundleName> output
path
airflowMetadata string // path to a pre-captured --airflow-metadata
manifest (JSON or YAML)
@@ -62,13 +63,13 @@ func runPack(stdout, stderr io.Writer, opts *packOptions)
error {
)
}
- // Resolve the DAG source file for both modes up front. --executable
requires
+ // Resolve the entrypoint source file for both modes up front.
--executable requires
// it explicitly; the build path falls back to discovery.
sourcePath := opts.source
if opts.executable != "" {
if sourcePath == "" {
return fmt.Errorf(
- "--executable requires --source: cannot infer
the DAG source for a pre-built binary",
+ "--executable requires --source: cannot infer
the entrypoint source for a pre-built binary",
)
}
} else if sourcePath == "" {
@@ -80,7 +81,7 @@ func runPack(stdout, stderr io.Writer, opts *packOptions)
error {
// when --source was not supplied, so an explicit --source
always wins.
discovered, err := discoverMainSource(opts.pkg)
if err != nil {
- return fmt.Errorf("locating DAG source file: %w", err)
+ return fmt.Errorf("locating the entrypoint source file:
%w", err)
}
sourcePath = discovered
}
@@ -155,7 +156,7 @@ func runPack(stdout, stderr io.Writer, opts *packOptions)
error {
return fmt.Errorf("executable %s: %w", execPath, err)
}
- if err := rejectOutputAlias(output, execPath, sourcePath,
opts.airflowMetadata); err != nil {
+ if err := rejectOutputAlias(output, execPath, []string{sourcePath},
opts.airflowMetadata); err != nil {
return err
}
@@ -163,7 +164,7 @@ func runPack(stdout, stderr io.Writer, opts *packOptions)
error {
if err != nil {
return err
}
- if len(meta.Dags) == 0 {
+ if len(meta.Dags) == 0 && len(meta.DagSourceFiles) == 0 {
return fmt.Errorf("bundle exposes no dags: nothing to pack")
}
for dagID, dag := range meta.Dags {
@@ -173,28 +174,46 @@ func runPack(stdout, stderr io.Writer, opts *packOptions)
error {
}
warnOnSuspiciousIDs(stderr, meta)
- manifest, err := renderManifest(meta, filepath.Base(sourcePath))
+ mod, err := moduleOf(opts, sourcePath)
if err != nil {
- return fmt.Errorf("rendering manifest: %w", err)
+ return err
+ }
+ layout, region, err := layoutSources(stderr, meta, sourcePath, mod)
+ if err != nil {
+ return err
+ }
+ if err := rejectOutputAlias(output, execPath, layout.diskPaths(),
opts.airflowMetadata); err != nil {
+ return err
}
- sourceBytes, err := os.ReadFile(sourcePath)
+ manifest, err := renderManifest(meta, layout)
if err != nil {
- return fmt.Errorf("reading source file: %w", err)
+ return fmt.Errorf("rendering manifest: %w", err)
}
// Assemble the bundle through a temp file and atomically move it into
// place: we never mutate the build artefact or the user-supplied
// --executable, and a failed pack never leaves a truncated or
half-written
// file at output.
- if err := writeBundle(execPath, output, sourceBytes, manifest); err !=
nil {
+ if err := writeBundle(execPath, output, region, manifest); err != nil {
return err
}
- fmt.Fprintf(stdout, "Wrote bundle %s (sdk=%s/%s, dags=%d)\n",
- output, meta.SDK.Language, meta.SDK.Version, len(meta.Dags))
+ fmt.Fprintf(stdout, "Wrote bundle %s (sdk=%s/%s, task_handler_dags=%d,
native_dags=%d)\n",
+ output, meta.SDK.Language, meta.SDK.Version, len(meta.Dags),
len(meta.DagSourceFiles))
return nil
}
+// moduleOf finds the module that owns the entrypoint. The build path asks the
go tool, which
+// knows about workspaces and vendoring. --executable has no package to ask
about, so it reads
+// the nearest go.mod above the source file.
+func moduleOf(opts *packOptions, sourcePath string) (goModule, error) {
+ dir := filepath.Dir(sourcePath)
+ if opts.executable != "" {
+ return findModule(dir)
+ }
+ return listModule(dir)
+}
+
// defaultOutputPath derives the default bundle output path from the directory
// that owns the DAG source file. That directory is the bundle's main package,
// so its base name is what `go build` itself would name the binary. On Windows
@@ -463,11 +482,10 @@ func runIntrospect(execPath string, flag string) ([]byte,
error) {
}
// renderManifest serialises the airflow-metadata manifest as deterministic,
-// sorted-key YAML matching airflow-metadata.schema.json. It injects the
schema's
-// source field (the filename the manifest is built from), which the producer's
-// Manifest omits because only the packer knows it; every other field is copied
-// from the introspected manifest verbatim.
-func renderManifest(meta airflowmetadata.Manifest, sourceName string) ([]byte,
error) {
+// sorted-key YAML matching airflow-metadata.schema.json. It adds the source
layout
+// (entrypoint_path, dag_source_paths and sources), which only the packer
knows; every
+// other field is copied from the introspected manifest verbatim.
+func renderManifest(meta airflowmetadata.Manifest, layout sourceLayout)
([]byte, error) {
version := meta.AirflowBundleMetadataVersion
if version == "" {
version = airflowmetadata.FormatVersion
@@ -487,7 +505,7 @@ func renderManifest(meta airflowmetadata.Manifest,
sourceName string) ([]byte, e
taskItems = append(taskItems, quotedScalar(t))
}
dagsNode.Content = append(dagsNode.Content,
- scalar(id),
+ quotedScalar(id),
&yaml.Node{
Kind: yaml.MappingNode,
Content: []*yaml.Node{
@@ -498,6 +516,33 @@ func renderManifest(meta airflowmetadata.Manifest,
sourceName string) ([]byte, e
)
}
+ dagPaths := make([]string, 0, len(layout.dagPaths))
+ for id := range layout.dagPaths {
+ dagPaths = append(dagPaths, id)
+ }
+ sort.Strings(dagPaths)
+ dagPathsNode := &yaml.Node{Kind: yaml.MappingNode}
+ for _, id := range dagPaths {
+ dagPathsNode.Content = append(
+ dagPathsNode.Content,
+ quotedScalar(id),
+ quotedScalar(layout.dagPaths[id]),
+ )
+ }
+
+ sourcesNode := &yaml.Node{Kind: yaml.SequenceNode}
+ for _, f := range layout.files {
+ sourcesNode.Content = append(sourcesNode.Content, &yaml.Node{
+ Kind: yaml.MappingNode,
+ Content: []*yaml.Node{
+ scalar("path"), quotedScalar(f.path),
+ scalar("offset"), intScalar(f.offset),
+ scalar("length"), intScalar(f.length),
+ scalar("sha256"), quotedScalar(f.sha256),
+ },
+ })
+ }
+
root := &yaml.Node{Kind: yaml.DocumentNode}
manifest := &yaml.Node{
Kind: yaml.MappingNode,
@@ -513,7 +558,9 @@ func renderManifest(meta airflowmetadata.Manifest,
sourceName string) ([]byte, e
quotedScalar(meta.SDK.SupervisorSchemaVersion),
},
},
- scalar("source"), quotedScalar(sourceName),
+ scalar("entrypoint_path"),
quotedScalar(layout.entrypoint),
+ scalar("dag_source_paths"), dagPathsNode,
+ scalar("sources"), sourcesNode,
scalar("dags"), dagsNode,
},
}
@@ -531,34 +578,40 @@ func renderManifest(meta airflowmetadata.Manifest,
sourceName string) ([]byte, e
return buf.Bytes(), nil
}
-// scalar emits a plain (unquoted) node. It is used for structural keys
-// (e.g. "sdk", "tasks") and for the Dag ID mapping keys.
+// scalar emits a plain (unquoted) node. It is for structural keys only (e.g.
"sdk", "tasks").
+// Dag IDs are data, so they go through quotedScalar.
func scalar(value string) *yaml.Node {
return &yaml.Node{Kind: yaml.ScalarNode, Value: value}
}
-// quotedScalar emits a double-quoted node. Data-bearing string *values* — task
-// IDs, the source filename, and the SDK fields — go through this so a value
-// that looks like a number, bool, or date (e.g. a task named "123" or "true")
-// round-trips as a string rather than being retyped by the YAML parser.
+func intScalar(value int) *yaml.Node {
+ return &yaml.Node{Kind: yaml.ScalarNode, Tag: "!!int", Value:
strconv.Itoa(value)}
+}
+
+// quotedScalar emits a double-quoted node. Data-bearing strings, such as Dag
IDs, task IDs,
+// the source paths, and the SDK fields, go through this so a value that looks
like a number,
+// bool, or date (e.g. a Dag named "2024" or "on") round-trips as a string
rather than being
+// retyped by the YAML parser.
func quotedScalar(value string) *yaml.Node {
return &yaml.Node{Kind: yaml.ScalarNode, Value: value, Style:
yaml.DoubleQuotedStyle}
}
// rejectOutputAlias fails if output resolves to the same file as any pack
-// input: the executable, the source, or a supplied --airflow-metadata file.
+// input: the executable, a source file, or a supplied --airflow-metadata file.
// Packing copies the executable to output with O_TRUNC and renames it into
// place, so an aliased output would clobber the input. metadataPath is empty
// when --airflow-metadata is not used and is skipped in that case.
-func rejectOutputAlias(output, execPath, sourcePath, metadataPath string)
error {
- for _, in := range []struct {
+func rejectOutputAlias(output, execPath string, sourcePaths []string,
metadataPath string) error {
+ type input struct {
path string
kind string
- }{
- {execPath, "executable"},
- {sourcePath, "source"},
- {metadataPath, "--airflow-metadata file"},
- } {
+ }
+ inputs := []input{{execPath, "executable"}}
+ for _, p := range sourcePaths {
+ inputs = append(inputs, input{p, "source"})
+ }
+ inputs = append(inputs, input{metadataPath, "--airflow-metadata file"})
+ for _, in := range inputs {
if in.path == "" {
continue
}
@@ -612,7 +665,7 @@ func sameFile(a, b string) (bool, error) {
}
// writeBundle assembles the bundle at output by copying the executable to a
-// temporary file in output's directory, appending the source+manifest footer
+// temporary file in output's directory, appending the source region and
manifest footer
// to that copy, then atomically renaming it into place. Writing through a
// temp file keeps a failed pack from leaving a truncated or half-written
// artefact at output, and guarantees the file being copied is never the same
diff --git a/go-sdk/cmd/airflow-go-pack/pack_integration_test.go
b/go-sdk/cmd/airflow-go-pack/pack_integration_test.go
index 7f5c6b4a373..019fc6c3fd2 100644
--- a/go-sdk/cmd/airflow-go-pack/pack_integration_test.go
+++ b/go-sdk/cmd/airflow-go-pack/pack_integration_test.go
@@ -26,6 +26,7 @@ import (
"path/filepath"
"regexp"
"runtime"
+ "strconv"
"testing"
"github.com/stretchr/testify/assert"
@@ -144,20 +145,26 @@ sdk:
language: "go"
version: "` + sdkVersion + `"
supervisor_schema_version: "` + execution.SupervisorSchemaVersion + `"
-source: "main.go"
+entrypoint_path: "example/bundle/main.go"
+dag_source_paths: {}
+sources:
+ - path: "example/bundle/main.go"
+ offset: 0
+ length: ` + strconv.Itoa(len(srcBytes)) + `
+ sha256: "` + sha256Hex(srcBytes) + `"
dags:
- concurrent_xcom_dag:
+ "concurrent_xcom_dag":
tasks:
- "pull_xcoms_concurrently"
- simple_dag:
+ "simple_dag":
tasks:
- "extract"
- "transform"
- "load"
- task_state_dag:
+ "task_state_dag":
tasks:
- "roundtrip_task_state"
- taskflow_binding_dag:
+ "taskflow_binding_dag":
tasks:
- "make_config"
- "make_numbers"
@@ -171,7 +178,7 @@ dags:
- "via_flat_map"
- "via_struct_map"
- "via_plain_map"
- variable_write_dag:
+ "variable_write_dag":
tasks:
- "write_and_delete_variable"
`
@@ -225,7 +232,7 @@ func TestPack_CrossCompileBuildModeForwardsFlags(t
*testing.T) {
source, metadata, err := bundlefooter.Read(outPath)
require.NoError(t, err)
- assert.Contains(t, string(metadata), "simple_dag:",
+ assert.Contains(t, string(metadata), `"simple_dag":`,
"manifest must be read from the host introspection build")
// Independently build the target-arch artefact with the same forwarded
diff --git a/go-sdk/cmd/airflow-go-pack/pack_test.go
b/go-sdk/cmd/airflow-go-pack/pack_test.go
index 999d5c0f017..992e1e6da87 100644
--- a/go-sdk/cmd/airflow-go-pack/pack_test.go
+++ b/go-sdk/cmd/airflow-go-pack/pack_test.go
@@ -18,21 +18,33 @@
package main
import (
+ "bytes"
"errors"
"io"
"os"
"path/filepath"
"runtime"
+ "strconv"
"testing"
"github.com/spf13/cobra"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
+ "gopkg.in/yaml.v3"
"github.com/apache/airflow/go-sdk/internal/airflowmetadata"
"github.com/apache/airflow/go-sdk/internal/bundlefooter"
)
+var testLayout = sourceLayout{
+ entrypoint: "cmd/bundle/main.go",
+ dagPaths: map[string]string{"zeta_dag": "cmd/bundle/main.go",
"alpha_dag": "dags/alpha.go"},
+ files: []sourceFile{
+ {path: "cmd/bundle/main.go", offset: 0, length: 10, sha256:
"aa"},
+ {path: "dags/alpha.go", offset: 10, length: 20, sha256: "bb"},
+ },
+}
+
func TestRenderManifest_DeterministicDagOrdering(t *testing.T) {
meta := airflowmetadata.Manifest{
AirflowBundleMetadataVersion: "1.0",
@@ -47,9 +59,9 @@ func TestRenderManifest_DeterministicDagOrdering(t
*testing.T) {
},
}
- got1, err := renderManifest(meta, "main.go")
+ got1, err := renderManifest(meta, testLayout)
require.NoError(t, err)
- got2, err := renderManifest(meta, "main.go")
+ got2, err := renderManifest(meta, testLayout)
require.NoError(t, err)
assert.Equal(t, got1, got2, "manifest should be byte-identical for
identical input")
@@ -59,12 +71,24 @@ sdk:
language: "go"
version: "0.1.0"
supervisor_schema_version: "2026-06-16"
-source: "main.go"
+entrypoint_path: "cmd/bundle/main.go"
+dag_source_paths:
+ "alpha_dag": "dags/alpha.go"
+ "zeta_dag": "cmd/bundle/main.go"
+sources:
+ - path: "cmd/bundle/main.go"
+ offset: 0
+ length: 10
+ sha256: "aa"
+ - path: "dags/alpha.go"
+ offset: 10
+ length: 20
+ sha256: "bb"
dags:
- alpha_dag:
+ "alpha_dag":
tasks:
- "x"
- zeta_dag:
+ "zeta_dag":
tasks:
- "a"
- "b"
@@ -72,9 +96,9 @@ dags:
assert.Equal(t, expected, string(got1))
}
-// Values (task IDs, source, SDK fields) are quoted so a scalar-looking value
-// stays a string; Dag ID keys stay plain scalars.
-func TestRenderManifest_QuotesValuesNotKeys(t *testing.T) {
+// Values (task IDs, source paths, SDK fields) and Dag ID keys are quoted so a
scalar-looking
+// string stays a string when a YAML parser reads the manifest back.
+func TestRenderManifest_QuotesDagIDs(t *testing.T) {
meta := airflowmetadata.Manifest{
AirflowBundleMetadataVersion: "1.0",
SDK: airflowmetadata.SDK{
@@ -83,19 +107,34 @@ func TestRenderManifest_QuotesValuesNotKeys(t *testing.T) {
SupervisorSchemaVersion: "2026-06-16",
},
Dags: map[string]airflowmetadata.Dag{
- "my_dag": {Tasks: []string{"123", "true"}},
+ "2024": {Tasks: []string{"123", "true"}},
+ "on": {Tasks: []string{"t1"}},
},
}
+ layout := sourceLayout{
+ entrypoint: "main.go",
+ dagPaths: map[string]string{"2024": "main.go", "on":
"main.go"},
+ files: []sourceFile{{path: "main.go", offset: 0, length:
10, sha256: "aa"}},
+ }
- got, err := renderManifest(meta, "main.go")
+ got, err := renderManifest(meta, layout)
require.NoError(t, err)
- // Task values that look like scalars are quoted.
assert.Contains(t, string(got), `- "123"`)
assert.Contains(t, string(got), `- "true"`)
- // The Dag ID key is a plain scalar, not quoted.
- assert.Contains(t, string(got), "\n my_dag:\n")
- assert.NotContains(t, string(got), `"my_dag"`)
+ assert.Contains(t, string(got), "\n \"2024\":\n")
+ assert.Contains(t, string(got), "\n \"on\":\n")
+ assert.Contains(t, string(got), `"2024": "main.go"`)
+ assert.Contains(t, string(got), `"on": "main.go"`)
+
+ var decoded struct {
+ Dags map[any]any `yaml:"dags"`
+ DagSourcePaths map[any]any `yaml:"dag_source_paths"`
+ }
+ require.NoError(t, yaml.Unmarshal(got, &decoded))
+ assert.Contains(t, decoded.Dags, "2024")
+ assert.Contains(t, decoded.Dags, "on")
+ assert.Equal(t, map[any]any{"2024": "main.go", "on": "main.go"},
decoded.DagSourcePaths)
}
func TestRenderManifest_EmptyDags(t *testing.T) {
@@ -108,7 +147,7 @@ func TestRenderManifest_EmptyDags(t *testing.T) {
},
Dags: map[string]airflowmetadata.Dag{},
}
- got, err := renderManifest(meta, "main.go")
+ got, err := renderManifest(meta, testLayout)
require.NoError(t, err)
assert.Contains(t, string(got), "dags: {}")
}
@@ -367,7 +406,7 @@ func TestRunPack_UsesMetadataFile(t *testing.T) {
assert.Contains(
t,
string(gotMeta),
- "my_dag:",
+ `"my_dag":`,
"Dag from --airflow-metadata must appear in the manifest",
)
@@ -377,6 +416,35 @@ func TestRunPack_UsesMetadataFile(t *testing.T) {
assert.Equal(t, exeBytes, binaryRegion, "the supplied --executable must
be packed verbatim")
}
+func TestRunPack_PacksABundleWithOnlyNativeDags(t *testing.T) {
+ dir := t.TempDir()
+ exe := filepath.Join(dir, "native")
+ require.NoError(t, os.WriteFile(exe, []byte("native-binary-bytes"),
0o755))
+ source := filepath.Join(dir, "main.go")
+ require.NoError(t, os.WriteFile(source, []byte("package main\nfunc
main() {}\n"), 0o644))
+ meta := filepath.Join(dir, "airflow-metadata.json")
+ require.NoError(t, os.WriteFile(meta, []byte(
+ `{"airflow_bundle_metadata_version":"1.0",`+
+
`"sdk":{"language":"go","version":"0.1.0","supervisor_schema_version":"2026-06-16"},`+
+
`"dags":{},"dag_source_files":{"native_dag":`+strconv.Quote(source)+`}}`,
+ ), 0o644))
+ out := filepath.Join(dir, "bundle")
+
+ var stdout bytes.Buffer
+ err := runPack(&stdout, io.Discard, &packOptions{
+ executable: exe,
+ source: source,
+ airflowMetadata: meta,
+ output: out,
+ })
+ require.NoError(t, err)
+
+ _, gotMeta, err := bundlefooter.Read(out)
+ require.NoError(t, err)
+ assert.Contains(t, string(gotMeta), `"native_dag": "main.go"`)
+ assert.Contains(t, stdout.String(), "task_handler_dags=0,
native_dags=1")
+}
+
// --airflow-metadata also accepts a YAML manifest, not only the JSON the
// binary prints.
func TestRunPack_AcceptsYAMLMetadataFile(t *testing.T) {
@@ -412,7 +480,7 @@ func TestRunPack_AcceptsYAMLMetadataFile(t *testing.T) {
assert.Contains(
t,
string(gotMeta),
- "yaml_dag:",
+ `"yaml_dag":`,
"Dag from a YAML --airflow-metadata file must appear in the
manifest",
)
}
diff --git a/go-sdk/cmd/airflow-go-pack/sources.go
b/go-sdk/cmd/airflow-go-pack/sources.go
new file mode 100644
index 00000000000..3455b73cb17
--- /dev/null
+++ b/go-sdk/cmd/airflow-go-pack/sources.go
@@ -0,0 +1,256 @@
+// 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 main
+
+import (
+ "bytes"
+ "crypto/sha256"
+ "encoding/hex"
+ "fmt"
+ "io"
+ "os"
+ "os/exec"
+ "path"
+ "path/filepath"
+ "sort"
+ "strings"
+
+ "github.com/apache/airflow/go-sdk/internal/airflowmetadata"
+ "github.com/apache/airflow/go-sdk/internal/bundlefooter"
+)
+
+// goModule is the module that owns the bundle's entrypoint.
+type goModule struct {
+ dir string // module root on disk
+ path string // module path from go.mod; empty when no go.mod was found
+}
+
+// sourceFile is one file embedded in the source region.
+type sourceFile struct {
+ path string // slash path relative to the module root
+ disk string // where the packer read the file from
+ offset int
+ length int
+ sha256 string // lowercase hex
+}
+
+// sourceLayout describes the source region: which file is the entrypoint,
which file each
+// native Dag came from, and where every file sits in the region.
+type sourceLayout struct {
+ entrypoint string
+ dagPaths map[string]string
+ files []sourceFile
+}
+
+// diskPaths lists the files the packer read, so the caller can keep the
output away from them.
+func (l sourceLayout) diskPaths() []string {
+ paths := make([]string, len(l.files))
+ for i, f := range l.files {
+ paths[i] = f.disk
+ }
+ return paths
+}
+
+// listModule asks the go tool for the module that owns dir. In a workspace it
lists every
+// workspace module, so it keeps the innermost one that contains dir.
+func listModule(dir string) (goModule, error) {
+ cmd := exec.Command("go", "list", "-m", "-f", "{{.Dir}}\n{{.Path}}")
+ cmd.Dir = dir
+ var stdout, stderr bytes.Buffer
+ cmd.Stdout = &stdout
+ cmd.Stderr = &stderr
+ if err := cmd.Run(); err != nil {
+ return goModule{}, fmt.Errorf(
+ "go list -m in %s: %w: %s", dir, err,
strings.TrimSpace(stderr.String()),
+ )
+ }
+ absDir, err := filepath.Abs(dir)
+ if err != nil {
+ return goModule{}, err
+ }
+ lines := splitNonEmpty(stdout.String())
+ var best goModule
+ for i := 0; i+1 < len(lines); i += 2 {
+ _, inside := relativeTo(lines[i], absDir)
+ if (inside || lines[i] == absDir) && len(lines[i]) >
len(best.dir) {
+ best = goModule{dir: lines[i], path: lines[i+1]}
+ }
+ }
+ if best.dir == "" {
+ return goModule{}, fmt.Errorf("go list -m in %s: no module
contains the directory", dir)
+ }
+ return best, nil
+}
+
+// findModule walks up from dir to the nearest go.mod and reads its module
path. It is for
+// --executable, where there is no package to hand to the go tool. Without a
go.mod the module
+// is dir itself with no path, which still resolves absolute paths under it.
+func findModule(dir string) (goModule, error) {
+ abs, err := filepath.Abs(dir)
+ if err != nil {
+ return goModule{}, err
+ }
+ for cur := abs; ; cur = filepath.Dir(cur) {
+ data, err := os.ReadFile(filepath.Join(cur, "go.mod"))
+ if err == nil {
+ return goModule{dir: cur, path: modulePath(data)}, nil
+ }
+ if !os.IsNotExist(err) {
+ return goModule{}, err
+ }
+ if filepath.Dir(cur) == cur {
+ return goModule{dir: abs}, nil
+ }
+ }
+}
+
+// modulePath returns the path on the module line of a go.mod file, or "" if
there is none.
+func modulePath(goMod []byte) string {
+ for line := range strings.SplitSeq(string(goMod), "\n") {
+ line, _, _ = strings.Cut(line, "//")
+ fields := strings.Fields(line)
+ if len(fields) == 2 && fields[0] == "module" {
+ return strings.Trim(fields[1], "\"`")
+ }
+ }
+ return ""
+}
+
+// relativeTo returns the slash path of file relative to dir. It reports false
if file is not
+// inside dir.
+func relativeTo(dir, file string) (string, bool) {
+ rel, err := filepath.Rel(dir, file)
+ if err != nil {
+ return "", false
+ }
+ rel = filepath.ToSlash(rel)
+ if rel == "." || rel == ".." || strings.HasPrefix(rel, "../") {
+ return "", false
+ }
+ return rel, true
+}
+
+// moduleRelPath turns a source path reported by the bundle binary into a
slash path relative
+// to the module root. The binary reports an absolute build path, or "<module
path>/<rel>"
+// when it was built with -trimpath. It reports false for a file that is not a
regular file
+// inside the module, such as one from the module cache or another module.
+func moduleRelPath(reported string, mod goModule) (string, bool) {
+ var rel string
+ switch {
+ case reported == "":
+ return "", false
+ case filepath.IsAbs(filepath.FromSlash(reported)):
+ var ok bool
+ if rel, ok = relativeTo(mod.dir, filepath.FromSlash(reported));
!ok {
+ return "", false
+ }
+ case mod.path != "" && strings.HasPrefix(reported, mod.path+"/"):
+ rel = path.Clean(strings.TrimPrefix(reported, mod.path+"/"))
+ if rel == ".." || strings.HasPrefix(rel, "../") {
+ return "", false
+ }
+ default:
+ return "", false
+ }
+ info, err := os.Stat(filepath.Join(mod.dir, filepath.FromSlash(rel)))
+ if err != nil || !info.Mode().IsRegular() {
+ return "", false
+ }
+ return rel, true
+}
+
+// layoutSources decides which files the bundle embeds and what each native
Dag maps to. The
+// entrypoint always comes first, and every other file follows in sorted
order, so identical
+// inputs give an identical region. A Dag whose file lies outside the module
maps to the
+// entrypoint, with a warning. It returns the layout and the concatenated
region.
+func layoutSources(
+ stderr io.Writer,
+ meta airflowmetadata.Manifest,
+ entrypoint string,
+ mod goModule,
+) (sourceLayout, []byte, error) {
+ absEntry, err := filepath.Abs(entrypoint)
+ if err != nil {
+ return sourceLayout{}, nil, err
+ }
+ entryPath, ok := relativeTo(mod.dir, absEntry)
+ if !ok {
+ entryPath = filepath.Base(absEntry)
+ }
+
+ layout := sourceLayout{
+ entrypoint: entryPath,
+ dagPaths: make(map[string]string, len(meta.DagSourceFiles)),
+ }
+ disk := map[string]string{entryPath: absEntry}
+
+ dagIDs := make([]string, 0, len(meta.DagSourceFiles))
+ for id := range meta.DagSourceFiles {
+ dagIDs = append(dagIDs, id)
+ }
+ sort.Strings(dagIDs)
+ for _, id := range dagIDs {
+ reported := meta.DagSourceFiles[id]
+ rel, ok := moduleRelPath(reported, mod)
+ if !ok {
+ fmt.Fprintf(stderr,
+ "warning: source file %q of dag %q is outside
module %s; "+
+ "embedding the entrypoint %s for it
instead\n",
+ reported, id, mod.dir, entryPath)
+ layout.dagPaths[id] = entryPath
+ continue
+ }
+ layout.dagPaths[id] = rel
+ if _, seen := disk[rel]; !seen {
+ disk[rel] = filepath.Join(mod.dir,
filepath.FromSlash(rel))
+ }
+ }
+
+ paths := make([]string, 0, len(disk))
+ for p := range disk {
+ if p != entryPath {
+ paths = append(paths, p)
+ }
+ }
+ sort.Strings(paths)
+ paths = append([]string{entryPath}, paths...)
+
+ var region bytes.Buffer
+ for _, p := range paths {
+ data, err := os.ReadFile(disk[p])
+ if err != nil {
+ return sourceLayout{}, nil, fmt.Errorf("reading source
file: %w", err)
+ }
+ sum := sha256.Sum256(data)
+ layout.files = append(layout.files, sourceFile{
+ path: p,
+ disk: disk[p],
+ offset: region.Len(),
+ length: len(data),
+ sha256: hex.EncodeToString(sum[:]),
+ })
+ region.Write(data)
+ if int64(region.Len()) > bundlefooter.MaxRegionSize {
+ return sourceLayout{}, nil, fmt.Errorf(
+ "source region too large after embedding %s
(max %d bytes)",
+ p, int64(bundlefooter.MaxRegionSize),
+ )
+ }
+ }
+ return layout, region.Bytes(), nil
+}
diff --git a/go-sdk/cmd/airflow-go-pack/sources_test.go
b/go-sdk/cmd/airflow-go-pack/sources_test.go
new file mode 100644
index 00000000000..3abe07c8b94
--- /dev/null
+++ b/go-sdk/cmd/airflow-go-pack/sources_test.go
@@ -0,0 +1,430 @@
+// 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 main
+
+import (
+ "bytes"
+ "crypto/sha256"
+ "encoding/hex"
+ "io"
+ "os"
+ "os/exec"
+ "path/filepath"
+ "runtime"
+ "strconv"
+ "testing"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+ "gopkg.in/yaml.v3"
+
+ "github.com/apache/airflow/go-sdk/internal/airflowmetadata"
+ "github.com/apache/airflow/go-sdk/internal/bundlefooter"
+)
+
+func sha256Hex(b []byte) string {
+ sum := sha256.Sum256(b)
+ return hex.EncodeToString(sum[:])
+}
+
+// packedManifest is the part of a packed manifest that the source layout adds.
+type packedManifest struct {
+ EntrypointPath string `yaml:"entrypoint_path"`
+ DagSourcePaths map[string]string `yaml:"dag_source_paths"`
+ Sources []struct {
+ Path string `yaml:"path"`
+ Offset int `yaml:"offset"`
+ Length int `yaml:"length"`
+ SHA256 string `yaml:"sha256"`
+ } `yaml:"sources"`
+}
+
+func readPacked(t *testing.T, bundle string) (packedManifest, []byte, []byte) {
+ t.Helper()
+ region, metadata, err := bundlefooter.Read(bundle)
+ require.NoError(t, err)
+ var m packedManifest
+ require.NoError(t, yaml.Unmarshal(metadata, &m))
+ return m, region, metadata
+}
+
+// writeModule lays out files under a new directory with a go.mod, and returns
the directory.
+func writeModule(t *testing.T, modulePath string, files map[string]string)
string {
+ t.Helper()
+ dir := t.TempDir()
+ require.NoError(t, os.WriteFile(filepath.Join(dir, "go.mod"),
+ []byte("module "+modulePath+"\n\ngo 1.24\n"), 0o644))
+ for rel, content := range files {
+ full := filepath.Join(dir, filepath.FromSlash(rel))
+ require.NoError(t, os.MkdirAll(filepath.Dir(full), 0o755))
+ require.NoError(t, os.WriteFile(full, []byte(content), 0o644))
+ }
+ return dir
+}
+
+func TestListModule_EntrypointAtModuleRoot(t *testing.T) {
+ if _, err := exec.LookPath("go"); err != nil {
+ t.Skip("go toolchain not on PATH")
+ }
+ t.Setenv("GOWORK", "off")
+ dir := writeModule(t, "example.com/app", map[string]string{
+ "main.go": "package main\nfunc main() {}\n",
+ })
+
+ mod, err := listModule(dir)
+ require.NoError(t, err)
+ assert.Equal(t, "example.com/app", mod.path)
+}
+
+func TestModulePath(t *testing.T) {
+ for _, tc := range []struct {
+ name string
+ mod string
+ want string
+ }{
+ {name: "plain", mod: "module example.com/app\n\ngo 1.24\n",
want: "example.com/app"},
+ {name: "quoted", mod: "module \"example.com/app\"\n", want:
"example.com/app"},
+ {name: "comment", mod: "// header\nmodule example.com/app //
trailing\n", want: "example.com/app"},
+ {name: "missing", mod: "go 1.24\n", want: ""},
+ } {
+ t.Run(tc.name, func(t *testing.T) {
+ assert.Equal(t, tc.want, modulePath([]byte(tc.mod)))
+ })
+ }
+}
+
+func TestFindModule(t *testing.T) {
+ dir := writeModule(
+ t,
+ "example.com/app",
+ map[string]string{"cmd/bundle/main.go": "package main\n"},
+ )
+
+ mod, err := findModule(filepath.Join(dir, "cmd", "bundle"))
+ require.NoError(t, err)
+ assert.Equal(t, goModule{dir: dir, path: "example.com/app"}, mod)
+
+ bare := t.TempDir()
+ mod, err = findModule(bare)
+ require.NoError(t, err)
+ assert.Equal(t, goModule{dir: bare}, mod, "without a go.mod the
directory is the module root")
+}
+
+func TestModuleRelPath(t *testing.T) {
+ dir := writeModule(t, "example.com/app", map[string]string{
+ "main.go": "package main\n",
+ "dags/etl.go": "package dags\n",
+ "sub/nested.go": "package sub\n",
+ })
+ mod := goModule{dir: dir, path: "example.com/app"}
+
+ for _, tc := range []struct {
+ name string
+ reported string
+ want string
+ ok bool
+ }{
+ {name: "absolute build path", reported:
filepath.ToSlash(filepath.Join(dir, "dags", "etl.go")), want: "dags/etl.go",
ok: true},
+ {name: "trimpath", reported: "example.com/app/dags/etl.go",
want: "dags/etl.go", ok: true},
+ {name: "trimpath at root", reported: "example.com/app/main.go",
want: "main.go", ok: true},
+ {name: "absolute path outside the module", reported:
"/home/u/go/pkg/mod/example.com/[email protected]/x.go"},
+ {name: "another module", reported:
"example.com/[email protected]/x.go"},
+ {name: "module path prefix of another module", reported:
"example.com/app/[email protected]/x.go"},
+ {name: "missing file", reported: "example.com/app/gone.go"},
+ {name: "escapes the module", reported:
"example.com/app/../x.go"},
+ {name: "directory", reported: "example.com/app/dags"},
+ {name: "empty", reported: ""},
+ } {
+ t.Run(tc.name, func(t *testing.T) {
+ got, ok := moduleRelPath(tc.reported, mod)
+ assert.Equal(t, tc.ok, ok)
+ assert.Equal(t, tc.want, got)
+ })
+ }
+}
+
+func TestLayoutSources(t *testing.T) {
+ dir := writeModule(t, "example.com/app", map[string]string{
+ "cmd/bundle/main.go": "package main\n",
+ "dags/b.go": "package dags // b\n",
+ "dags/a.go": "package dags // a\n",
+ })
+ mod := goModule{dir: dir, path: "example.com/app"}
+ meta := airflowmetadata.Manifest{DagSourceFiles: map[string]string{
+ "on_b": "example.com/app/dags/b.go",
+ "also_b": filepath.ToSlash(filepath.Join(dir, "dags", "b.go")),
+ "on_a": "example.com/app/dags/a.go",
+ "on_main": "example.com/app/cmd/bundle/main.go",
+ "foreign": "example.com/[email protected]/x.go",
+ }}
+
+ var stderr bytes.Buffer
+ layout, region, err := layoutSources(
+ &stderr,
+ meta,
+ filepath.Join(dir, "cmd", "bundle", "main.go"),
+ mod,
+ )
+ require.NoError(t, err)
+
+ assert.Equal(t, "cmd/bundle/main.go", layout.entrypoint)
+ assert.Equal(t, map[string]string{
+ "on_b": "dags/b.go",
+ "also_b": "dags/b.go",
+ "on_a": "dags/a.go",
+ "on_main": "cmd/bundle/main.go",
+ "foreign": "cmd/bundle/main.go",
+ }, layout.dagPaths)
+
+ var paths []string
+ offset := 0
+ for _, f := range layout.files {
+ paths = append(paths, f.path)
+ data, err := os.ReadFile(f.disk)
+ require.NoError(t, err)
+ assert.Equal(t, offset, f.offset, f.path)
+ assert.Equal(t, len(data), f.length, f.path)
+ assert.Equal(t, sha256Hex(data), f.sha256, f.path)
+ assert.Equal(t, data, region[f.offset:f.offset+f.length],
f.path)
+ offset += f.length
+ }
+ assert.Equal(t, []string{"cmd/bundle/main.go", "dags/a.go",
"dags/b.go"}, paths,
+ "entrypoint first, then the rest sorted, each file once")
+ assert.Len(t, region, offset)
+
+ assert.Contains(
+ t,
+ stderr.String(),
+ `warning: source file "example.com/[email protected]/x.go" of dag
"foreign"`,
+ )
+ assert.Equal(t, 1, bytes.Count(stderr.Bytes(), []byte("warning:")))
+}
+
+func TestLayoutSources_EntrypointOutsideModule(t *testing.T) {
+ mod := goModule{dir: t.TempDir()}
+ elsewhere := filepath.Join(t.TempDir(), "main.go")
+ require.NoError(t, os.WriteFile(elsewhere, []byte("package main\n"),
0o644))
+
+ layout, _, err := layoutSources(&bytes.Buffer{},
airflowmetadata.Manifest{}, elsewhere, mod)
+ require.NoError(t, err)
+ assert.Equal(t, "main.go", layout.entrypoint)
+ assert.Empty(t, layout.dagPaths)
+}
+
+func TestRunPack_EmbedsEntrypointWhenBundleHasNoDagSources(t *testing.T) {
+ dir := writeModule(
+ t,
+ "example.com/app",
+ map[string]string{"main.go": "package main\nfunc main() {}\n"},
+ )
+ exe := filepath.Join(dir, "prebuilt")
+ require.NoError(t, os.WriteFile(exe, []byte("prebuilt-binary-bytes"),
0o755))
+ meta := filepath.Join(dir, "airflow-metadata.json")
+ require.NoError(t, os.WriteFile(meta, []byte(
+ `{"airflow_bundle_metadata_version":"1.0",`+
+
`"sdk":{"language":"go","version":"0.1.0","supervisor_schema_version":"2026-06-16"},`+
+ `"dags":{"my_dag":{"tasks":["t1"]}}}`,
+ ), 0o644))
+ out := filepath.Join(t.TempDir(), "bundle")
+
+ require.NoError(t, runPack(&bytes.Buffer{}, &bytes.Buffer{},
&packOptions{
+ executable: exe,
+ source: filepath.Join(dir, "main.go"),
+ airflowMetadata: meta,
+ output: out,
+ }))
+
+ m, region, _ := readPacked(t, out)
+ assert.Equal(t, "main.go", m.EntrypointPath)
+ assert.Empty(t, m.DagSourcePaths)
+ require.Len(t, m.Sources, 1)
+ assert.Equal(t, "package main\nfunc main() {}\n", string(region))
+}
+
+func TestRunPack_DagFileOutsideModuleFallsBackToEntrypoint(t *testing.T) {
+ dir := writeModule(
+ t,
+ "example.com/app",
+ map[string]string{"main.go": "package main\nfunc main() {}\n"},
+ )
+ exe := filepath.Join(dir, "prebuilt")
+ require.NoError(t, os.WriteFile(exe, []byte("prebuilt-binary-bytes"),
0o755))
+ meta := filepath.Join(dir, "airflow-metadata.json")
+ require.NoError(t, os.WriteFile(meta, []byte(
+ `{"airflow_bundle_metadata_version":"1.0",`+
+
`"sdk":{"language":"go","version":"0.1.0","supervisor_schema_version":"2026-06-16"},`+
+ `"dags":{"py_dag":{"tasks":["t1"]}},`+
+
`"dag_source_files":{"vendored":"/root/go/pkg/mod/example.com/[email protected]/dag.go"}}`,
+ ), 0o644))
+ out := filepath.Join(t.TempDir(), "bundle")
+
+ var stderr bytes.Buffer
+ require.NoError(t, runPack(&bytes.Buffer{}, &stderr, &packOptions{
+ executable: exe,
+ source: filepath.Join(dir, "main.go"),
+ airflowMetadata: meta,
+ output: out,
+ }))
+
+ m, _, _ := readPacked(t, out)
+ assert.Equal(t, map[string]string{"vendored": "main.go"},
m.DagSourcePaths)
+ require.Len(t, m.Sources, 1)
+ assert.Contains(
+ t,
+ stderr.String(),
+ `warning: source file
"/root/go/pkg/mod/example.com/[email protected]/dag.go" of dag "vendored"`,
+ )
+}
+
+func TestRunPack_RejectsOutputAliasingDagSourceFile(t *testing.T) {
+ dir := writeModule(t, "example.com/app", map[string]string{
+ "main.go": "package main\nfunc main() {}\n",
+ "dags/reports.go": "package dags\n",
+ })
+ exe := filepath.Join(dir, "prebuilt")
+ require.NoError(t, os.WriteFile(exe, []byte("prebuilt-binary-bytes"),
0o755))
+ reports := filepath.Join(dir, "dags", "reports.go")
+ meta := filepath.Join(dir, "airflow-metadata.json")
+ require.NoError(t, os.WriteFile(meta, []byte(
+ `{"airflow_bundle_metadata_version":"1.0",`+
+
`"sdk":{"language":"go","version":"0.1.0","supervisor_schema_version":"2026-06-16"},`+
+
`"dags":{},"dag_source_files":{"reports":`+strconv.Quote(reports)+`}}`,
+ ), 0o644))
+
+ err := runPack(io.Discard, io.Discard, &packOptions{
+ executable: exe,
+ source: filepath.Join(dir, "main.go"),
+ airflowMetadata: meta,
+ output: reports,
+ })
+ require.ErrorContains(t, err, "same file as the source")
+
+ got, readErr := os.ReadFile(reports)
+ require.NoError(t, readErr)
+ assert.Equal(t, "package dags\n", string(got))
+}
+
+// Packs the multidag fixture through the real build path: three Dag files
across two imported
+// packages plus the entrypoint, with and without -trimpath.
+func TestPack_MultiDagBundle(t *testing.T) {
+ if testing.Short() {
+ t.Skip("shells out to `go build`")
+ }
+ if _, err := exec.LookPath("go"); err != nil {
+ t.Skip("go toolchain not on PATH")
+ }
+
+ for _, tc := range []struct {
+ name string
+ build []string
+ }{
+ {name: "build paths"},
+ {name: "trimpath", build: []string{"--", "-trimpath"}},
+ } {
+ t.Run(tc.name, func(t *testing.T) {
+ const fixture = "cmd/airflow-go-pack/testdata/multidag"
+ out := filepath.Join(t.TempDir(), "bundle")
+ var stderr bytes.Buffer
+ cmd := newRootCmd()
+ cmd.SetArgs(append([]string{"./testdata/multidag",
"--output", out}, tc.build...))
+ cmd.SetOut(&bytes.Buffer{})
+ cmd.SetErr(&stderr)
+ require.NoError(t, cmd.Execute())
+ assert.NotContains(t, stderr.String(), "warning")
+
+ m, region, _ := readPacked(t, out)
+ assert.Equal(t, fixture+"/main.go", m.EntrypointPath)
+ assert.Equal(t, map[string]string{
+ "orders": fixture + "/main.go",
+ "reports": fixture + "/reports/reports.go",
+ "billing": fixture + "/factory/factory.go",
+ "shipping": fixture + "/factory/factory.go",
+ }, m.DagSourcePaths)
+
+ var paths []string
+ for _, s := range m.Sources {
+ paths = append(paths, s.Path)
+ data, err := os.ReadFile(filepath.Join("..",
"..", filepath.FromSlash(s.Path)))
+ require.NoError(t, err)
+ assert.Equal(t, data,
region[s.Offset:s.Offset+s.Length], s.Path)
+ assert.Equal(t, sha256Hex(data), s.SHA256,
s.Path)
+ }
+ assert.Equal(t, []string{
+ fixture + "/main.go",
+ fixture + "/factory/factory.go",
+ fixture + "/reports/reports.go",
+ }, paths)
+ })
+ }
+}
+
+// A bundle with task handlers only has no Dag source of its own: just the
entrypoint.
+func TestPack_TaskHandlersOnlyBundle(t *testing.T) {
+ if testing.Short() {
+ t.Skip("shells out to `go build`")
+ }
+ if _, err := exec.LookPath("go"); err != nil {
+ t.Skip("go toolchain not on PATH")
+ }
+
+ out := filepath.Join(t.TempDir(), "bundle")
+ cmd := newRootCmd()
+ cmd.SetArgs([]string{"./testdata/handlersonly", "--output", out})
+ cmd.SetOut(&bytes.Buffer{})
+ cmd.SetErr(&bytes.Buffer{})
+ require.NoError(t, cmd.Execute())
+
+ m, _, metadata := readPacked(t, out)
+ assert.Equal(t, "cmd/airflow-go-pack/testdata/handlersonly/main.go",
m.EntrypointPath)
+ assert.Empty(t, m.DagSourcePaths)
+ assert.Len(t, m.Sources, 1)
+ assert.Contains(t, string(metadata), "dag_source_paths: {}")
+}
+
+// Identical inputs give a byte-identical bundle.
+func TestPack_IsDeterministic(t *testing.T) {
+ if testing.Short() {
+ t.Skip("shells out to `go build`")
+ }
+ if _, err := exec.LookPath("go"); err != nil {
+ t.Skip("go toolchain not on PATH")
+ }
+
+ exe := filepath.Join(t.TempDir(), "multidag")
+ goBuild(t, "./testdata/multidag", exe, runtime.GOOS, runtime.GOARCH)
+
+ pack := func() string {
+ out := filepath.Join(t.TempDir(), "bundle")
+ cmd := newRootCmd()
+ cmd.SetArgs([]string{
+ "--executable", exe,
+ "--source", "testdata/multidag/main.go",
+ "--output", out,
+ })
+ cmd.SetOut(&bytes.Buffer{})
+ cmd.SetErr(&bytes.Buffer{})
+ require.NoError(t, cmd.Execute())
+ return out
+ }
+ first, second := pack(), pack()
+ a, err := os.ReadFile(first)
+ require.NoError(t, err)
+ b, err := os.ReadFile(second)
+ require.NoError(t, err)
+ assert.Equal(t, a, b)
+}
diff --git a/go-sdk/internal/bundle/doc.go
b/go-sdk/cmd/airflow-go-pack/testdata/handlersonly/main.go
similarity index 65%
copy from go-sdk/internal/bundle/doc.go
copy to go-sdk/cmd/airflow-go-pack/testdata/handlersonly/main.go
index 8a77698fe70..e16dab50f88 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/cmd/airflow-go-pack/testdata/handlersonly/main.go
@@ -15,10 +15,22 @@
// 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, the Dag and task ids it lists in the manifest,
and
-// the serialized Dags it sends to the Dag processor.
-//
-// Package airflow builds the tasks and ids from the task handlers a bundle
-// registers, and the serialized Dags from its Dags.
-package bundle
+// Command handlersonly is a bundle fixture for the packer tests. It registers
task handlers and
+// no Dag from airflow.Dag.
+package main
+
+import (
+ "log"
+
+ "github.com/apache/airflow/go-sdk/airflow"
+)
+
+func extract(airflow.Context) error { return nil }
+
+func main() {
+ bundle := airflow.Bundle()
+ bundle.Register(airflow.TaskHandler("py_etl", "extract", extract))
+ if err := bundle.Serve(); err != nil {
+ log.Fatal(err)
+ }
+}
diff --git a/go-sdk/internal/bundle/doc.go
b/go-sdk/cmd/airflow-go-pack/testdata/multidag/factory/factory.go
similarity index 68%
copy from go-sdk/internal/bundle/doc.go
copy to go-sdk/cmd/airflow-go-pack/testdata/multidag/factory/factory.go
index 8a77698fe70..833e8e510aa 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/cmd/airflow-go-pack/testdata/multidag/factory/factory.go
@@ -15,10 +15,15 @@
// 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, the Dag and task ids it lists in the manifest,
and
-// the serialized Dags it sends to the Dag processor.
-//
-// Package airflow builds the tasks and ids from the task handlers a bundle
-// registers, and the serialized Dags from its Dags.
-package bundle
+// Package factory builds Dags for several dag_ids from one file.
+package factory
+
+import "github.com/apache/airflow/go-sdk/airflow"
+
+func run(airflow.Context) error { return nil }
+
+func New(dagID string) *airflow.DagRef {
+ dag := airflow.Dag(dagID)
+ dag.Task(run)
+ return dag
+}
diff --git a/go-sdk/internal/bundle/doc.go
b/go-sdk/cmd/airflow-go-pack/testdata/multidag/main.go
similarity index 50%
copy from go-sdk/internal/bundle/doc.go
copy to go-sdk/cmd/airflow-go-pack/testdata/multidag/main.go
index 8a77698fe70..64db46f2626 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/cmd/airflow-go-pack/testdata/multidag/main.go
@@ -15,10 +15,35 @@
// 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, the Dag and task ids it lists in the manifest,
and
-// the serialized Dags it sends to the Dag processor.
-//
-// Package airflow builds the tasks and ids from the task handlers a bundle
-// registers, and the serialized Dags from its Dags.
-package bundle
+// Command multidag is a bundle fixture for the packer tests. It declares Dags
in its own file,
+// in an imported package, and through a factory in a third package.
+package main
+
+import (
+ "log"
+
+ "github.com/apache/airflow/go-sdk/airflow"
+
"github.com/apache/airflow/go-sdk/cmd/airflow-go-pack/testdata/multidag/factory"
+
"github.com/apache/airflow/go-sdk/cmd/airflow-go-pack/testdata/multidag/reports"
+)
+
+func extract(airflow.Context) error { return nil }
+
+func main() {
+ bundle := airflow.Bundle()
+
+ orders := airflow.Dag("orders")
+ orders.Task(extract)
+
+ bundle.Register(
+ airflow.TaskHandler("py_etl", "extract", extract),
+ orders,
+ reports.Dag(),
+ factory.New("billing"),
+ factory.New("shipping"),
+ )
+
+ if err := bundle.Serve(); err != nil {
+ log.Fatal(err)
+ }
+}
diff --git a/go-sdk/internal/bundle/doc.go
b/go-sdk/cmd/airflow-go-pack/testdata/multidag/reports/reports.go
similarity index 68%
copy from go-sdk/internal/bundle/doc.go
copy to go-sdk/cmd/airflow-go-pack/testdata/multidag/reports/reports.go
index 8a77698fe70..7b432b41c73 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/cmd/airflow-go-pack/testdata/multidag/reports/reports.go
@@ -15,10 +15,15 @@
// 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, the Dag and task ids it lists in the manifest,
and
-// the serialized Dags it sends to the Dag processor.
-//
-// Package airflow builds the tasks and ids from the task handlers a bundle
-// registers, and the serialized Dags from its Dags.
-package bundle
+// Package reports declares a Dag in a file other than the bundle's entrypoint.
+package reports
+
+import "github.com/apache/airflow/go-sdk/airflow"
+
+func build(airflow.Context) error { return nil }
+
+func Dag() *airflow.DagRef {
+ dag := airflow.Dag("reports")
+ dag.Task(build)
+ return dag
+}
diff --git a/go-sdk/example/bundle/Justfile b/go-sdk/example/bundle/Justfile
index 07df82a2280..97ae2025c14 100644
--- a/go-sdk/example/bundle/Justfile
+++ b/go-sdk/example/bundle/Justfile
@@ -26,8 +26,8 @@ build: pack
# One-step build + pack. The single `go tool airflow-go-pack`
# invocation runs `go build` internally, queries the binary for its
-# DAG/task identity via --airflow-metadata, and appends the source plus
-# airflow-metadata.yaml plus AFBNDL01 trailer. The output is a single
+# DAG/task identity via --airflow-metadata, and appends the source files
+# plus airflow-metadata.yaml plus AFBNDL01 trailer. The output is a single
# self-contained executable bundle, named after the bundle's package
# directory and written to the current directory. Drop it into the
# Dag bundle named by ExecutableCoordinator's task_handler_bundle_name
diff --git a/go-sdk/internal/airflowmetadata/airflowmetadata.go
b/go-sdk/internal/airflowmetadata/airflowmetadata.go
index d11412d722c..cc0a8ac91c4 100644
--- a/go-sdk/internal/airflowmetadata/airflowmetadata.go
+++ b/go-sdk/internal/airflowmetadata/airflowmetadata.go
@@ -36,6 +36,9 @@ type Manifest struct {
AirflowBundleMetadataVersion string
`json:"airflow_bundle_metadata_version" yaml:"airflow_bundle_metadata_version"`
SDK SDK `json:"sdk"
yaml:"sdk"`
Dags map[string]Dag `json:"dags"
yaml:"dags"`
+ // DagSourceFiles maps the dag_id of each Dag authored with airflow.Dag
to the source file
+ // that declared it, as the binary was built. Only the packer reads it,
to embed the files.
+ DagSourceFiles map[string]string `json:"dag_source_files,omitempty"
yaml:"dag_source_files,omitempty"`
}
// SDK identifies the SDK that produced the bundle.
diff --git a/go-sdk/internal/bundle/doc.go b/go-sdk/internal/bundle/doc.go
index 8a77698fe70..0222ffffe0c 100644
--- a/go-sdk/internal/bundle/doc.go
+++ b/go-sdk/internal/bundle/doc.go
@@ -16,9 +16,10 @@
// under the License.
// 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.
+// tasks it looks up and runs, the Dag and task ids and the Dag source files it
+// lists in the manifest, and the serialized Dags it sends to the Dag
processor.
//
-// Package airflow builds the tasks and ids from the task handlers a bundle
-// registers, and the serialized Dags from its Dags.
+// Package airflow builds the tasks from both the task handlers and the Dags a
+// bundle registers, the ids from its task handlers, and the source files and
+// serialized Dags from its Dags.
package bundle
diff --git a/go-sdk/internal/bundle/task.go b/go-sdk/internal/bundle/task.go
index c3cb8cc6557..a8aff19e9af 100644
--- a/go-sdk/internal/bundle/task.go
+++ b/go-sdk/internal/bundle/task.go
@@ -71,6 +71,11 @@ type DagSerializer interface {
SerializeDags(fileloc, relativeFileloc string) []SerializedDag
}
+// DagSourceLister reports the source file that declared each Dag of a bundle,
by dag_id.
+type DagSourceLister interface {
+ ListDagSourceFiles() map[string]string
+}
+
type taskFunction struct {
fn reflect.Value
fullName string
diff --git a/go-sdk/pkg/execution/metadata.go b/go-sdk/pkg/execution/metadata.go
index 05dbebd70a9..807db34fe55 100644
--- a/go-sdk/pkg/execution/metadata.go
+++ b/go-sdk/pkg/execution/metadata.go
@@ -95,6 +95,11 @@ func collectManifest(b bundle.EnumerableBundle)
airflowmetadata.Manifest {
dag.Tasks = append(dag.Tasks, handler.TaskID)
meta.Dags[handler.DagID] = dag
}
+ if lister, ok := b.(bundle.DagSourceLister); ok {
+ if files := lister.ListDagSourceFiles(); len(files) > 0 {
+ meta.DagSourceFiles = files
+ }
+ }
return meta
}
diff --git a/go-sdk/pkg/execution/metadata_test.go
b/go-sdk/pkg/execution/metadata_test.go
index d105b24d898..9c0928f49b9 100644
--- a/go-sdk/pkg/execution/metadata_test.go
+++ b/go-sdk/pkg/execution/metadata_test.go
@@ -26,8 +26,20 @@ import (
"gopkg.in/yaml.v3"
"github.com/apache/airflow/go-sdk/internal/airflowmetadata"
+ "github.com/apache/airflow/go-sdk/internal/bundle"
)
+type handlersOnly []bundle.TaskHandlerInfo
+
+func (h handlersOnly) ListTaskHandlers() []bundle.TaskHandlerInfo { return h }
+
+type handlersAndSources struct {
+ handlersOnly
+ files map[string]string
+}
+
+func (h handlersAndSources) ListDagSourceFiles() map[string]string { return
h.files }
+
func sampleManifest() airflowmetadata.Manifest {
return airflowmetadata.Manifest{
// "1.0" and a task named "123" must survive a YAML round-trip
as strings.
@@ -117,3 +129,27 @@ func TestEncodeManifest_UnsupportedFormat(t *testing.T) {
_, err := encodeManifest(sampleManifest(), MetadataFormat("xml"))
require.Error(t, err)
}
+
+func TestCollectManifest_DagSourceFiles(t *testing.T) {
+ handlers := handlersOnly{{DagID: "d", TaskID: "t"}}
+
+ t.Run("bundle without a lister", func(t *testing.T) {
+ meta := collectManifest(handlers)
+ assert.Nil(t, meta.DagSourceFiles)
+ data, err := encodeManifest(meta, MetadataFormatYAML)
+ require.NoError(t, err)
+ assert.NotContains(t, string(data), "dag_source_files")
+ })
+
+ t.Run("lister with files", func(t *testing.T) {
+ files := map[string]string{"native": "/src/native.go"}
+ meta := collectManifest(handlersAndSources{handlers, files})
+ assert.Equal(t, files, meta.DagSourceFiles)
+ assert.Equal(t, []string{"t"}, meta.Dags["d"].Tasks)
+ })
+
+ t.Run("lister without files", func(t *testing.T) {
+ meta := collectManifest(handlersAndSources{handlers, nil})
+ assert.Nil(t, meta.DagSourceFiles)
+ })
+}