Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pkg/airflowrt/include/gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -15,4 +15,4 @@ airflow.db
.astro/*.local.yaml

**/.astro/dbt_metadata.json
**/.astro/manifest.slim.json
**/.astro/*manifest*.slim.json
16 changes: 9 additions & 7 deletions pkg/cosmosboost/precompute/cleanup.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"io/fs"
"os"
"path/filepath"
"strings"
"time"
)

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
51 changes: 51 additions & 0 deletions pkg/cosmosboost/precompute/cleanup_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 -
Expand Down
23 changes: 16 additions & 7 deletions pkg/cosmosboost/precompute/discover.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//
Expand Down Expand Up @@ -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

Expand All @@ -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)] {
Expand Down
15 changes: 11 additions & 4 deletions pkg/cosmosboost/precompute/hash.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"os"
"path/filepath"
"sort"
"strings"
)

// excludedDirs are directory names skipped during *project* discovery
Expand Down Expand Up @@ -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})
Expand Down Expand Up @@ -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
Expand Down
16 changes: 9 additions & 7 deletions pkg/cosmosboost/precompute/hash_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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{
Expand Down
47 changes: 24 additions & 23 deletions pkg/cosmosboost/precompute/metadata.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
pankajastro marked this conversation as resolved.

// algoProjectTree hashes a whole dbt project directory (source files).
// v2 excludes .git, which never ships in a deploy payload, so VCS activity
Expand All @@ -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"`
Expand All @@ -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.
Expand Down
Loading
Loading