jason810496 commented on code in PR #74342: URL: https://github.com/apache/airflow/pull/74342#discussion_r4209195351
########## 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: Fixed in 3750a31e3dd: `listModule` now accepts the module root itself. `TestListModule_EntrypointAtModuleRoot` covers `main.go` at the module root. ########## 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: Fixed in 4bc863d0f2c: Dag ids are quoted in both `dags` and `dag_source_paths`. A `!!str` tag alone was not enough, since yaml.v3 still writes `on` plain. `TestRenderManifest_QuotesDagIDs` round-trips `2024` and `on`. ########## 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: The summary now prints `task_handler_dags` and `native_dags` (f0004dc79c5). For the schema and spec, a follow-up removes the dag_id metadata from the artifact, which resolves the contract, so both stay unchanged here. ########## 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: The Go SDK releases so far are betas, and we don't keep compatibility code for them. Instead, 0d5f0e31813 makes the error tell the user to repack the bundle with the current airflow-go-pack. ########## 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: Agreed, the call site stays the rule for Go, since Go has no module body to attribute the Dag to. 0e71e05dcf8 updates ADR-0015 to say so. ########## 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: Replaced with `TestRunPack_RejectsOutputAliasingDagSourceFile` in 43d3315a03e. It goes through `runPack` and fails if the second `rejectOutputAlias` call is removed. -- 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]
