From 0f2998db59833f6173bbf9676d2d5a257e8636e5 Mon Sep 17 00:00:00 2001 From: eitsupi Date: Sat, 29 Aug 2026 14:41:19 +0000 Subject: [PATCH 1/2] perf(manifest)!: cache compact DAG analysis --- .../dlin-core/src/parser/manifest/report.rs | 6 +- crates/dlin-core/src/parser/manifest_cache.rs | 413 +++++++++--------- crates/dlin/src/commands/graph.rs | 14 +- crates/dlin/src/commands/manifest.rs | 46 +- crates/dlin/src/commands/shared.rs | 157 +++---- .../dlin/tests/integration/manifest_only.rs | 119 +++++ 6 files changed, 436 insertions(+), 319 deletions(-) diff --git a/crates/dlin-core/src/parser/manifest/report.rs b/crates/dlin-core/src/parser/manifest/report.rs index 93eae33..ed572c2 100644 --- a/crates/dlin-core/src/parser/manifest/report.rs +++ b/crates/dlin-core/src/parser/manifest/report.rs @@ -31,7 +31,7 @@ pub struct ManifestCapabilities { /// A stable, machine-readable diagnostic produced while loading a manifest. #[non_exhaustive] -#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] #[serde(rename_all = "snake_case")] pub enum ManifestDiagnosticKind { ParseError, @@ -83,7 +83,7 @@ impl std::fmt::Display for ManifestDiagnosticKind { /// Severity of a manifest load diagnostic. #[non_exhaustive] -#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] #[serde(rename_all = "snake_case")] pub enum ManifestDiagnosticSeverity { Error, @@ -101,7 +101,7 @@ impl ManifestDiagnosticSeverity { /// Diagnostic details retained by [`ManifestLoadReport`]. #[non_exhaustive] -#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub struct ManifestDiagnostic { pub kind: ManifestDiagnosticKind, pub severity: ManifestDiagnosticSeverity, diff --git a/crates/dlin-core/src/parser/manifest_cache.rs b/crates/dlin-core/src/parser/manifest_cache.rs index 7151e80..15d56f7 100644 --- a/crates/dlin-core/src/parser/manifest_cache.rs +++ b/crates/dlin-core/src/parser/manifest_cache.rs @@ -1,46 +1,93 @@ use std::path::{Path, PathBuf}; -use std::time::SystemTime; use serde::{Deserialize, Serialize}; use crate::graph::types::LineageGraph; +use crate::parser::manifest::ManifestDiagnostic; const CACHE_DIR: &str = ".dlin_cache"; const CACHE_FILENAME: &str = "manifest_graph_cache.json"; -#[derive(Debug, Serialize, Deserialize)] +/// Compact model-level information persisted by the manifest cache. +/// +/// The typed manifest is intentionally excluded. Consumers that need compiled +/// SQL can parse the current artifact lazily after restoring this analysis. +#[non_exhaustive] +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ManifestAnalysis { + graph: LineageGraph, + diagnostics: Vec, + project_name: Option, + referenced_file_paths: Vec, +} + +impl ManifestAnalysis { + pub fn new( + graph: LineageGraph, + diagnostics: Vec, + project_name: Option, + referenced_file_paths: Vec, + ) -> Self { + Self { + graph, + diagnostics, + project_name, + referenced_file_paths, + } + } + + pub fn into_parts( + self, + ) -> ( + LineageGraph, + Vec, + Option, + Vec, + ) { + ( + self.graph, + self.diagnostics, + self.project_name, + self.referenced_file_paths, + ) + } +} + +#[derive(Debug, Deserialize)] struct ManifestCacheFile { #[serde(default)] version: String, + #[serde(default)] entry: Option, } -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Deserialize)] struct ManifestCacheEntry { - #[serde(default)] - manifest_identity: String, - mtime_secs: u64, - #[serde(default)] - mtime_nanos: u32, - file_size: u64, - #[serde(default)] - content_hash: u64, - graph: LineageGraph, + input_hash: u64, + analysis: ManifestAnalysis, +} + +#[derive(Serialize)] +struct ManifestCacheFileRef<'a> { + version: &'a str, + entry: Option>, +} + +#[derive(Serialize)] +struct ManifestCacheEntryRef<'a> { + input_hash: u64, + analysis: &'a ManifestAnalysis, } -pub struct ManifestGraphCache { +/// Persistent cache for model-level manifest analysis. +pub struct ManifestAnalysisCache { version: String, entry: Option, cache_path: Option, dirty: bool, } -// `CARGO_PKG_VERSION`, used by `load` and `fresh`, is the compatibility and -// invalidation boundary for both cache format and graph semantics. We avoid a -// second cache schema/semantics version: a release/version bump automatically -// discards older caches. Compatibility between development builds sharing a -// package version is not guaranteed; use `--refresh-cache` when needed. -impl ManifestGraphCache { +impl ManifestAnalysisCache { pub fn disabled() -> Self { Self { version: String::new(), @@ -51,17 +98,13 @@ impl ManifestGraphCache { } pub fn load(project_dir: &Path, cache_dir: Option<&Path>) -> Self { - let cache_path = match cache_dir { - Some(dir) => dir.join(CACHE_FILENAME), - None => project_dir.join(CACHE_DIR).join(CACHE_FILENAME), - }; + let cache_path = cache_path(project_dir, cache_dir); let version = env!("CARGO_PKG_VERSION").to_string(); let entry = std::fs::read_to_string(&cache_path) .ok() .and_then(|content| serde_json::from_str::(&content).ok()) - .filter(|cf| cf.version == version) - .and_then(|cf| cf.entry); - + .filter(|file| file.version == version) + .and_then(|file| file.entry); Self { version, entry, @@ -71,170 +114,100 @@ impl ManifestGraphCache { } pub fn fresh(project_dir: &Path, cache_dir: Option<&Path>) -> Self { - let cache_path = match cache_dir { - Some(dir) => dir.join(CACHE_FILENAME), - None => project_dir.join(CACHE_DIR).join(CACHE_FILENAME), - }; Self { version: env!("CARGO_PKG_VERSION").to_string(), entry: None, - cache_path: Some(cache_path), + cache_path: Some(cache_path(project_dir, cache_dir)), dirty: false, } } - pub fn get(&self, manifest_path: &Path) -> Option<&LineageGraph> { - let entry = self.entry.as_ref()?; - let stat = file_stat(manifest_path)?; - if entry_matches_fingerprint( - entry, - manifest_path, - ( - stat.mtime_secs, - stat.mtime_nanos, - stat.file_size, - stat.content_hash, - ), - ) { - Some(&entry.graph) - } else { - None - } - } - - /// Look up a cached graph using metadata and content hash computed by the - /// caller. This avoids rereading the manifest when its bytes were already - /// loaded for parsing and fingerprinting. - pub fn get_with_fingerprint( - &self, - manifest_path: &Path, - fingerprint: (u64, u32, u64, u64), - ) -> Option<&LineageGraph> { - let entry = self.entry.as_ref()?; - entry_matches_fingerprint(entry, manifest_path, fingerprint).then_some(&entry.graph) + /// Take the candidate entry for this manifest. Since this cache stores a + /// single slot, a miss also discards the obsolete candidate before a + /// replacement is inserted. The ownership-taking API keeps warm command + /// paths from cloning the graph and diagnostics. + pub fn take_for_manifest(&mut self, manifest_bytes: &[u8]) -> Option { + self.cache_path.as_ref()?; + let input_hash = hash_manifest_bytes(manifest_bytes); + self.entry + .take() + .filter(|entry| entry.input_hash == input_hash) + .map(|entry| entry.analysis) } - pub fn insert_if_fingerprint_matches( - &mut self, - manifest_path: &Path, - graph: &LineageGraph, - expected: (u64, u32, u64, u64), - ) -> bool { - let Some(stat) = file_stat(manifest_path) else { - return false; - }; - if ( - stat.mtime_secs, - stat.mtime_nanos, - stat.file_size, - stat.content_hash, - ) != expected - { - return false; - } + pub fn insert_for_manifest(&mut self, manifest_bytes: &[u8], analysis: ManifestAnalysis) { + let input_hash = hash_manifest_bytes(manifest_bytes); self.entry = Some(ManifestCacheEntry { - manifest_identity: manifest_identity(manifest_path), - mtime_secs: stat.mtime_secs, - mtime_nanos: stat.mtime_nanos, - file_size: stat.file_size, - content_hash: stat.content_hash, - graph: graph.clone(), + input_hash, + analysis, }); self.dirty = true; - true } pub fn save(&self) { - let cache_path = match (&self.cache_path, self.dirty) { - (Some(p), true) => p, - _ => return, + let Some(cache_path) = self.cache_path.as_ref().filter(|_| self.dirty) else { + return; }; - let cf = ManifestCacheFile { - version: self.version.clone(), - entry: self.entry.clone(), + let file = ManifestCacheFileRef { + version: &self.version, + entry: self.entry.as_ref().map(|entry| ManifestCacheEntryRef { + input_hash: entry.input_hash, + analysis: &entry.analysis, + }), }; - if let Some(parent) = cache_path.parent() { - if std::fs::create_dir_all(parent).is_err() { - crate::warn!("could not create cache directory: {}", parent.display()); - return; - } - let gitignore = parent.join(".gitignore"); - if !gitignore.exists() - && let Err(e) = std::fs::write(&gitignore, "# Automatically created by dlin\n*\n") - { - crate::warn!("could not create {}: {}", gitignore.display(), e); - } + let Some(parent) = cache_path.parent() else { + return; + }; + if let Err(error) = std::fs::create_dir_all(parent) { + crate::warn!( + "could not create cache directory {}: {}", + parent.display(), + error + ); + return; + } + let gitignore = parent.join(".gitignore"); + if !gitignore.exists() + && let Err(error) = std::fs::write(&gitignore, "# Automatically created by dlin\n*\n") + { + crate::warn!("could not create {}: {}", gitignore.display(), error); } - match serde_json::to_string(&cf) { - Ok(json) => { - if let Err(e) = std::fs::write(cache_path, json) { - crate::warn!("could not write cache file {}: {}", cache_path.display(), e); - } - } - Err(e) => { - crate::warn!("could not serialize manifest graph cache: {}", e); - } + let Ok(json) = serde_json::to_string(&file) else { + crate::warn!("could not serialize manifest analysis cache"); + return; + }; + // Keep the active file valid if serialization is interrupted. + let temporary = cache_path.with_extension("json.tmp"); + if let Err(error) = std::fs::write(&temporary, json).and_then(|()| { + #[cfg(windows)] + let _ = std::fs::remove_file(cache_path); + std::fs::rename(&temporary, cache_path) + }) { + crate::warn!( + "could not write cache file {}: {}", + cache_path.display(), + error + ); + let _ = std::fs::remove_file(temporary); } } } -struct FileStat { - mtime_secs: u64, - mtime_nanos: u32, - file_size: u64, - content_hash: u64, -} - -fn file_stat(path: &Path) -> Option { - let meta = std::fs::metadata(path).ok()?; - let content = std::fs::read(path).ok()?; - let mtime_secs = meta - .modified() - .ok()? - .duration_since(SystemTime::UNIX_EPOCH) - .ok()? - .as_secs(); - let mtime_nanos = meta - .modified() - .ok()? - .duration_since(SystemTime::UNIX_EPOCH) - .ok()? - .subsec_nanos(); - Some(FileStat { - mtime_secs, - mtime_nanos, - file_size: meta.len(), - content_hash: hash_bytes(&content), - }) -} - -fn entry_matches_fingerprint( - entry: &ManifestCacheEntry, - manifest_path: &Path, - fingerprint: (u64, u32, u64, u64), -) -> bool { - entry.manifest_identity == manifest_identity(manifest_path) - && ( - entry.mtime_secs, - entry.mtime_nanos, - entry.file_size, - entry.content_hash, - ) == fingerprint +fn cache_path(project_dir: &Path, cache_dir: Option<&Path>) -> PathBuf { + cache_dir + .map(Path::to_path_buf) + .unwrap_or_else(|| project_dir.join(CACHE_DIR)) + .join(CACHE_FILENAME) } -fn manifest_identity(path: &Path) -> String { - path.canonicalize() - .unwrap_or_else(|_| path.to_path_buf()) - .to_string_lossy() - .to_string() -} - -fn hash_bytes(bytes: &[u8]) -> u64 { +/// Deterministic content hash used to validate the current manifest bytes. +/// The package version remains the sole cache compatibility boundary. +fn hash_manifest_bytes(bytes: &[u8]) -> u64 { const FNV_OFFSET_BASIS: u64 = 0xcbf29ce484222325; const FNV_PRIME: u64 = 0x100000001b3; let mut hash = FNV_OFFSET_BASIS; - for &b in bytes { - hash ^= b as u64; + for &byte in bytes { + hash ^= byte as u64; hash = hash.wrapping_mul(FNV_PRIME); } hash @@ -243,42 +216,92 @@ fn hash_bytes(bytes: &[u8]) -> u64 { #[cfg(test)] mod tests { use super::*; - use std::time::SystemTime; + use crate::parser::manifest::{ManifestDiagnosticKind, ManifestDiagnosticSeverity}; + + fn analysis() -> ManifestAnalysis { + ManifestAnalysis::new( + LineageGraph::new(), + vec![ManifestDiagnostic { + kind: ManifestDiagnosticKind::MissingSchemaVersion, + severity: ManifestDiagnosticSeverity::Warning, + message: "missing".to_string(), + hint: None, + raw_resource: None, + raw_type: None, + schema: None, + }], + Some("project".to_string()), + vec!["models/orders.sql".to_string()], + ) + } + + #[test] + fn content_hash_lookup_does_not_use_path_or_metadata() { + let temp = tempfile::tempdir().unwrap(); + let mut cache = ManifestAnalysisCache::fresh(temp.path(), None); + cache.insert_for_manifest(b"a", analysis()); + assert!(cache.take_for_manifest(b"a").is_some()); + } + + #[test] + fn mismatched_slot_is_a_miss() { + let temp = tempfile::tempdir().unwrap(); + let mut cache = ManifestAnalysisCache::fresh(temp.path(), None); + cache.insert_for_manifest(b"a", analysis()); + assert!(cache.take_for_manifest(b"b").is_none()); + } + + #[test] + fn compact_analysis_round_trips_through_persistence() { + let temp = tempfile::tempdir().unwrap(); + let mut cache = ManifestAnalysisCache::fresh(temp.path(), None); + cache.insert_for_manifest(b"manifest", analysis()); + cache.save(); + + let mut loaded = ManifestAnalysisCache::load(temp.path(), None); + let restored = loaded.take_for_manifest(b"manifest").unwrap(); + let (_, diagnostics, project_name, paths) = restored.into_parts(); + assert_eq!(project_name.as_deref(), Some("project")); + assert_eq!(paths, ["models/orders.sql"]); + assert_eq!(diagnostics.len(), 1); + } + + #[test] + fn corrupt_cache_fails_open() { + let temp = tempfile::tempdir().unwrap(); + let path = temp.path().join(CACHE_DIR).join(CACHE_FILENAME); + std::fs::create_dir_all(path.parent().unwrap()).unwrap(); + std::fs::write(path, b"not json").unwrap(); + let mut loaded = ManifestAnalysisCache::load(temp.path(), None); + assert!(loaded.take_for_manifest(b"manifest").is_none()); + } + + #[test] + fn disabled_cache_skips_lookup_and_persistent_io() { + let mut cache = ManifestAnalysisCache::disabled(); + assert!(cache.take_for_manifest(b"manifest").is_none()); + cache.insert_for_manifest(b"manifest", analysis()); + cache.save(); + } #[test] - fn fingerprint_lookup_uses_caller_fingerprint_without_reread() { + fn package_version_mismatch_fails_open() { let temp = tempfile::tempdir().unwrap(); - let manifest_path = temp.path().join("manifest.json"); - let bytes = br#"{}"#; - std::fs::write(&manifest_path, bytes).unwrap(); - let metadata = std::fs::metadata(&manifest_path).unwrap(); - let modified = metadata - .modified() - .unwrap() - .duration_since(SystemTime::UNIX_EPOCH) - .unwrap(); - let fingerprint = ( - modified.as_secs(), - modified.subsec_nanos(), - metadata.len(), - hash_bytes(bytes), - ); + let mut cache = ManifestAnalysisCache::fresh(temp.path(), None); + cache.insert_for_manifest(b"manifest", analysis()); + cache.save(); - let mut cache = ManifestGraphCache::fresh(temp.path(), None); - let graph = LineageGraph::new(); - assert!(cache.insert_if_fingerprint_matches(&manifest_path, &graph, fingerprint)); - // Keep the path identity stable while changing the on-disk content. - std::fs::write(&manifest_path, br#"{"changed":true}"#).unwrap(); + let path = temp.path().join(CACHE_DIR).join(CACHE_FILENAME); + let mut file: serde_json::Value = + serde_json::from_str(&std::fs::read_to_string(path).unwrap()).unwrap(); + file["version"] = serde_json::Value::String("different-release".to_string()); + std::fs::write( + temp.path().join(CACHE_DIR).join(CACHE_FILENAME), + serde_json::to_vec(&file).unwrap(), + ) + .unwrap(); - // The caller-computed old fingerprint is trusted and does not reread - // the changed manifest. - assert!( - cache - .get_with_fingerprint(&manifest_path, fingerprint) - .is_some() - ); - // The compatibility API reads and hashes the current content, so it - // correctly misses the stale entry. - assert!(cache.get(&manifest_path).is_none()); + let mut loaded = ManifestAnalysisCache::load(temp.path(), None); + assert!(loaded.take_for_manifest(b"manifest").is_none()); } } diff --git a/crates/dlin/src/commands/graph.rs b/crates/dlin/src/commands/graph.rs index 277babb..04f3694 100644 --- a/crates/dlin/src/commands/graph.rs +++ b/crates/dlin/src/commands/graph.rs @@ -20,7 +20,7 @@ pub(crate) fn run_graph_command(args: GraphArgs) -> Result<()> { // Validate flag combinations before building DAG validate_source_flags(&args.source, args.manifest_path.as_ref())?; - let (dag, project_opt, manifest_diagnostics, manifest_opt) = build_dag( + let (dag, project_opt, manifest_diagnostics, manifest_context) = build_dag( &project_dir, &args.source, args.manifest_path.as_ref(), @@ -128,7 +128,9 @@ pub(crate) fn run_graph_command(args: GraphArgs) -> Result<()> { Some(collect_sql_contents_for_source( &args.source, &project_dir, - manifest_opt.as_ref(), + manifest_context + .as_ref() + .and_then(|context| context.manifest.as_ref()), args.manifest_path.as_ref(), &filtered, )) @@ -193,7 +195,7 @@ pub(crate) fn run_list_command(args: ListArgs) -> Result<()> { validate_source_flags(&args.source, args.manifest_path.as_ref())?; - let (dag, project_opt, manifest_diagnostics, manifest_opt) = build_dag( + let (dag, project_opt, manifest_diagnostics, manifest_context) = build_dag( &project_dir, &args.source, args.manifest_path.as_ref(), @@ -280,7 +282,9 @@ pub(crate) fn run_list_command(args: ListArgs) -> Result<()> { Some(collect_sql_contents_for_source( &args.source, &project_dir, - manifest_opt.as_ref(), + manifest_context + .as_ref() + .and_then(|context| context.manifest.as_ref()), args.manifest_path.as_ref(), &filtered, )) @@ -395,7 +399,7 @@ pub(crate) fn run_impact_command( .unwrap_or_else(|_| project_dir.to_path_buf()); validate_source_flags(source, manifest_path)?; - let (dag, project_opt, manifest_diagnostics, _manifest) = build_dag( + let (dag, project_opt, manifest_diagnostics, _manifest_context) = build_dag( &project_dir, source, manifest_path, diff --git a/crates/dlin/src/commands/manifest.rs b/crates/dlin/src/commands/manifest.rs index 342379e..c18b0ce 100644 --- a/crates/dlin/src/commands/manifest.rs +++ b/crates/dlin/src/commands/manifest.rs @@ -22,7 +22,7 @@ pub(crate) fn run_summary_command(args: SummaryArgs) -> Result<()> { validate_source_flags(&args.source, args.manifest_path.as_ref())?; - let (dag, project_opt, manifest_diagnostics, manifest_opt) = build_dag( + let (dag, project_opt, manifest_diagnostics, manifest_context) = build_dag( &project_dir, &args.source, args.manifest_path.as_ref(), @@ -34,16 +34,18 @@ pub(crate) fn run_summary_command(args: SummaryArgs) -> Result<()> { let (project_name, vars_count, manifest_status) = match args.source { SourceType::Manifest => { - let name = manifest_opt + let name = manifest_context .as_ref() - .and_then(|manifest| manifest.metadata.project_name.clone()) + .and_then(|context| context.project_name.clone()) .unwrap_or_else(|| "(unknown)".to_string()); let status = match parser::project::DbtProject::load(&project_dir) { - Ok(project) => check_manifest_freshness( + Ok(project) => check_manifest_freshness_with_referenced_paths( &project_dir, args.manifest_path.as_ref(), &project, - manifest_opt.as_ref(), + manifest_context + .as_ref() + .map(|context| context.referenced_file_paths.as_slice()), ), Err(e) => { let is_not_found = e @@ -99,10 +101,18 @@ fn find_deleted_manifest_files( manifest: &parser::manifest::Manifest, project_dir: &Path, ) -> Vec { - let mut deleted: Vec = manifest - .collect_file_paths() + let paths = manifest.collect_file_paths(); + find_deleted_manifest_paths(paths.iter(), project_dir) +} + +fn find_deleted_manifest_paths<'a>( + paths: impl IntoIterator, + project_dir: &Path, +) -> Vec { + let mut deleted: Vec = paths .into_iter() .filter(|p| !project_dir.join(p).exists()) + .cloned() .collect(); deleted.sort(); deleted @@ -164,6 +174,24 @@ pub(crate) fn check_manifest_freshness( manifest_path: Option<&PathBuf>, project: &parser::project::DbtProject, manifest: Option<&parser::manifest::Manifest>, +) -> Option { + let paths = manifest.map(|value| value.collect_file_paths().into_iter().collect::>()); + check_manifest_freshness_with_referenced_paths( + project_dir, + manifest_path, + project, + paths.as_deref(), + ) +} + +/// Check manifest freshness using compact referenced paths from the model-level +/// cache. This avoids reparsing the typed manifest on warm summary paths. +#[cfg(not(tarpaulin_include))] +pub(crate) fn check_manifest_freshness_with_referenced_paths( + project_dir: &Path, + manifest_path: Option<&PathBuf>, + project: &parser::project::DbtProject, + referenced_file_paths: Option<&[String]>, ) -> Option { let not_found = render::summary::ManifestStatus { found: false, @@ -199,8 +227,8 @@ pub(crate) fn check_manifest_freshness( stale_files.sort(); // Check for deleted files: paths referenced in manifest but missing on disk - let deleted_files = match manifest { - Some(manifest) => find_deleted_manifest_files(manifest, project_dir), + let deleted_files = match referenced_file_paths { + Some(paths) => find_deleted_manifest_paths(paths.iter(), project_dir), None => match parser::manifest::load_manifest(&resolved) { Ok(manifest) => find_deleted_manifest_files(&manifest, project_dir), Err(e) => { diff --git a/crates/dlin/src/commands/shared.rs b/crates/dlin/src/commands/shared.rs index c933b21..c92daf9 100644 --- a/crates/dlin/src/commands/shared.rs +++ b/crates/dlin/src/commands/shared.rs @@ -1,5 +1,4 @@ use std::path::{Path, PathBuf}; -use std::time::SystemTime; use anyhow::Result; @@ -8,12 +7,19 @@ use dlin_core::graph; use dlin_core::graph::column_lineage::{DialectClassification, DlinDialect}; use dlin_core::parser; +#[cfg(not(tarpaulin_include))] +pub(crate) struct ManifestDagContext { + pub(crate) manifest: Option, + pub(crate) project_name: Option, + pub(crate) referenced_file_paths: Vec, +} + #[cfg(not(tarpaulin_include))] type DagBuildResult = Result<( graph::types::LineageGraph, Option, Vec, - Option, + Option, )>; /// Build the lineage DAG from either a manifest file or by parsing SQL files. @@ -30,75 +36,60 @@ pub(crate) fn build_dag( SourceType::Manifest => { let resolved = resolve_manifest_path_or_default(manifest_path, project_dir)?; let mut cache = if no_cache { - parser::manifest_cache::ManifestGraphCache::disabled() + parser::manifest_cache::ManifestAnalysisCache::disabled() } else if refresh_cache { - parser::manifest_cache::ManifestGraphCache::fresh(project_dir, cache_dir) + parser::manifest_cache::ManifestAnalysisCache::fresh(project_dir, cache_dir) } else { - parser::manifest_cache::ManifestGraphCache::load(project_dir, cache_dir) + parser::manifest_cache::ManifestAnalysisCache::load(project_dir, cache_dir) }; - // Parse exactly once so diagnostics are available on cache hits. - // Graph construction is deferred until after the cache lookup. - let (load_report, parsed_fingerprint) = - load_manifest_report_with_fingerprint(&resolved)?; - let (manifest, diagnostics) = load_report.into_parts(); - let manifest = manifest.ok_or_else(|| { - let message = diagnostics - .iter() - .find(|diagnostic| { - diagnostic.severity == parser::manifest::ManifestDiagnosticSeverity::Error - }) - .map(|diagnostic| diagnostic.message.clone()) - .unwrap_or_else(|| "manifest could not be loaded".to_string()); - anyhow::anyhow!(message) - })?; - let cached_graph = parsed_fingerprint.as_ref().and_then(|fingerprint| { - cache.get_with_fingerprint( - &resolved, - ( - fingerprint.mtime_secs, - fingerprint.mtime_nanos, - fingerprint.size_bytes, - fingerprint.content_hash, - ), - ) - }); - if let Some(cached_graph) = cached_graph { - let graph = cached_graph.clone(); - return Ok((graph, None, diagnostics, Some(manifest))); + // Read and hash bytes before parsing. A cache hit restores the + // compact model-level analysis without deserializing Manifest. + let manifest_bytes = load_manifest_bytes(&resolved)?; + if let Some(analysis) = cache.take_for_manifest(&manifest_bytes) { + let (graph, diagnostics, project_name, referenced_file_paths) = + analysis.into_parts(); + let context = ManifestDagContext { + manifest: None, + project_name, + referenced_file_paths, + }; + return Ok((graph, None, diagnostics, Some(context))); } - let graph_report = parser::manifest::build_graph_from_load_report( - parser::manifest::ManifestLoadReport::from_parts(Some(manifest), diagnostics), - )?; + + let load_report = + parser::manifest::load_manifest_report_from_bytes(&manifest_bytes, &resolved); + let graph_report = parser::manifest::build_graph_from_load_report(load_report)?; let parser::manifest::ManifestGraphReport { graph, diagnostics, manifest, .. } = graph_report; - if let Some(parsed_fingerprint) = parsed_fingerprint { - let expected = ( - parsed_fingerprint.mtime_secs, - parsed_fingerprint.mtime_nanos, - parsed_fingerprint.size_bytes, - parsed_fingerprint.content_hash, - ); - if cache.insert_if_fingerprint_matches(&resolved, &graph, expected) { - cache.save(); - } else { - dlin_core::warn!( - "manifest changed after read/parse; skipping manifest graph cache write for {}", - resolved.display() - ); - } - } else { - dlin_core::warn!( - "could not compute manifest fingerprint; skipping manifest graph cache write for {}", - resolved.display() + let project_name = manifest.metadata.project_name.clone(); + let referenced_file_paths: Vec = + manifest.collect_file_paths().into_iter().collect(); + if !no_cache { + let analysis = parser::manifest_cache::ManifestAnalysis::new( + graph.clone(), + diagnostics.clone(), + project_name.clone(), + referenced_file_paths.clone(), ); + cache.insert_for_manifest(&manifest_bytes, analysis); + cache.save(); } - Ok((graph, None, diagnostics, Some(manifest))) + Ok(( + graph, + None, + diagnostics, + Some(ManifestDagContext { + project_name, + referenced_file_paths, + manifest: Some(manifest), + }), + )) } SourceType::Sql => { let project = parser::project::DbtProject::load(project_dir)?; @@ -137,62 +128,14 @@ pub(crate) fn warn_manifest_diagnostics(diagnostics: &[parser::manifest::Manifes } } -#[derive(Debug, Clone, PartialEq, Eq)] -struct ManifestFingerprint { - mtime_secs: u64, - mtime_nanos: u32, - size_bytes: u64, - content_hash: u64, -} - -fn manifest_fingerprint(meta: &std::fs::Metadata, bytes: &[u8]) -> Option { - let modified = meta - .modified() - .ok()? - .duration_since(SystemTime::UNIX_EPOCH) - .ok()?; - Some(ManifestFingerprint { - mtime_secs: modified.as_secs(), - mtime_nanos: modified.subsec_nanos(), - size_bytes: meta.len(), - content_hash: hash_bytes(bytes), - }) -} - -fn load_manifest_report_with_fingerprint( - manifest_path: &Path, -) -> Result<( - parser::manifest::ManifestLoadReport, - Option, -)> { - let meta = std::fs::metadata(manifest_path).map_err(|e| { - dlin_core::error::DbtLineageError::FileReadError { - path: manifest_path.to_path_buf(), - source: e, - } - })?; +fn load_manifest_bytes(manifest_path: &Path) -> Result> { let bytes = std::fs::read(manifest_path).map_err(|e| { dlin_core::error::DbtLineageError::FileReadError { path: manifest_path.to_path_buf(), source: e, } })?; - let load_report = parser::manifest::load_manifest_report_from_bytes(&bytes, manifest_path); - let fingerprint = manifest_fingerprint(&meta, &bytes); - Ok((load_report, fingerprint)) -} - -fn hash_bytes(bytes: &[u8]) -> u64 { - // FNV-1a 64-bit constants (offset basis + prime), used for lightweight - // non-cryptographic fingerprinting of manifest content. - const FNV_OFFSET_BASIS: u64 = 0xcbf29ce484222325; - const FNV_PRIME: u64 = 0x100000001b3; - let mut hash = FNV_OFFSET_BASIS; - for &b in bytes { - hash ^= b as u64; - hash = hash.wrapping_mul(FNV_PRIME); - } - hash + Ok(bytes) } pub(crate) fn warn_sql_mode_test_limitation(source: &SourceType, has_tests: bool) { diff --git a/crates/dlin/tests/integration/manifest_only.rs b/crates/dlin/tests/integration/manifest_only.rs index d8e2069..0fe76c0 100644 --- a/crates/dlin/tests/integration/manifest_only.rs +++ b/crates/dlin/tests/integration/manifest_only.rs @@ -343,6 +343,125 @@ fn test_summary_manifest_mode_json_output() { ); } +#[test] +fn test_manifest_cache_warm_summary_restores_metadata_and_deleted_paths() { + let tmp = minimal_manifest_dir(Some("cached_project")); + fs::write( + tmp.path().join("dbt_project.yml"), + "name: cached_project\nmodel-paths: [models]\n", + ) + .unwrap(); + let manifest_path = tmp.path().join("target/manifest.json"); + let cache_dir = tmp.path().join("cache"); + let args = [ + "summary", + "--source", + "manifest", + "--manifest-path", + manifest_path.to_str().unwrap(), + "--project-dir", + tmp.path().to_str().unwrap(), + "--cache-dir", + cache_dir.to_str().unwrap(), + "-o", + "json", + ]; + + let cold = Command::new(binary_path()).args(args).output().unwrap(); + assert!(cold.status.success(), "cold summary failed: {:?}", cold); + let warm = Command::new(binary_path()).args(args).output().unwrap(); + assert!(warm.status.success(), "warm summary failed: {:?}", warm); + assert_eq!(cold.stdout, warm.stdout); + let summary: serde_json::Value = serde_json::from_slice(&warm.stdout).unwrap(); + assert_eq!(summary["project_name"], "cached_project"); + assert_eq!(summary["manifest_status"]["deleted_file_count"], 2); + assert_eq!( + summary["manifest_status"]["deleted_files"] + .as_array() + .unwrap() + .len(), + 2 + ); +} + +#[test] +fn test_manifest_cache_warm_sql_content_loads_typed_manifest_lazily() { + let tmp = minimal_manifest_dir(None); + let manifest_path = tmp.path().join("target/manifest.json"); + let content = fs::read_to_string(&manifest_path).unwrap(); + let content = content.replace( + "\"description\": \"Staged orders\"", + "\"description\": \"Staged orders\", \"compiled_code\": \"select 1 as order_id\"", + ); + fs::write(&manifest_path, content).unwrap(); + let cache_dir = tmp.path().join("cache"); + let args = [ + "graph", + "--source", + "manifest", + "--manifest-path", + manifest_path.to_str().unwrap(), + "--project-dir", + tmp.path().to_str().unwrap(), + "--cache-dir", + cache_dir.to_str().unwrap(), + "--json-fields", + "sql_content", + "-o", + "json", + ]; + let cold = Command::new(binary_path()).args(args).output().unwrap(); + assert!(cold.status.success(), "cold graph failed: {:?}", cold); + let warm = Command::new(binary_path()).args(args).output().unwrap(); + assert!(warm.status.success(), "warm graph failed: {:?}", warm); + assert_eq!(cold.stdout, warm.stdout); + let graph: serde_json::Value = serde_json::from_slice(&warm.stdout).unwrap(); + assert!( + graph["nodes"] + .as_array() + .unwrap() + .iter() + .any(|node| node["sql_content"] == "select 1 as order_id"), + "unexpected graph JSON: {}", + String::from_utf8_lossy(&warm.stdout) + ); + + let list_args = [ + "list", + "--source", + "manifest", + "--manifest-path", + manifest_path.to_str().unwrap(), + "--project-dir", + tmp.path().to_str().unwrap(), + "--cache-dir", + cache_dir.to_str().unwrap(), + "--json-fields", + "sql_content", + "-o", + "json", + ]; + let list_cold = Command::new(binary_path()) + .args(list_args) + .output() + .unwrap(); + assert!( + list_cold.status.success(), + "cold list failed: {:?}", + list_cold + ); + let list_warm = Command::new(binary_path()) + .args(list_args) + .output() + .unwrap(); + assert!( + list_warm.status.success(), + "warm list failed: {:?}", + list_warm + ); + assert_eq!(list_cold.stdout, list_warm.stdout); +} + #[test] fn test_summary_manifest_mode_with_malformed_project_yml() { let tmp = minimal_manifest_dir(Some("malformed_test")); From 91a1a0eddaa20e35df5ff653cb10956dc501d32a Mon Sep 17 00:00:00 2001 From: eitsupi Date: Sat, 29 Aug 2026 14:52:46 +0000 Subject: [PATCH 2/2] fix(manifest): reuse input bytes for SQL content --- crates/dlin/src/commands/graph.rs | 75 ++++++++++++++++++++++++++++++ crates/dlin/src/commands/shared.rs | 3 ++ 2 files changed, 78 insertions(+) diff --git a/crates/dlin/src/commands/graph.rs b/crates/dlin/src/commands/graph.rs index 04f3694..84bce6f 100644 --- a/crates/dlin/src/commands/graph.rs +++ b/crates/dlin/src/commands/graph.rs @@ -131,6 +131,9 @@ pub(crate) fn run_graph_command(args: GraphArgs) -> Result<()> { manifest_context .as_ref() .and_then(|context| context.manifest.as_ref()), + manifest_context + .as_ref() + .and_then(|context| context.manifest_bytes.as_deref()), args.manifest_path.as_ref(), &filtered, )) @@ -285,6 +288,9 @@ pub(crate) fn run_list_command(args: ListArgs) -> Result<()> { manifest_context .as_ref() .and_then(|context| context.manifest.as_ref()), + manifest_context + .as_ref() + .and_then(|context| context.manifest_bytes.as_deref()), args.manifest_path.as_ref(), &filtered, )) @@ -330,6 +336,7 @@ fn collect_sql_contents_for_source( source: &SourceType, project_dir: &Path, manifest: Option<&parser::manifest::Manifest>, + manifest_bytes: Option<&[u8]>, manifest_path: Option<&PathBuf>, graph: &graph::types::LineageGraph, ) -> HashMap { @@ -337,6 +344,18 @@ fn collect_sql_contents_for_source( SourceType::Manifest => { if let Some(manifest) = manifest { manifest.collect_sql_contents() + } else if let Some(bytes) = manifest_bytes { + // Use the already-read bytes even when the path has changed + // or disappeared since DAG construction. The path is only + // needed for diagnostics from the byte parser. + let resolved = manifest_path + .cloned() + .unwrap_or_else(|| project_dir.join("target").join("manifest.json")); + let Ok(manifest) = parser::manifest::load_manifest_from_bytes(bytes, &resolved) + else { + return HashMap::new(); + }; + manifest.collect_sql_contents() } else { let Ok(resolved) = resolve_manifest_path_or_default(manifest_path, project_dir) else { @@ -452,3 +471,59 @@ pub(crate) fn run_impact_command( Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn manifest_sql_content_uses_initial_bytes_when_path_content_differs() { + let temp = tempfile::tempdir().unwrap(); + let manifest_path = temp.path().join("deleted-manifest.json"); + let bytes = serde_json::to_vec(&serde_json::json!({ + "metadata": {}, + "nodes": { + "model.project.orders": { + "unique_id": "model.project.orders", + "name": "orders", + "resource_type": "model", + "depends_on": {"nodes": []}, + "config": {}, + "description": null, + "path": null, + "original_file_path": null, + "columns": {}, + "compiled_code": "select 1", + "database": null, + "schema": null + } + } + })) + .unwrap(); + let conflicting_bytes = String::from_utf8(bytes.clone()) + .unwrap() + .replace("select 1", "select 2"); + std::fs::write(&manifest_path, conflicting_bytes).unwrap(); + let contents = collect_sql_contents_for_source( + &SourceType::Manifest, + temp.path(), + None, + Some(&bytes), + Some(&manifest_path), + &graph::types::LineageGraph::new(), + ); + + assert_eq!( + parser::manifest::load_manifest_from_bytes(&bytes, &manifest_path) + .unwrap() + .collect_sql_contents() + .get("model.project.orders") + .map(String::as_str), + Some("select 1") + ); + assert_eq!( + contents.get("model.project.orders").map(String::as_str), + Some("select 1") + ); + } +} diff --git a/crates/dlin/src/commands/shared.rs b/crates/dlin/src/commands/shared.rs index c92daf9..710842c 100644 --- a/crates/dlin/src/commands/shared.rs +++ b/crates/dlin/src/commands/shared.rs @@ -10,6 +10,7 @@ use dlin_core::parser; #[cfg(not(tarpaulin_include))] pub(crate) struct ManifestDagContext { pub(crate) manifest: Option, + pub(crate) manifest_bytes: Option>, pub(crate) project_name: Option, pub(crate) referenced_file_paths: Vec, } @@ -51,6 +52,7 @@ pub(crate) fn build_dag( analysis.into_parts(); let context = ManifestDagContext { manifest: None, + manifest_bytes: Some(manifest_bytes), project_name, referenced_file_paths, }; @@ -88,6 +90,7 @@ pub(crate) fn build_dag( project_name, referenced_file_paths, manifest: Some(manifest), + manifest_bytes: None, }), )) }