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]
