kaxil commented on code in PR #74342:
URL: https://github.com/apache/airflow/pull/74342#discussion_r4206219578


##########
go-sdk/cmd/airflow-go-pack/sources.go:
##########
@@ -0,0 +1,255 @@
+// 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 {
+               if _, ok := relativeTo(lines[i], absDir); ok && len(lines[i]) > 
len(best.dir) {

Review Comment:
   When the entrypoint's directory is the module root itself (a plain `go mod 
init` repo with `main.go` at the top), `relativeTo` returns false for `rel == 
"."`, so the only module `go list -m` reports gets dropped and the pack fails 
with "no module contains the directory". The same module packs fine on main. 
Could this accept the equal case, e.g. `(ok || lines[i] == absDir)`? The other 
`relativeTo` callers compare a file against a directory, so the helper itself 
can stay as is. Every current fixture sits below the go-sdk module root, so a 
pack test with `main.go` at the root of a temp module would pin it.



##########
go-sdk/cmd/airflow-go-pack/pack.go:
##########
@@ -155,15 +156,15 @@ 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
        }
 
        meta, err := obtainMetadata(opts, introspectPath)
        if err != nil {
                return err
        }
-       if len(meta.Dags) == 0 {
+       if len(meta.Dags) == 0 && len(meta.DagSourceFiles) == 0 {

Review Comment:
   With only native Dags this now writes `dags: {}`, but the schema still has 
`minProperties: 1` on `dags` and the spec says every dag_id the bundle exposes 
MUST appear there, and #74341 leaves both unchanged. The reader only `.get`s 
`dags`, so nothing breaks at runtime, but the two PRs publish contradicting 
contracts. Should #74341 relax the schema and spec so `dags` lists only 
task-handler Dags? The `Wrote bundle ... dags=%d` summary at the end of 
`runPack` has the same blind spot and prints `dags=0` for a native-only bundle.



##########
go-sdk/cmd/airflow-go-pack/inspect.go:
##########
@@ -52,6 +59,45 @@ 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 {

Review Comment:
   Bundles packed by the released `go-sdk/v1.0.0-beta2` and `beta3` packers 
have a source region and a `source:` name but no `sources` index, so `inspect 
--source` on them now errors here instead of printing the file. Could this fall 
back to one file covering the whole region, named by the legacy `source` field?



##########
go-sdk/cmd/airflow-go-pack/sources_test.go:
##########
@@ -0,0 +1,402 @@
+// 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"
+       "os"
+       "os/exec"
+       "path/filepath"
+       "runtime"
+       "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 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 TestRejectOutputAlias_EmbeddedSourceFile(t *testing.T) {
+       dir := t.TempDir()
+       exe := filepath.Join(dir, "exe")
+       entry := filepath.Join(dir, "main.go")
+       other := filepath.Join(dir, "dag.go")
+       for _, p := range []string{exe, entry, other} {
+               require.NoError(t, os.WriteFile(p, []byte("x"), 0o644))
+       }
+
+       err := rejectOutputAlias(other, exe, []string{entry, other}, "")

Review Comment:
   This calls `rejectOutputAlias` directly, so removing the second call in 
`runPack` (the one that covers the per-Dag files) leaves the suite green. That 
call is what stops `--output dags/reports.go` from replacing a Dag file with 
the binary, so could this go through `runPack` the way 
`TestRunPack_RejectsOutputAliasingMetadataFile` does?



##########
go-sdk/cmd/airflow-go-pack/pack.go:
##########
@@ -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,
+                       scalar(id),

Review Comment:
   The keys here go through `scalar(id)` while the values are quoted, so a 
dag_id like `2024`, `on` or `null` is written unquoted. yaml.v3 emits `2024: 
"main.go"`, and PyYAML `safe_load` (what `_bundle_metadata.py` uses) reads it 
back as the int `2024` (`on` becomes `True`), so `resolve_source_path` in 
#74341 misses and the Code view quietly shows the entrypoint. 
`quotedScalar(id)` would fix it. The `dags` keys above have the same problem, 
but that predates this PR.



##########
go-sdk/airflow/dag.go:
##########
@@ -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.

Review Comment:
   ADR-0015 says Go adopts the TypeScript shape so the Code tab behaves the 
same across SDKs, and TypeScript records a factory-built Dag against the module 
whose body was running (the caller), while this records the factory's file. So 
`factory.New("billing")` called from `main.go` shows `factory/factory.go`, 
which doesn't contain the id. The call site of `Dag` seems like the right rule 
for Go, but could ADR-0015 say so, rather than documenting the two SDKs as 
identical?



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

To unsubscribe, e-mail: [email protected]

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

Reply via email to