diff --git a/pkg/airflowrt/include/gitignore b/pkg/airflowrt/include/gitignore index a549900b8..64c686f31 100644 --- a/pkg/airflowrt/include/gitignore +++ b/pkg/airflowrt/include/gitignore @@ -15,4 +15,4 @@ airflow.db .astro/*.local.yaml **/.astro/dbt_metadata.json -**/.astro/manifest.slim.json +**/.astro/*manifest*.slim.json diff --git a/pkg/cosmosboost/precompute/cleanup.go b/pkg/cosmosboost/precompute/cleanup.go index c50c64871..586f23563 100644 --- a/pkg/cosmosboost/precompute/cleanup.go +++ b/pkg/cosmosboost/precompute/cleanup.go @@ -7,6 +7,7 @@ import ( "io/fs" "os" "path/filepath" + "strings" "time" ) @@ -39,8 +40,9 @@ type CleanupSummary struct { } // Cleanup removes every .astro/dbt_metadata.json sidecar and every -// .astro/manifest.slim.json under the given roots that this tool wrote, -// pruning each containing .astro directory when removal leaves it empty. +// .astro/*.slim.json slim manifest under the given roots that this tool +// wrote, pruning each containing .astro directory when removal leaves it +// empty. // // Each file is judged by its own producer marker, never by its neighbor's, // because the mixed states are real: a slim manifest outlives its sidecar when @@ -70,10 +72,10 @@ func Cleanup(roots []string) (CleanupSummary, error) { if filepath.Base(filepath.Dir(path)) != sidecarDir { return nil } - switch d.Name() { - case sidecarName: + switch { + case d.Name() == sidecarName: recordOnce(seen, &results, path, removeSidecar) - case slimManifestName: + case strings.HasSuffix(d.Name(), slimManifestSuffix): recordOnce(seen, &results, path, removeSlimManifest) } return nil @@ -143,8 +145,8 @@ func removeSidecar(path string) CleanupResult { }) } -// removeSlimManifest removes one .astro/manifest.slim.json, keeping it unless -// its own _generated_by marker names this tool. +// removeSlimManifest removes one .astro/*.slim.json, keeping it unless its +// own _generated_by marker names this tool. func removeSlimManifest(path string) CleanupResult { return removeArtifact(path, func(data []byte) bool { var marker struct { diff --git a/pkg/cosmosboost/precompute/cleanup_test.go b/pkg/cosmosboost/precompute/cleanup_test.go index ab25e75aa..181ecaf73 100644 --- a/pkg/cosmosboost/precompute/cleanup_test.go +++ b/pkg/cosmosboost/precompute/cleanup_test.go @@ -96,6 +96,57 @@ func TestCleanupRemovesSlimManifestAlongsideSidecar(t *testing.T) { } } +// TestCleanupRemovesCustomNamedSlimManifest: Cleanup matches by the +// .slim.json suffix, not a fixed literal name. +func TestCleanupRemovesCustomNamedSlimManifest(t *testing.T) { + dir := t.TempDir() + writeFiles(t, dir, map[string]string{ + "manifest_full.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + }) + summary, err := Run([]string{dir}, "test", Options{SlimManifest: true}) + if err != nil || summary.CountFailed() > 0 { + t.Fatalf("stamping fixture manifest failed: err=%v failed=%d", err, summary.CountFailed()) + } + slimPath := filepath.Join(dir, sidecarDir, "manifest_full.slim.json") + if _, err := os.Stat(slimPath); err != nil { + t.Fatalf("fixture setup: slim manifest not written: %v", err) + } + + cleanupSummary, err := Cleanup([]string{dir}) + if err != nil { + t.Fatalf("Cleanup: %v", err) + } + if got := len(cleanupSummary.Results); got != 2 { + t.Fatalf("results = %d, want 2 (sidecar + slim manifest)", got) + } + if cleanupSummary.CountFailed() != 0 || cleanupSummary.CountKept() != 0 { + t.Fatalf("failed=%d kept=%d, want 0/0", cleanupSummary.CountFailed(), cleanupSummary.CountKept()) + } + if _, err := os.Stat(slimPath); !os.IsNotExist(err) { + t.Fatalf("custom-named slim manifest still present after cleanup: %v", err) + } +} + +// TestCleanupKeepsForeignSlimJSONSuffixedFile: the broader suffix match must +// not weaken the ownership check for a *.slim.json file we didn't write. +func TestCleanupKeepsForeignSlimJSONSuffixedFile(t *testing.T) { + dir := t.TempDir() + writeFiles(t, dir, map[string]string{ + ".astro/some_other_tool.slim.json": `{"_generated_by": {"application": "someone-else"}}`, + }) + + summary, err := Cleanup([]string{dir}) + if err != nil { + t.Fatalf("Cleanup: %v", err) + } + if got := summary.CountKept(); got != 1 { + t.Fatalf("kept = %d, want 1", got) + } + if _, err := os.Stat(filepath.Join(dir, sidecarDir, "some_other_tool.slim.json")); err != nil { + t.Fatalf("foreign *.slim.json file was removed: %v", err) + } +} + // TestCleanupJudgesEachArtifactSeparately: provenance is read per file, never // inferred from the neighbor. Deleting on the neighbor's marker would destroy // a file we do not own in one direction, and strand a stale artifact of ours - diff --git a/pkg/cosmosboost/precompute/discover.go b/pkg/cosmosboost/precompute/discover.go index d5190817b..08abfd5ea 100644 --- a/pkg/cosmosboost/precompute/discover.go +++ b/pkg/cosmosboost/precompute/discover.go @@ -4,14 +4,13 @@ import ( "io/fs" "os" "path/filepath" + "strings" ) // dbt requires the project file to be named exactly dbt_project.yml (the .yaml // extension is not accepted), so we match only this name. const dbtProjectFile = "dbt_project.yml" -const manifestFile = "manifest.json" - // findProjects walks root and returns every directory that contains a // dbt_project.yml. // @@ -52,13 +51,23 @@ var manifestSkipDirs = map[string]bool{ gitDir: true, // VCS internals can't hold a project's manifest } -// findManifests walks root and returns manifest.json file paths. +// isManifestCandidateName reports whether a *.json file could be a dbt +// manifest by name alone (case-insensitive "manifest" anywhere); content is +// validated separately (isDbtManifest). +func isManifestCandidateName(name string) bool { + lower := strings.ToLower(name) + return strings.HasSuffix(lower, ".json") && strings.Contains(lower, "manifest") +} + +// findManifests walks root and returns every file matching +// isManifestCandidateName. // // A manifest whose parent directory is itself a discovered project root is // omitted: that project's folder hash already covers a manifest sitting in its -// root. Manifests elsewhere — most importantly a standalone one shipped for a -// manifest-only (DBT_MANIFEST) deployment, or a project's target/manifest.json — -// each get their own sidecar. +// root, and processProject slims it there instead. Manifests elsewhere — +// most importantly a standalone one shipped for a manifest-only +// (DBT_MANIFEST) deployment, or a project's target/manifest.json — each get +// their own sidecar. func findManifests(root string, projectDirs map[string]bool) ([]string, error) { var manifests []string @@ -72,7 +81,7 @@ func findManifests(root string, projectDirs map[string]bool) ([]string, error) { } return nil } - if d.Name() != manifestFile { + if !isManifestCandidateName(d.Name()) { return nil } if projectDirs[filepath.Dir(path)] { diff --git a/pkg/cosmosboost/precompute/hash.go b/pkg/cosmosboost/precompute/hash.go index b99d1a5e8..7331f0cf3 100644 --- a/pkg/cosmosboost/precompute/hash.go +++ b/pkg/cosmosboost/precompute/hash.go @@ -10,6 +10,7 @@ import ( "os" "path/filepath" "sort" + "strings" ) // excludedDirs are directory names skipped during *project* discovery @@ -122,6 +123,9 @@ func hashProject(dir string, cfg dbtConfig) (hash string, files int, totalBytes digest, n, err := hashFile(path) if err != nil { + if !strings.Contains(rel, "/") && isManifestCandidateName(d.Name()) { + return nil // processProject re-reads and warns on this one itself + } return err } entries = append(entries, entry{rel, digest}) @@ -177,16 +181,19 @@ func stripEntityCreatedAt(doc map[string]any) { } } -// isDbtManifest reports whether doc looks like a dbt manifest. dbt always writes -// metadata.dbt_schema_version, which unrelated manifest.json files (web app/PWA, -// tooling, etc.) do not have — so we only stamp files that carry it. +// isDbtManifest reports whether doc is a dbt manifest specifically - every +// dbt artifact (run_results.json, catalog.json, semantic_manifest.json, ...) +// carries the same metadata.dbt_schema_version, just under its own schema +// URL, so presence alone would also match those - most importantly +// semantic_manifest.json, whose name (like discovery's) also contains +// "manifest". func isDbtManifest(doc map[string]any) bool { meta, ok := doc["metadata"].(map[string]any) if !ok { return false } v, ok := meta["dbt_schema_version"].(string) - return ok && v != "" + return ok && strings.Contains(v, "/manifest/") } // readManifestDoc reads path and parses it as JSON. isDbt reports whether the diff --git a/pkg/cosmosboost/precompute/hash_test.go b/pkg/cosmosboost/precompute/hash_test.go index 8d3b7dd04..7947e9714 100644 --- a/pkg/cosmosboost/precompute/hash_test.go +++ b/pkg/cosmosboost/precompute/hash_test.go @@ -220,15 +220,17 @@ func TestHashManifestIgnoresVolatileMetadata(t *testing.T) { } } -// TestHashManifestSkipsNonDBT verifies that a manifest.json lacking the dbt shape -// (e.g. a web-app/PWA manifest) or invalid JSON is not treated as a dbt manifest, -// so it won't be stamped. +// TestHashManifestSkipsNonDBT: a non-dbt manifest.json, invalid JSON, or +// another dbt artifact (same metadata.dbt_schema_version field) isn't +// treated as a dbt manifest. func TestHashManifestSkipsNonDBT(t *testing.T) { dir := t.TempDir() cases := map[string]string{ - "webapp.json": `{"name":"My App","short_name":"App","icons":[]}`, // no metadata.dbt_schema_version - "nometa.json": `{"nodes":{"model.x":{}}}`, // nodes but no metadata - "invalid.json": `{not json`, // not JSON at all + "webapp.json": `{"name":"My App","short_name":"App","icons":[]}`, // no metadata.dbt_schema_version + "nometa.json": `{"nodes":{"model.x":{}}}`, // nodes but no metadata + "invalid.json": `{not json`, // not JSON at all + "run_results.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/run-results/v6.json"}}`, // a dbt artifact, but not the manifest + "catalog.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/catalog/v1.json"},"nodes":{}}`, // ditto, and it even has "nodes" } for name, content := range cases { p := filepath.Join(dir, name) @@ -484,7 +486,7 @@ func TestHashManifestStableAcrossFullParses(t *testing.T) { func TestHashManifestKeepsUserMetaCreatedAt(t *testing.T) { // The scalar top-level key pins that the created_at stripper tolerates // non-collection values in the document root. - base := `{"metadata": {"dbt_schema_version": "v12"}, "unrelated_scalar": 7, + base := `{"metadata": {"dbt_schema_version": "https://schemas.getdbt.com/dbt/manifest/v12.json"}, "unrelated_scalar": 7, "nodes": {"model.shop.a": {"name": "a", "created_at": 1.0, "meta": {"created_at": "%s"}}}}` dir := t.TempDir() writeFiles(t, dir, map[string]string{ diff --git a/pkg/cosmosboost/precompute/metadata.go b/pkg/cosmosboost/precompute/metadata.go index f10247451..0b74846bf 100644 --- a/pkg/cosmosboost/precompute/metadata.go +++ b/pkg/cosmosboost/precompute/metadata.go @@ -11,8 +11,9 @@ import ( const application = "astro" const ( - // schemaVersion lets the read-side cope with future format changes. - schemaVersion = 1 + // schemaVersion lets the read-side cope with future format changes. v2: + // filtered_manifest -> manifests, keyed by source filename. + schemaVersion = 2 // algoProjectTree hashes a whole dbt project directory (source files). // v2 excludes .git, which never ships in a deploy payload, so VCS activity @@ -37,25 +38,27 @@ const ( sidecarPerm = 0o644 ) -// Metadata is the content of the .astro/dbt_metadata.json sidecar. The Cosmos -// Boost plugin reads the version.hash field and uses it as the cache version key; -// it never recomputes the hash itself. +// Metadata is the content of the .astro/dbt_metadata.json sidecar. type Metadata struct { Schema int `json:"schema"` Version ProjectVersion `json:"version"` GeneratedBy GeneratedBy `json:"generated_by"` - // FilteredManifest points at the slim manifest beside this sidecar, absent - // when none was written. The sidecar is the consumer's entry point: it - // already reads and schema-gates this file, so it discovers the slim - // manifest here and falls back to the full one when the section is missing. - FilteredManifest *FilteredManifest `json:"filtered_manifest,omitempty"` + // Manifests is keyed by manifest filename, so a consumer looks up its own + // entry instead of trusting Version when a directory has several. + Manifests map[string]ManifestVersion `json:"manifests,omitempty"` } -// FilteredManifest describes a slim manifest sitting next to the sidecar. Path -// is relative to the sidecar's directory, and Version hashes the slim file's -// own bytes, so a consumer can confirm the two are a matched pair without -// re-hashing the full manifest. -type FilteredManifest struct { +// ManifestVersion is one manifest's own hash, plus its slim companion, if +// one was written. +type ManifestVersion struct { + Version ProjectVersion `json:"version"` + Slim *SlimManifest `json:"slim_manifest,omitempty"` +} + +// SlimManifest describes a slim manifest sitting next to the sidecar. Version +// hashes the slim file's own bytes, so a consumer can confirm the two are a +// matched pair without re-hashing the full manifest. +type SlimManifest struct { Schema int `json:"schema"` Path string `json:"path"` Version ProjectVersion `json:"version"` @@ -74,15 +77,13 @@ type GeneratedBy struct { Version string `json:"version"` } -// writeSidecar writes .astro/dbt_metadata.json inside dir. version is the -// producer's version, recorded in generated_by. filtered describes a slim -// manifest written beside it, or is nil when none was. -func writeSidecar(dir, algo, hash, version string, filtered *FilteredManifest) error { +// writeSidecar writes .astro/dbt_metadata.json inside dir. +func writeSidecar(dir, algo, hash, version string, manifests map[string]ManifestVersion) error { meta := Metadata{ - Schema: schemaVersion, - Version: ProjectVersion{Algo: algo, Hash: hash}, - GeneratedBy: GeneratedBy{Application: application, Version: version}, - FilteredManifest: filtered, + Schema: schemaVersion, + Version: ProjectVersion{Algo: algo, Hash: hash}, + GeneratedBy: GeneratedBy{Application: application, Version: version}, + Manifests: manifests, } // Metadata holds only strings and ints, so marshaling cannot fail. diff --git a/pkg/cosmosboost/precompute/precompute.go b/pkg/cosmosboost/precompute/precompute.go index 8b6866885..10b5796d8 100644 --- a/pkg/cosmosboost/precompute/precompute.go +++ b/pkg/cosmosboost/precompute/precompute.go @@ -4,6 +4,7 @@ import ( "encoding/json" "fmt" "io" + "os" "path/filepath" "runtime" "sort" @@ -18,18 +19,20 @@ const ( kindManifest = "manifest" ) +type unit struct{ kind, path string } + // Result records what happened for one unit of work: either a dbt project -// directory or a standalone manifest.json. +// directory or a standalone manifest. type Result struct { Kind string // "project" or "manifest" - Path string // project directory, or manifest.json path + Path string // project directory, or manifest path Hash string // version hash (empty if Err != nil or Skipped) Files int // files hashed (1 for a manifest) Bytes int64 // total bytes hashed Duration time.Duration // time spent on this unit - Skipped bool // a manifest.json that isn't a dbt manifest (no sidecar written) + Skipped bool // a manifest-like file that isn't a dbt manifest (no sidecar written) Warning string // non-fatal note (sidecar still written), e.g. an unresolved template - Err error // non-nil if hashing, writing the sidecar, or writing the slim manifest failed + Err error // non-nil if hashing, writing the sidecar, or writing a slim manifest failed } // Summary is the structured outcome of a precompute run. It backs both the @@ -43,26 +46,24 @@ type Summary struct { // hash sidecars. type Options struct { // SlimManifest also writes a slim, field-filtered copy of each discovered - // manifest.json (see buildSlimManifest) next to its sidecar. + // manifest (see buildSlimManifest) next to its sidecar. SlimManifest bool } -// Run finds every dbt project (a directory with dbt_project.yml) and standalone -// dbt manifest.json under the given roots, and writes a .astro/dbt_metadata.json -// hash sidecar next to each. Units are processed concurrently — one worker each, -// bounded by GOMAXPROCS — and each is hashed over sorted input, so results are -// deterministic with no cross-worker coordination. +// Run finds every dbt project (a directory with dbt_project.yml) and every +// standalone dbt manifest under the given roots, and writes a +// .astro/dbt_metadata.json hash sidecar next to each. version is recorded in +// each sidecar's generated_by. // -// Per-unit failures are best-effort: a unit that fails is recorded in its Result -// and does not stop the others. Run only returns a non-nil error for a top-level -// problem, such as a root that cannot be walked. version is recorded in each -// sidecar's generated_by. +// Per-unit failures are best-effort: a unit that fails is recorded in its +// Result and does not stop the others. Run only returns a non-nil error for a +// top-level problem, such as a root that cannot be walked. // -// With opts.SlimManifest set, every dbt manifest.json also gets a slim, -// field-filtered copy (see buildSlimManifest) written into the .astro/ beside -// it, for the Cosmos Boost plugin to load in place of the full manifest at -// DAG-parse time. That includes one in a project's own root, which is not a -// discovery unit of its own and is handled by processProject. +// Manifest units are computed concurrently, then written in a second pass +// grouped by directory: two manifests sharing a directory (see +// isManifestCandidateName) share one dbt_metadata.json, so writing it twice +// from two goroutines would race. Project units need no such grouping - each +// owns its directory exclusively. func Run(roots []string, version string, opts Options) (Summary, error) { start := time.Now() @@ -88,7 +89,6 @@ func Run(roots []string, version string, opts Options) (Summary, error) { } } - type unit struct{ kind, path string } var units []unit for d := range projectDirs { units = append(units, unit{kindProject, d}) @@ -103,6 +103,7 @@ func Run(roots []string, version string, opts Options) (Summary, error) { }) results := make([]Result, len(units)) + computations := make([]manifestComputation, len(units)) sem := make(chan struct{}, max(1, runtime.GOMAXPROCS(0))) var wg sync.WaitGroup for i, u := range units { @@ -114,12 +115,22 @@ func Run(roots []string, version string, opts Options) (Summary, error) { if u.kind == kindProject { results[i] = processProject(u.path, version, opts) } else { - results[i] = processManifest(u.path, version, opts) + computations[i] = computeManifest(u.path, version, opts) } }(i, u) } wg.Wait() + byDir := map[string][]int{} + for i, u := range units { + if u.kind == kindManifest { + byDir[filepath.Dir(u.path)] = append(byDir[filepath.Dir(u.path)], i) + } + } + for dir, idxs := range byDir { + writeManifestGroup(dir, idxs, units, computations, results, version) + } + return Summary{Duration: time.Since(start), Results: results}, nil } @@ -140,78 +151,165 @@ func processProject(dir, version string, opts Options) Result { " in dbt_project.yml hold unresolved Jinja templates; using the dbt default directories for exclusion (the real ones may add cache churn)" } - // A manifest.json in the project root is not a unit of its own - its .astro/ - // is this project's, so it would collide with the sidecar written below - and - // findManifests skips it for that reason. Slim it here instead, leaving the - // project's own hash as the anchor the pointer hangs off. - var filtered *FilteredManifest + // A manifest-like file in the project root is not a unit of its own - + // its .astro/ is this project's - so findManifests skips it. Slim every + // one found directly here instead, leaving the project's own hash as the + // anchor. + manifests := map[string]ManifestVersion{} if opts.SlimManifest { - if doc, _, isDbt, readErr := readManifestDoc(filepath.Join(dir, manifestFile)); readErr == nil && isDbt { - // Nothing mutates doc afterward here, unlike processManifest. - data, _ := json.Marshal(buildSlimManifest(doc, version)) - if filtered, r.Err = writeSlimManifest(dir, data); r.Err != nil { - r.Duration = time.Since(start) - return r + entries, readErr := os.ReadDir(dir) + if readErr != nil { + r.Warning = joinNotes(r.Warning, "could not scan for manifests to slim: "+readErr.Error()) + } else { + for _, e := range entries { + if e.IsDir() || !isManifestCandidateName(e.Name()) { + continue + } + doc, _, isDbt, readErr := readManifestDoc(filepath.Join(dir, e.Name())) + if readErr != nil { + r.Warning = joinNotes(r.Warning, e.Name()+": could not read as a manifest ("+readErr.Error()+")") + continue + } + if !isDbt { + continue + } + // Marshal before hashDocument mutates doc below (see computeManifest). + data, _ := json.Marshal(buildSlimManifest(doc, version)) + slim, writeErr := writeSlimManifest(dir, e.Name(), data) + if writeErr != nil { + // A failure slimming one candidate doesn't cost the + // project its tree-hash sidecar, or a sibling candidate + // its own already-written slim file. + r.Warning = joinNotes(r.Warning, e.Name()+": could not write its slim manifest ("+writeErr.Error()+")") + continue + } + manifests[e.Name()] = ManifestVersion{ + Version: ProjectVersion{Algo: algoManifestJSON, Hash: hashDocument(doc)}, + Slim: slim, + } } } } - r.Err = writeSidecar(dir, algoProjectTree, hash, version, filtered) + r.Err = writeSidecar(dir, algoProjectTree, hash, version, manifests) r.Duration = time.Since(start) return r } -// processManifest hashes one manifest.json and writes a sidecar next to it, -// plus a slim, field-filtered copy of the manifest when opts asks for one (see -// buildSlimManifest). A file that isn't a dbt manifest is skipped (nothing is -// written) so unrelated manifest.json files in the project aren't stamped. -func processManifest(path, version string, opts Options) Result { - start := time.Now() +// manifestComputation is what one manifest unit produces before any writes - +// Run groups these by directory so siblings share one sidecar write instead +// of racing each other for it (see writeManifestGroup). +type manifestComputation struct { + start time.Time + bytes int64 + isDbt bool + hash string + slimData []byte + err error +} + +// computeManifest reads and hashes one manifest-like file, and builds its +// slim copy when opts asks for one. A file that isn't actually a dbt manifest +// is left unhashed, so an unrelated *.json file whose name happens to contain +// "manifest" isn't stamped. +func computeManifest(path, version string, opts Options) manifestComputation { + c := manifestComputation{start: time.Now()} doc, bytes, isDbt, err := readManifestDoc(path) - var hash string - var slimData []byte - if err == nil && isDbt { - if opts.SlimManifest { - // Marshal before hashDocument mutates doc: the slim manifest shares - // doc's nested values, so only turning it into bytes here decouples - // the two. It holds JSON-native types only, so this cannot fail. - slimData, _ = json.Marshal(buildSlimManifest(doc, version)) - } - hash = hashDocument(doc) + c.bytes = bytes + c.isDbt = isDbt + if err != nil { + c.err = err + return c } - r := Result{Kind: kindManifest, Path: path, Hash: hash, Files: 1, Bytes: bytes, Duration: time.Since(start)} - switch { - case err != nil: - r.Err = err - case !isDbt: - r.Skipped = true - default: - dir := filepath.Dir(path) - // The sidecar goes last: it carries the filtered_manifest pointer, so it - // must never exist before the file it points at. Stopping short of it - // looks like "nothing was stamped", which BestEffortPreDeploy treats as - // safe. - var filtered *FilteredManifest - if slimData != nil { - filtered, r.Err = writeSlimManifest(dir, slimData) + if !isDbt { + return c + } + if opts.SlimManifest { + // Marshal before hashDocument mutates doc: the slim manifest shares + // doc's nested values, so only turning it into bytes here decouples + // the two. It holds JSON-native types only, so this cannot fail. + c.slimData, _ = json.Marshal(buildSlimManifest(doc, version)) + } + c.hash = hashDocument(doc) + return c +} + +// writeManifestGroup writes the slim file for each manifest in idxs that +// computed successfully, then one shared sidecar naming all of them - +// idxs are every kindManifest unit found in dir, so this is the directory's +// only writer. Units whose computation failed, or weren't dbt manifests, get +// their Result set here too and are left out of the sidecar. +func writeManifestGroup(dir string, idxs []int, units []unit, computations []manifestComputation, results []Result, version string) { + manifests := map[string]ManifestVersion{} + for _, i := range idxs { + path, c := units[i].path, computations[i] + r := Result{Kind: kindManifest, Path: path, Files: 1, Bytes: c.bytes} + switch { + case c.err != nil: + r.Err = c.err + case !c.isDbt: + r.Skipped = true + default: + r.Hash = c.hash + var slim *SlimManifest + if c.slimData != nil { + slim, r.Err = writeSlimManifest(dir, filepath.Base(path), c.slimData) + } + if r.Err == nil { + manifests[filepath.Base(path)] = ManifestVersion{ + Version: ProjectVersion{Algo: algoManifestJSON, Hash: c.hash}, + Slim: slim, + } + } } - if r.Err == nil { - r.Err = writeSidecar(dir, algoManifestJSON, hash, version, filtered) + r.Duration = time.Since(c.start) + results[i] = r + } + + if len(manifests) == 0 { + return + } + // The sidecar's top-level Version mirrors one manifest for a reader that + // doesn't yet look at Manifests; sorted keys make that choice stable + // across runs rather than whichever unit's goroutine finished last. + names := make([]string, 0, len(manifests)) + for name := range manifests { + names = append(names, name) + } + sort.Strings(names) + primary := manifests[names[0]] + + sidecarErr := writeSidecar(dir, primary.Version.Algo, primary.Version.Hash, version, manifests) + if sidecarErr == nil { + return + } + for _, i := range idxs { + if results[i].Err == nil && !results[i].Skipped { + results[i].Err = sidecarErr } } - return r } -// writeSlimManifest writes data as dir's slim manifest and returns the sidecar -// pointer describing it. data must already be marshaled, so a caller that later -// mutates the source doc cannot leak into it (see processManifest). -func writeSlimManifest(dir string, data []byte) (*FilteredManifest, error) { - if err := writeArtifact(dir, slimManifestName, data); err != nil { +// joinNotes appends add to existing, semicolon-separated. +func joinNotes(existing, add string) string { + if existing == "" { + return add + } + return existing + "; " + add +} + +// writeSlimManifest writes data as the slim companion of manifestFilename +// (see slimNameFor) inside dir, and returns the sidecar entry describing it. +// data must already be marshaled, so a caller that later mutates the source +// doc cannot leak into it. +func writeSlimManifest(dir, manifestFilename string, data []byte) (*SlimManifest, error) { + name := slimNameFor(manifestFilename) + if err := writeArtifact(dir, name, data); err != nil { return nil, err } - return &FilteredManifest{ + return &SlimManifest{ Schema: slimSchemaVersion, - Path: slimManifestName, + Path: name, Version: ProjectVersion{Algo: algoFilteredManifest, Hash: sha256Hex(data)}, }, nil } @@ -227,8 +325,8 @@ func (s Summary) CountFailed() int { return n } -// CountSkipped returns the number of units skipped (manifest.json files that aren't -// dbt manifests). +// CountSkipped returns the number of units skipped (manifest-like files that +// aren't dbt manifests). func (s Summary) CountSkipped() int { n := 0 for _, r := range s.Results { diff --git a/pkg/cosmosboost/precompute/precompute_test.go b/pkg/cosmosboost/precompute/precompute_test.go index 1177db052..702c38034 100644 --- a/pkg/cosmosboost/precompute/precompute_test.go +++ b/pkg/cosmosboost/precompute/precompute_test.go @@ -126,6 +126,236 @@ func TestRunSkipsNonDBTManifest(t *testing.T) { } } +// TestIsManifestCandidateName pins the name-matching rule discovery is built +// on: case-insensitive, "manifest" anywhere in a *.json filename. +func TestIsManifestCandidateName(t *testing.T) { + cases := map[string]bool{ + "manifest.json": true, + "manifest_full.json": true, + "MANIFEST.JSON": true, + "dbt_manifest_v2.json": true, + "run_results.json": false, + "catalog.json": false, + "manifest.txt": false, + "manifest": false, + } + for name, want := range cases { + if got := isManifestCandidateName(name); got != want { + t.Errorf("isManifestCandidateName(%q) = %v, want %v", name, got, want) + } + } +} + +// TestRunSharedDirectorySidecarDescribesBothManifests: two manifests sharing +// a directory outside a project write one shared sidecar, but each keeps its +// own entry (keyed by filename) in Manifests, with its own hash and slim +// pointer - neither is lost to the other. +func TestRunSharedDirectorySidecarDescribesBothManifests(t *testing.T) { + root := t.TempDir() + writeFiles(t, root, map[string]string{ + "shared/manifest_a.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{"model.a":{"name":"a"}}}`, + "shared/manifest_b.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{"model.b":{"name":"b"}}}`, + }) + + summary, err := Run([]string{root}, "test", Options{SlimManifest: true}) + if err != nil { + t.Fatal(err) + } + for _, r := range summary.Results { + if r.Err != nil { + t.Fatalf("unexpected error: %+v", r) + } + } + mustExist(t, filepath.Join(root, "shared", sidecarDir, "manifest_a.slim.json")) + mustExist(t, filepath.Join(root, "shared", sidecarDir, "manifest_b.slim.json")) + + var meta Metadata + readJSON(t, filepath.Join(root, "shared", sidecarDir, sidecarName), &meta) + if len(meta.Manifests) != 2 { + t.Fatalf("manifests = %+v, want 2 entries", meta.Manifests) + } + for name, want := range map[string]string{"manifest_a.json": "manifest_a.slim.json", "manifest_b.json": "manifest_b.slim.json"} { + entry, ok := meta.Manifests[name] + if !ok || entry.Version.Hash == "" || entry.Slim == nil || entry.Slim.Path != want { + t.Fatalf("Manifests[%q] = %+v, want a hash and slim path %q", name, entry, want) + } + } + if meta.Manifests["manifest_a.json"].Version.Hash == meta.Manifests["manifest_b.json"].Version.Hash { + t.Fatal("the two manifests' own hashes must differ (their content does)") + } + // The top-level Version is for a reader that doesn't yet look at + // Manifests; it should consistently pick the alphabetically-first + // manifest, not whichever goroutine happened to finish last. + for i := 0; i < 20; i++ { + if _, err := Run([]string{root}, "test", Options{SlimManifest: true}); err != nil { + t.Fatal(err) + } + readJSON(t, filepath.Join(root, "shared", sidecarDir, sidecarName), &meta) + if meta.Version.Hash != meta.Manifests["manifest_a.json"].Version.Hash { + t.Fatalf("top-level Version = %+v, want manifest_a.json's (alphabetically first)", meta.Version) + } + } +} + +// TestRunDoesNotFlagOtherDbtArtifactsAsManifests: run_results.json and +// catalog.json (which "dbt build"/"dbt docs generate" always produce beside +// manifest.json) don't contain "manifest" in their name; semantic_manifest.json +// does, but its schema URL isn't a manifest one, so content rejects it too. +func TestRunDoesNotFlagOtherDbtArtifactsAsManifests(t *testing.T) { + root := t.TempDir() + writeFiles(t, root, map[string]string{ + "proj/dbt_project.yml": "name: shop\n", + "proj/models/a.sql": "select 1", + "proj/target/manifest.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + "proj/target/run_results.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/run-results/v6.json"}}`, + "proj/target/catalog.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/catalog/v1.json"},"nodes":{}}`, + "proj/target/semantic_manifest.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/semantic-manifest/v1.json"}}`, + }) + + summary, err := Run([]string{root}, "test", Options{SlimManifest: true}) + if err != nil { + t.Fatal(err) + } + var succeeded, skipped int + for _, r := range summary.Results { + if r.Err != nil { + t.Fatalf("unexpected error: %+v", r) + } + if r.Skipped { + if filepath.Base(r.Path) != "semantic_manifest.json" { + t.Fatalf("unexpected skip: %+v", r) + } + skipped++ + continue + } + succeeded++ + } + if succeeded != 2 || skipped != 1 { // 1 project + 1 manifest succeed, semantic_manifest.json is skipped + t.Fatalf("want 2 successes + 1 skip, got %d successes, %d skipped: %+v", succeeded, skipped, summary.Results) + } + mustExist(t, filepath.Join(root, "proj", "target", sidecarDir, slimManifestName)) + if _, err := os.Stat(filepath.Join(root, "proj", "target", sidecarDir, "semantic_manifest.slim.json")); !os.IsNotExist(err) { + t.Fatalf("semantic_manifest.json must not be treated as a manifest: %v", err) + } +} + +// TestRunSlimsEveryManifestInProjectRoot: a project root with two manifests +// gets both slimmed, plus its own unaffected tree-hash sidecar. +func TestRunSlimsEveryManifestInProjectRoot(t *testing.T) { + root := t.TempDir() + writeFiles(t, root, map[string]string{ + "proj/dbt_project.yml": "name: shop\n", + "proj/models/a.sql": "select 1", + "proj/manifest.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + "proj/manifest_full.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + }) + + summary, err := Run([]string{root}, "test", Options{SlimManifest: true}) + if err != nil { + t.Fatal(err) + } + if len(summary.Results) != 1 || summary.Results[0].Kind != kindProject || summary.Results[0].Err != nil { + t.Fatalf("want 1 project-only result, got %+v", summary.Results) + } + mustExist(t, filepath.Join(root, "proj", sidecarDir, sidecarName)) + mustExist(t, filepath.Join(root, "proj", sidecarDir, slimManifestName)) + mustExist(t, filepath.Join(root, "proj", sidecarDir, "manifest_full.slim.json")) +} + +// TestRunProjectRootSlimFailureIsolatedToOneManifest: one candidate's slim +// write failing (its target path is pre-occupied by a directory, so this +// works even run as root) must not cost the project its tree-hash sidecar, +// or an earlier sibling its already-written slim file. +func TestRunProjectRootSlimFailureIsolatedToOneManifest(t *testing.T) { + root := t.TempDir() + writeFiles(t, root, map[string]string{ + "proj/dbt_project.yml": "name: shop\n", + "proj/models/a.sql": "select 1", + // Alphabetically first, so it's slimmed before the failing one. + "proj/manifest_a.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + "proj/manifest_b.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + }) + // manifest_b's slim write will fail: its target path is already a directory. + if err := os.MkdirAll(filepath.Join(root, "proj", sidecarDir, "manifest_b.slim.json"), 0o755); err != nil { + t.Fatal(err) + } + + summary, err := Run([]string{root}, "test", Options{SlimManifest: true}) + if err != nil { + t.Fatal(err) + } + if len(summary.Results) != 1 || summary.Results[0].Err != nil || summary.Results[0].Warning == "" { + t.Fatalf("want 1 successful project result with a warning, got %+v", summary.Results) + } + mustExist(t, filepath.Join(root, "proj", sidecarDir, sidecarName)) + mustExist(t, filepath.Join(root, "proj", sidecarDir, "manifest_a.slim.json")) + + var meta Metadata + readJSON(t, filepath.Join(root, "proj", sidecarDir, sidecarName), &meta) + if _, ok := meta.Manifests["manifest_a.json"]; !ok { + t.Fatalf("manifest_a.json missing from sidecar despite slimming successfully: %+v", meta.Manifests) + } + if _, ok := meta.Manifests["manifest_b.json"]; ok { + t.Fatalf("manifest_b.json should not be listed - its slim write failed: %+v", meta.Manifests) + } +} + +// TestRunWarnsOnUnreadableProjectRootManifest: a candidate in the project +// root that can't even be read (unlike one that reads fine but isn't a dbt +// manifest) must surface a warning, not vanish silently - the project still +// gets stamped on its own tree hash either way. +func TestRunWarnsOnUnreadableProjectRootManifest(t *testing.T) { + if runtime.GOOS == "windows" { + t.Skip("file permission semantics differ on windows") + } + root := t.TempDir() + writeFiles(t, root, map[string]string{ + "proj/dbt_project.yml": "name: shop\n", + "proj/models/a.sql": "select 1", + "proj/manifest_broken.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + }) + broken := filepath.Join(root, "proj", "manifest_broken.json") + if err := os.Chmod(broken, 0o000); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(broken, 0o644) }) + + summary, err := Run([]string{root}, "test", Options{SlimManifest: true}) + if err != nil { + t.Fatal(err) + } + if len(summary.Results) != 1 || summary.Results[0].Err != nil { + t.Fatalf("want 1 successful project result, got %+v", summary.Results) + } + if summary.Results[0].Warning == "" { + t.Fatal("want a warning naming the unreadable manifest, got none") + } + mustExist(t, filepath.Join(root, "proj", sidecarDir, sidecarName)) +} + +// TestRunHandlesMultipleProjectsWithDifferentManifestNames: three dbt +// projects, each naming its manifest differently, all get discovered and +// slimmed with no configuration. +func TestRunHandlesMultipleProjectsWithDifferentManifestNames(t *testing.T) { + root := t.TempDir() + writeFiles(t, root, map[string]string{ + "dbt1/dbt_project.yml": "name: one\n", + "dbt1/manifest_custom.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + "dbt2/dbt_project.yml": "name: two\n", + "dbt2/manifest_by_run.json": `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}`, + }) + + summary, err := Run([]string{root}, "test", Options{SlimManifest: true}) + if err != nil { + t.Fatal(err) + } + if len(summary.Results) != 2 { + t.Fatalf("want 2 project results, got %+v", summary.Results) + } + mustExist(t, filepath.Join(root, "dbt1", sidecarDir, "manifest_custom.slim.json")) + mustExist(t, filepath.Join(root, "dbt2", sidecarDir, "manifest_by_run.slim.json")) +} + // TestRunWarnsOnTemplatedPackagesPath verifies a project whose packages-install-path // is a Jinja template still gets stamped, but the Result carries a non-fatal warning. func TestRunWarnsOnTemplatedPackagesPath(t *testing.T) { @@ -457,10 +687,8 @@ func TestRunReportsFailedUnits(t *testing.T) { // TestRunWritesSlimManifestAlongsideSidecar: with the option set, every // discovered manifest.json gets a slim copy next to its hash sidecar, and the -// sidecar carries the filtered_manifest pointer at it. The pointer hashes the -// slim file's own bytes, so a consumer can confirm the two are a matched pair - -// which adjacency alone no longer implies, now that cleanup judges each -// artifact independently. +// sidecar's Manifests entry for it carries a pointer hashing the slim file's +// own bytes, so a consumer can confirm the two are a matched pair. func TestRunWritesSlimManifestAlongsideSidecar(t *testing.T) { root := t.TempDir() writeFiles(t, root, map[string]string{ @@ -484,22 +712,23 @@ func TestRunWritesSlimManifestAlongsideSidecar(t *testing.T) { var meta Metadata readJSON(t, filepath.Join(astroDir, sidecarName), &meta) - fm := meta.FilteredManifest - if fm == nil { - t.Fatalf("sidecar has no filtered_manifest section: %+v", meta) + entry, ok := meta.Manifests["manifest.json"] + fm := entry.Slim + if !ok || fm == nil { + t.Fatalf("sidecar has no slim_manifest for manifest.json: %+v", meta) } if fm.Path != slimManifestName || fm.Schema != slimSchemaVersion || fm.Version.Algo != algoFilteredManifest { - t.Fatalf("filtered_manifest = %+v", fm) + t.Fatalf("slim_manifest = %+v", fm) } slimData, err := os.ReadFile(filepath.Join(astroDir, slimManifestName)) if err != nil { t.Fatal(err) } if want := sha256Hex(slimData); fm.Version.Hash != want { - t.Fatalf("filtered_manifest hash = %q, want the slim file's own hash %q", fm.Version.Hash, want) + t.Fatalf("slim_manifest hash = %q, want the slim file's own hash %q", fm.Version.Hash, want) } if fm.Version.Hash == meta.Version.Hash { - t.Fatal("filtered_manifest hash must not be the full manifest's hash") + t.Fatal("slim_manifest hash must not be the full manifest's hash") } } @@ -524,8 +753,8 @@ func TestRunSkipsSlimManifestWhenNotRequested(t *testing.T) { if err != nil { t.Fatal(err) } - if strings.Contains(string(raw), "filtered_manifest") { - t.Fatalf("filtered_manifest emitted with the slim manifest disabled: %s", raw) + if strings.Contains(string(raw), "slim_manifest") { + t.Fatalf("slim_manifest emitted with the slim manifest disabled: %s", raw) } } @@ -563,8 +792,9 @@ func TestRunSlimsManifestInProjectRoot(t *testing.T) { if meta.Version.Algo != algoProjectTree { t.Fatalf("project sidecar algo = %q, want %q", meta.Version.Algo, algoProjectTree) } - if meta.FilteredManifest == nil || meta.FilteredManifest.Path != slimManifestName { - t.Fatalf("project sidecar does not point at the slim manifest: %+v", meta.FilteredManifest) + entry, ok := meta.Manifests["manifest.json"] + if !ok || entry.Slim == nil || entry.Slim.Path != slimManifestName { + t.Fatalf("project sidecar does not point at the slim manifest: %+v", meta.Manifests) } } @@ -612,8 +842,8 @@ func TestRunSkipsNonDbtManifestInProjectRoot(t *testing.T) { } var meta Metadata readJSON(t, filepath.Join(astroDir, sidecarName), &meta) - if meta.FilteredManifest != nil { - t.Fatalf("project sidecar points at a slim manifest that was never written: %+v", meta.FilteredManifest) + if len(meta.Manifests) != 0 { + t.Fatalf("project sidecar lists a manifest that was never slimmed: %+v", meta.Manifests) } } diff --git a/pkg/cosmosboost/precompute/slim.go b/pkg/cosmosboost/precompute/slim.go index c202c95ac..fb8201165 100644 --- a/pkg/cosmosboost/precompute/slim.go +++ b/pkg/cosmosboost/precompute/slim.go @@ -1,11 +1,28 @@ package precompute -const slimManifestName = "manifest.slim.json" +import "strings" + +// slimManifestSuffix marks a slim companion file, and lets Cleanup recognize +// one regardless of which manifest it came from. +const slimManifestSuffix = ".slim.json" + +// slimManifestName is the slim companion for the default "manifest.json". +const slimManifestName = "manifest" + slimManifestSuffix // slimSchemaVersion identifies the allowlist that produced a slim manifest, so // a reader that doesn't recognize it can fall back to the full manifest. const slimSchemaVersion = 1 +// slimNameFor returns manifestFilename's slim companion name, e.g. +// "manifest_full.json" -> "manifest_full.slim.json". The trim is +// case-insensitive, matching isManifestCandidateName. +func slimNameFor(manifestFilename string) string { + if len(manifestFilename) >= 5 && strings.EqualFold(manifestFilename[len(manifestFilename)-5:], ".json") { + manifestFilename = manifestFilename[:len(manifestFilename)-5] + } + return manifestFilename + slimManifestSuffix +} + // slimSections are the only top-level collections Cosmos loads nodes from // (cosmos/dbt/graph.py::_load_nodes_from_manifest_data). var slimSections = []string{"nodes", "sources", "exposures"} diff --git a/pkg/cosmosboost/precompute/slim_test.go b/pkg/cosmosboost/precompute/slim_test.go index 02d9dbe88..b60c0cd4b 100644 --- a/pkg/cosmosboost/precompute/slim_test.go +++ b/pkg/cosmosboost/precompute/slim_test.go @@ -17,6 +17,26 @@ func parseDoc(t *testing.T, raw string) map[string]any { return doc } +// TestSlimNameFor pins the naming convention a consumer needs to compute a +// manifest's slim companion from its own filename alone, with no sidecar +// lookup required. +func TestSlimNameFor(t *testing.T) { + cases := map[string]string{ + "manifest.json": "manifest.slim.json", + "manifest_full.json": "manifest_full.slim.json", + "manifest_by_schedule.json": "manifest_by_schedule.slim.json", + "MANIFEST.JSON": "MANIFEST.slim.json", // matches isManifestCandidateName's case-insensitivity + } + for in, want := range cases { + if got := slimNameFor(in); got != want { + t.Errorf("slimNameFor(%q) = %q, want %q", in, got, want) + } + } + if slimNameFor("manifest.json") != slimManifestName { + t.Errorf("slimNameFor(%q) must equal slimManifestName, the default case's constant", "manifest.json") + } +} + // TestBuildSlimManifestKeepsOnlyAllowedResourceFields pins the field allowlist // Cosmos actually reads from a manifest node - both the DbtNode build // (cosmos/dbt/graph.py::_build_dbt_node_from_manifest_resource) and the outlet diff --git a/pkg/cosmosboost/predeploy.go b/pkg/cosmosboost/predeploy.go index dff788924..1cb237ada 100644 --- a/pkg/cosmosboost/predeploy.go +++ b/pkg/cosmosboost/predeploy.go @@ -32,7 +32,7 @@ func slimManifestEnabled() bool { // (dbt_project.yml) gets a .astro/dbt_metadata.json sidecar carrying its // content hash, which the Cosmos Boost plugin uses as a cache version key at // parse time instead of hashing the project tree itself. Every standalone dbt -// manifest.json gets a hash sidecar too, plus a slim, field-filtered copy for +// manifest gets a hash sidecar too, plus a slim, field-filtered copy for // the plugin to load in place of the full manifest at DAG-parse time. func PreDeploy(path string) error { opts := precompute.Options{SlimManifest: slimManifestEnabled()} diff --git a/pkg/cosmosboost/predeploy_test.go b/pkg/cosmosboost/predeploy_test.go index 7398abf14..0385c3efd 100644 --- a/pkg/cosmosboost/predeploy_test.go +++ b/pkg/cosmosboost/predeploy_test.go @@ -45,7 +45,7 @@ func TestPreDeployWritesArtifact(t *testing.T) { } `json:"generated_by"` } require.NoError(t, json.Unmarshal(data, &meta)) - require.Equal(t, 1, meta.Schema, "schema is the plugin's compatibility gate and must stay 1") + require.Equal(t, 2, meta.Schema, "schema is the plugin's compatibility gate - bump deliberately, not by accident") require.NotEmpty(t, meta.Version.Hash, "version.hash is what the plugin consumes") require.NotEmpty(t, meta.Version.Algo) require.Equal(t, "astro", meta.GeneratedBy.Application) @@ -101,6 +101,19 @@ func TestSlimManifestEnabled(t *testing.T) { } } +// TestPreDeployDiscoversCustomNamedManifest: a manifest not literally named +// manifest.json is found and stamped automatically - by name and content, +// with nothing to configure. +func TestPreDeployDiscoversCustomNamedManifest(t *testing.T) { + dir := t.TempDir() + manifest := `{"metadata":{"dbt_schema_version":"https://schemas.getdbt.com/dbt/manifest/v12.json"},"nodes":{}}` + require.NoError(t, os.WriteFile(filepath.Join(dir, "manifest_full.json"), []byte(manifest), 0o644)) + + require.NoError(t, PreDeploy(dir)) + require.FileExists(t, filepath.Join(dir, artifactRelPath)) + require.FileExists(t, filepath.Join(dir, ".astro", "manifest_full.slim.json")) +} + func TestPreDeployNoDbtContentIsANoOp(t *testing.T) { dir := t.TempDir() require.NoError(t, os.WriteFile(filepath.Join(dir, "app.py"), []byte("print('hi')"), 0o644))