diff --git a/cli/src/cache/mod.rs b/cli/src/cache/mod.rs index 4397f1a..d418ae7 100644 --- a/cli/src/cache/mod.rs +++ b/cli/src/cache/mod.rs @@ -21,8 +21,10 @@ pub use fingerprint::{ }; pub use location::{CacheLocation, ProjectKey}; pub use schema::SCHEMA_VERSION; -pub(crate) use store::CacheLoadFailure; pub use store::{CacheGraphRead, CacheStore, SnapshotSummary}; +pub(crate) use store::{ + CacheLoadFailure, CachePublicationFailure, PublicationConflict, escaped_publication_identifier, +}; #[cfg(test)] pub(crate) use store::{reset_whole_graph_loads, whole_graph_loads}; #[cfg(test)] diff --git a/cli/src/cache/store.rs b/cli/src/cache/store.rs index cfc8f6e..b506324 100644 --- a/cli/src/cache/store.rs +++ b/cli/src/cache/store.rs @@ -30,6 +30,14 @@ use super::{ restore_subgraph, }; +mod candidate_publication; +mod graph_publication; +mod publication_error; + +pub(crate) use publication_error::{ + CachePublicationFailure, PublicationConflict, PublicationField, escaped_publication_identifier, +}; + const LOCK_WAIT_CAP: Duration = Duration::from_secs(2); #[cfg(test)] @@ -68,18 +76,6 @@ pub struct CacheGraphRead<'store, 'deadline> { deadline: &'deadline Deadline, } -#[derive(Debug)] -struct CandidateFileRow { - language: String, - content_hash: Vec, - size_bytes: i64, - mtime_seconds: Option, - mtime_nanoseconds: Option, - package_assignment: String, - file_facts: Vec, - file_subgraph: Option>, -} - #[derive(Debug)] struct LoadedCandidateFileRow { path: String, @@ -349,8 +345,17 @@ impl CacheStore { candidate: &CandidateSnapshot, deadline: &Deadline, ) -> Result<(), CacheError> { + self.publish_candidate_detailed(candidate, deadline) + .map_err(Into::into) + } + + pub(crate) fn publish_candidate_detailed( + &self, + candidate: &CandidateSnapshot, + deadline: &Deadline, + ) -> Result<(), CachePublicationFailure> { if !self.writable { - return Err(CacheError::ReadOnly); + return Err(CacheError::ReadOnly.into()); } ensure_time(deadline)?; let encoded = PreparedCandidate::new(candidate, deadline)?; @@ -376,10 +381,24 @@ impl CacheStore { // Publication timestamps are store-owned. Fingerprint components, // unlike timestamps, are exact compatibility content and conflicts // must be rejected even when the derived compatibility id matches. - if stored_compatibility.0 != encoded.language_fingerprint - || stored_compatibility.1 != encoded.package_fingerprint - { - return Err(CacheError::CandidateConflict); + for (differs, field) in [ + ( + stored_compatibility.0 != encoded.language_fingerprint, + PublicationField::LanguageFingerprint, + ), + ( + stored_compatibility.1 != encoded.package_fingerprint, + PublicationField::PackageFingerprint, + ), + ] { + if differs { + return Err(CachePublicationFailure::conflict( + encoded.candidate_id, + field, + None, + None, + )); + } } let existing: Option = self.connection.query_row( "SELECT compatibility_id, input_digest, completeness, inventory_file_count, inventory_total_bytes FROM candidates WHERE candidate_id = ?1", @@ -395,13 +414,36 @@ impl CacheStore { }, ).optional().map_err(|error| map_sqlite_error(error, deadline))?; if let Some(existing) = existing { - if existing.compatibility_id != encoded.compatibility_id - || existing.input_digest != encoded.input_digest - || existing.completeness != encoded.completeness - || existing.inventory_file_count != encoded.inventory_file_count - || existing.inventory_total_bytes != encoded.inventory_total_bytes - { - return Err(CacheError::CandidateConflict); + for (differs, field) in [ + ( + existing.compatibility_id != encoded.compatibility_id, + PublicationField::CompatibilityId, + ), + ( + existing.input_digest != encoded.input_digest, + PublicationField::InputDigest, + ), + ( + existing.completeness != encoded.completeness, + PublicationField::Completeness, + ), + ( + existing.inventory_file_count != encoded.inventory_file_count, + PublicationField::InventoryFileCount, + ), + ( + existing.inventory_total_bytes != encoded.inventory_total_bytes, + PublicationField::InventoryTotalBytes, + ), + ] { + if differs { + return Err(CachePublicationFailure::conflict( + encoded.candidate_id, + field, + None, + None, + )); + } } self.verify_existing_candidate(&encoded, deadline)?; } else { @@ -456,7 +498,7 @@ impl CacheStore { params![encoded.candidate_id.as_slice(), graph.tier], |row| row.get(0), ).optional().map_err(|error| map_sqlite_error(error, deadline))?; let snapshot_id = if let Some(snapshot_id) = snapshot_id { - self.verify_existing_graph(snapshot_id, graph, deadline)?; + self.verify_existing_graph(encoded.candidate_id, snapshot_id, graph, deadline)?; snapshot_id } else { self.connection.execute( @@ -567,7 +609,7 @@ impl CacheStore { [], ) .map_err(|error| map_sqlite_error(error, deadline))?; - ensure_time(deadline) + ensure_time(deadline).map_err(CachePublicationFailure::from) })(); match result { Ok(()) => match self.connection.execute_batch("COMMIT") { @@ -580,7 +622,7 @@ impl CacheStore { Err(error) => { let mapped = map_sqlite_error(error, deadline); let _ = self.connection.execute_batch("ROLLBACK"); - Err(mapped) + Err(mapped.into()) } }, Err(error) => { @@ -610,82 +652,6 @@ impl CacheStore { })(); } - fn verify_existing_graph( - &self, - snapshot_id: i64, - graph: &PreparedGraph, - deadline: &Deadline, - ) -> Result<(), CacheError> { - let stored_symbols = - self.load_graph_payloads(snapshot_id, "graph_symbols", "symbol", deadline)?; - // Edges have no serialized copy to compare, so compare the identity the - // columns carry: `edge_key` is the lossless edge identity and - // `confidence` is the one attribute it deliberately excludes. - let stored_edges = - self.load_graph_payloads(snapshot_id, "graph_edges", "edge_key", deadline)?; - let stored_confidence = - self.load_graph_text(snapshot_id, "graph_edges", "confidence", deadline)?; - if stored_symbols.len() != graph.symbols.len() - || stored_edges.len() != graph.edges.len() - || stored_confidence.len() != graph.edges.len() - || stored_symbols - .iter() - .zip(&graph.symbols) - .any(|(stored, row)| *stored != row.payload) - || stored_edges - .iter() - .zip(&graph.edges) - .any(|(stored, row)| *stored != row.edge_key) - || stored_confidence - .iter() - .zip(&graph.edges) - .any(|(stored, row)| *stored != row.confidence) - { - return Err(CacheError::CandidateConflict); - } - Ok(()) - } - - fn load_graph_payloads( - &self, - snapshot_id: i64, - table: &str, - column: &str, - deadline: &Deadline, - ) -> Result>, CacheError> { - let sql = - format!("SELECT {column} FROM {table} WHERE snapshot_id = ?1 ORDER BY ordinal ASC"); - let mut statement = self - .connection - .prepare(&sql) - .map_err(|error| map_sqlite_error(error, deadline))?; - statement - .query_map([snapshot_id], |row| row.get::<_, Vec>(0)) - .map_err(|error| map_sqlite_error(error, deadline))? - .collect::, _>>() - .map_err(|error| map_sqlite_error(error, deadline)) - } - - fn load_graph_text( - &self, - snapshot_id: i64, - table: &str, - column: &str, - deadline: &Deadline, - ) -> Result, CacheError> { - let sql = - format!("SELECT {column} FROM {table} WHERE snapshot_id = ?1 ORDER BY ordinal ASC"); - let mut statement = self - .connection - .prepare(&sql) - .map_err(|error| map_sqlite_error(error, deadline))?; - statement - .query_map([snapshot_id], |row| row.get::<_, String>(0)) - .map_err(|error| map_sqlite_error(error, deadline))? - .collect::, _>>() - .map_err(|error| map_sqlite_error(error, deadline)) - } - fn load_graph_rows( &self, snapshot_id: i64, @@ -1169,95 +1135,6 @@ impl CacheStore { .map(Some) } - fn verify_existing_candidate( - &self, - candidate: &PreparedCandidate, - deadline: &Deadline, - ) -> Result<(), CacheError> { - let count: i64 = self - .connection - .query_row( - "SELECT count(*) FROM candidate_files WHERE candidate_id = ?1", - [candidate.candidate_id.as_slice()], - |row| row.get(0), - ) - .map_err(|error| map_sqlite_error(error, deadline))?; - if count - != i64::try_from(candidate.files.len()).map_err(|_| CacheError::InvalidCandidate)? - { - return Err(CacheError::CandidateConflict); - } - let omissions = { - let mut statement = self - .connection - .prepare("SELECT path, reason, detail FROM candidate_omissions WHERE candidate_id = ?1 ORDER BY path ASC, reason ASC, detail ASC") - .map_err(|error| map_sqlite_error(error, deadline))?; - statement - .query_map([candidate.candidate_id.as_slice()], |row| { - Ok(super::CacheOmission { - path: row.get(0)?, - reason: row.get(1)?, - detail: row.get(2)?, - }) - }) - .map_err(|error| map_sqlite_error(error, deadline))? - .collect::, _>>() - .map_err(|error| map_sqlite_error(error, deadline))? - }; - if omissions != candidate.omissions { - return Err(CacheError::CandidateConflict); - } - for file in &candidate.files { - ensure_time(deadline)?; - let found: Option = self - .connection - .query_row( - "SELECT language, content_hash, size_bytes, mtime_seconds, mtime_nanoseconds, package_assignment, file_facts, file_subgraph FROM candidate_files WHERE candidate_id = ?1 AND path = ?2", - params![candidate.candidate_id.as_slice(), file.path], - |row| { - Ok(CandidateFileRow { - language: row.get(0)?, - content_hash: row.get(1)?, - size_bytes: row.get(2)?, - mtime_seconds: row.get(3)?, - mtime_nanoseconds: row.get(4)?, - package_assignment: row.get(5)?, - file_facts: row.get(6)?, - file_subgraph: row.get(7)?, - }) - }, - ) - .optional() - .map_err(|error| map_sqlite_error(error, deadline))?; - let Some(found) = found else { - return Err(CacheError::CandidateConflict); - }; - if found.language != file.language - || found.content_hash != file.content_hash - || found.size_bytes != file.size_bytes - || found.mtime_seconds != file.mtime_seconds - || found.mtime_nanoseconds != file.mtime_nanoseconds - || found.package_assignment != file.package_assignment - || found.file_facts != file.facts - { - return Err(CacheError::CandidateConflict); - } - match (found.file_subgraph, &file.subgraph) { - (None, Some(subgraph)) => { - self.connection.execute( - "UPDATE candidate_files SET file_subgraph = ?1 WHERE candidate_id = ?2 AND path = ?3 AND file_subgraph IS NULL", - params![subgraph, candidate.candidate_id.as_slice(), file.path], - ).map_err(|error| map_sqlite_error(error, deadline))?; - } - (Some(stored), Some(incoming)) if stored != *incoming => { - return Err(CacheError::CandidateConflict); - } - _ => {} - } - } - Ok(()) - } - fn load_candidate_inner( &self, candidate_id: CandidateId, @@ -2582,6 +2459,12 @@ impl GraphRead for CacheGraphRead<'_, '_> { #[cfg(test)] mod tests { + mod candidate_conflicts; + mod metadata_refresh; + mod publication_diagnostics; + mod publication_slots; + mod publication_transactions; + use super::*; use code2graph::{ Confidence, Descriptor, Occurrence, Provenance, RefRole, SymbolKind, Visibility, @@ -3009,85 +2892,6 @@ mod tests { ); } - #[test] - fn candidate_publication_keeps_complete_and_partial_slots_independent() { - let temp = tempdir().expect("tempdir"); - let root = temp.path().join("project"); - fs::create_dir(&root).expect("project"); - let cache_location = location(&root, temp.path()); - let store = - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); - let complete = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); - let partial = candidate(CacheCompleteness::Partial, ResolverCacheTier::Name); - store - .publish_candidate(&complete, &Deadline::new(None)) - .expect("publish complete"); - store - .publish_candidate(&partial, &Deadline::new(None)) - .expect("publish partial"); - assert_eq!( - store - .load_active( - ResolverCacheTier::Name, - CacheCompleteness::Complete, - complete.compatibility.id, - &Deadline::new(None) - ) - .expect("load") - .expect("active") - .candidate_id, - complete.candidate_id - ); - assert_eq!( - store - .load_active( - ResolverCacheTier::Name, - CacheCompleteness::Partial, - partial.compatibility.id, - &Deadline::new(None) - ) - .expect("load") - .expect("active") - .candidate_id, - partial.candidate_id - ); - let incompatible = CompatibilityFingerprint::new( - super::super::LanguageFeatureFingerprint::current(), - super::super::PackageFingerprint::from_normalized(["different-package"]), - ); - let loaded_complete = store - .load_active( - ResolverCacheTier::Name, - CacheCompleteness::Complete, - complete.compatibility.id, - &Deadline::new(None), - ) - .expect("load") - .expect("active"); - assert_eq!( - loaded_complete.compatibility.language_fingerprint, - complete.compatibility.language_fingerprint - ); - assert_eq!( - loaded_complete.compatibility.package_fingerprint, - complete.compatibility.package_fingerprint - ); - assert!( - store - .load_active( - ResolverCacheTier::Name, - CacheCompleteness::Complete, - incompatible, - &Deadline::new(None), - ) - .expect("compatibility miss") - .is_none() - ); - store - .publish_candidate(&complete, &Deadline::new(None)) - .expect("idempotent publish"); - } - #[test] fn fresh_writable_cache_enables_incremental_auto_vacuum() { // Locks the auto_vacuum mode set in `initialize_or_join_v1`: it must be @@ -3110,76 +2914,6 @@ mod tests { ); } - #[test] - fn superseding_a_slot_garbage_collects_the_prior_snapshot_and_candidate() { - let temp = tempdir().expect("tempdir"); - let root = temp.path().join("project"); - fs::create_dir(&root).expect("project"); - let cache_location = location(&root, temp.path()); - let store = - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); - // Two distinct candidates (different input digests) target the same - // (tier, completeness) slot; publishing B flips active away from A. - let a = candidate_with_hash( - CacheCompleteness::Complete, - ResolverCacheTier::Name, - [3; 32], - ); - let b = candidate_with_hash( - CacheCompleteness::Complete, - ResolverCacheTier::Name, - [7; 32], - ); - assert_ne!(a.candidate_id, b.candidate_id); - store - .publish_candidate(&a, &Deadline::new(None)) - .expect("publish a"); - store - .publish_candidate(&b, &Deadline::new(None)) - .expect("publish b"); - - // Only B's snapshot survives; A's snapshot and candidate rows are gone. - let snapshot_count: i64 = store - .connection - .query_row("SELECT count(*) FROM graph_snapshots", [], |row| row.get(0)) - .expect("snapshot count"); - assert_eq!(snapshot_count, 1); - let surviving_candidate: Vec = store - .connection - .query_row("SELECT candidate_id FROM graph_snapshots", [], |row| { - row.get(0) - }) - .expect("surviving candidate"); - assert_eq!( - surviving_candidate.as_slice(), - b.candidate_id.as_bytes().as_slice() - ); - let a_candidate_count: i64 = store - .connection - .query_row( - "SELECT count(*) FROM candidates WHERE candidate_id = ?1", - [a.candidate_id.as_bytes().as_slice()], - |row| row.get(0), - ) - .expect("a candidate count"); - assert_eq!(a_candidate_count, 0); - - // B remains the queryable active snapshot for the slot. - assert_eq!( - store - .load_active( - ResolverCacheTier::Name, - CacheCompleteness::Complete, - b.compatibility.id, - &Deadline::new(None), - ) - .expect("load") - .expect("active") - .candidate_id, - b.candidate_id - ); - } - #[test] fn latest_active_loads_full_snapshot_without_compatibility_or_mutation() { let temp = tempdir().expect("tempdir"); @@ -3329,109 +3063,6 @@ mod tests { )); } - #[test] - fn rejects_inconsistent_candidates_and_conflicting_republication() { - let temp = tempdir().expect("tempdir"); - let root = temp.path().join("project"); - fs::create_dir(&root).expect("project"); - let cache_location = location(&root, temp.path()); - let store = - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); - - let mut unsorted = candidate(CacheCompleteness::Partial, ResolverCacheTier::Name); - unsorted.omissions = vec![ - super::super::CacheOmission { - path: "z".into(), - reason: "x".into(), - detail: "detail".into(), - }, - super::super::CacheOmission { - path: "a".into(), - reason: "x".into(), - detail: "detail".into(), - }, - ]; - unsorted.candidate_id = CandidateId::new( - unsorted.compatibility.id, - unsorted.input_digest, - unsorted.completeness, - &unsorted.omissions, - ); - assert!(matches!( - store.publish_candidate(&unsorted, &Deadline::new(None)), - Err(CacheError::InvalidCandidate) - )); - - let mut overflow = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); - overflow.created_at_ns = u64::MAX; - assert!(matches!( - store.publish_candidate(&overflow, &Deadline::new(None)), - Err(CacheError::InvalidCandidate) - )); - - let original = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); - store - .publish_candidate(&original, &Deadline::new(None)) - .expect("publish"); - let mut republished = original.clone(); - republished.created_at_ns += 1; - republished.compatibility.created_at_ns += 1; - store - .publish_candidate(&republished, &Deadline::new(None)) - .expect("timestamps are store-owned and do not conflict"); - assert_eq!( - store - .load_active( - ResolverCacheTier::Name, - CacheCompleteness::Complete, - original.compatibility.id, - &Deadline::new(None), - ) - .expect("load") - .expect("active") - .created_at_ns, - original.created_at_ns - ); - } - - #[test] - fn scope_publication_requires_and_restores_every_owned_subgraph() { - let temp = tempdir().expect("tempdir"); - let root = temp.path().join("project"); - fs::create_dir(&root).expect("project"); - let cache_location = location(&root, temp.path()); - let store = - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); - let mut snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Scope); - assert!(matches!( - store.publish_candidate(&snapshot, &Deadline::new(None)), - Err(CacheError::InvalidCandidate) - )); - // A Name snapshot may be published first; a later Scope publication - // for the identical candidate augments its per-file subgraphs. - let mut name = snapshot.clone(); - name.tier_graphs = vec![( - ResolverCacheTier::Name, - CodeGraph { - symbols: Vec::new(), - edges: Vec::new(), - }, - )]; - store - .publish_candidate(&name, &Deadline::new(None)) - .expect("publish name"); - let mut incremental = IncrementalGraph::new(); - incremental.upsert(&snapshot.files[0].facts); - snapshot.files[0].subgraph = incremental.subgraph("src/a.rs").cloned(); - store - .publish_candidate(&snapshot, &Deadline::new(None)) - .expect("augment with scope"); - let restored = store - .hydrate_scope_subgraphs(snapshot.candidate_id, &Deadline::new(None)) - .expect("hydrate"); - assert!(restored.subgraph("src/a.rs").is_some()); - } - #[test] fn missing_normalized_graph_snapshot_is_typed() { let temp = tempdir().expect("tempdir"); @@ -3469,101 +3100,6 @@ mod tests { )); } - #[test] - fn failed_graph_write_rolls_back_candidate_and_active_publication() { - let temp = tempdir().expect("tempdir"); - let root = temp.path().join("project"); - fs::create_dir(&root).expect("project"); - let cache_location = location(&root, temp.path()); - let store = - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); - let candidate = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); - store.connection.execute_batch( - "CREATE TEMP TRIGGER fail_graph BEFORE INSERT ON graph_snapshots BEGIN SELECT RAISE(ABORT, 'injected graph failure'); END", - ).expect("failure trigger"); - assert!(matches!( - store.publish_candidate(&candidate, &Deadline::new(None)), - Err(CacheError::Access) - )); - let candidate_count: i64 = store - .connection - .query_row( - "SELECT count(*) FROM candidates WHERE candidate_id = ?1", - [candidate.candidate_id.as_bytes().as_slice()], - |row| row.get(0), - ) - .expect("candidate count"); - let active_count: i64 = store - .connection - .query_row("SELECT count(*) FROM active_snapshots", [], |row| { - row.get(0) - }) - .expect("active count"); - assert_eq!((candidate_count, active_count), (0, 0)); - store - .connection - .execute_batch("DROP TRIGGER fail_graph") - .expect("drop trigger"); - store - .publish_candidate(&candidate, &Deadline::new(None)) - .expect("retry"); - } - - #[test] - fn concurrent_publishers_commit_whole_candidates() { - use std::sync::{Arc, Barrier}; - - let temp = tempdir().expect("tempdir"); - let root = temp.path().join("project"); - fs::create_dir(&root).expect("project"); - let cache_location = location(&root, temp.path()); - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)) - .expect("initialize"); - let barrier = Arc::new(Barrier::new(2)); - let handles: Vec<_> = [CacheCompleteness::Complete, CacheCompleteness::Partial] - .into_iter() - .map(|completeness| { - let barrier = Arc::clone(&barrier); - let root = root.clone(); - let cache_location = cache_location.clone(); - std::thread::spawn(move || { - let store = - CacheStore::open_writable(&cache_location, &root, &Deadline::new(None))?; - let candidate = candidate(completeness, ResolverCacheTier::Name); - barrier.wait(); - store.publish_candidate(&candidate, &Deadline::new(None))?; - Ok::<_, CacheError>(candidate.candidate_id) - }) - }) - .collect(); - let ids: Vec<_> = handles - .into_iter() - .map(|handle| handle.join().expect("publisher thread").expect("publish")) - .collect(); - let store = - CacheStore::open_frozen(&cache_location, &root, &Deadline::new(None)).expect("frozen"); - for (completeness, id) in [CacheCompleteness::Complete, CacheCompleteness::Partial] - .into_iter() - .zip(ids) - { - assert_eq!( - store - .load_active( - ResolverCacheTier::Name, - completeness, - candidate(completeness, ResolverCacheTier::Name) - .compatibility - .id, - &Deadline::new(None), - ) - .expect("load") - .expect("active") - .candidate_id, - id - ); - } - } - #[test] fn frozen_missing_cache_creates_nothing() { let temp = tempdir().expect("tempdir"); diff --git a/cli/src/cache/store/candidate_publication.rs b/cli/src/cache/store/candidate_publication.rs new file mode 100644 index 0000000..a5485b2 --- /dev/null +++ b/cli/src/cache/store/candidate_publication.rs @@ -0,0 +1,146 @@ +// SPDX-License-Identifier: Apache-2.0 + +use rusqlite::{OptionalExtension, params}; + +use super::{ + CacheError, CachePublicationFailure, CacheStore, CandidateId, Deadline, PreparedCandidate, + PublicationField, ensure_time, map_sqlite_error, +}; + +#[derive(Debug)] +struct CandidateFileRow { + language: String, + content_hash: Vec, + size_bytes: i64, + mtime_seconds: Option, + mtime_nanoseconds: Option, + package_assignment: String, + file_facts: Vec, + file_subgraph: Option>, +} + +impl CacheStore { + pub(super) fn verify_existing_candidate( + &self, + candidate: &PreparedCandidate, + deadline: &Deadline, + ) -> Result<(), CachePublicationFailure> { + let count: i64 = self + .connection + .query_row( + "SELECT count(*) FROM candidate_files WHERE candidate_id = ?1", + [candidate.candidate_id.as_slice()], + |row| row.get(0), + ) + .map_err(|error| map_sqlite_error(error, deadline))?; + if count + != i64::try_from(candidate.files.len()).map_err(|_| CacheError::InvalidCandidate)? + { + return Err(CachePublicationFailure::conflict( + candidate.candidate_id, + PublicationField::FileCount, + None, + None, + )); + } + let omissions = + self.load_omissions(CandidateId::from_bytes(candidate.candidate_id), deadline)?; + if omissions != candidate.omissions { + return Err(CachePublicationFailure::conflict( + candidate.candidate_id, + PublicationField::Omissions, + None, + None, + )); + } + let mut update_mtime = self.connection.prepare( + "UPDATE candidate_files SET mtime_seconds = ?1, mtime_nanoseconds = ?2 WHERE candidate_id = ?3 AND path = ?4", + ).map_err(|error| map_sqlite_error(error, deadline))?; + for file in &candidate.files { + ensure_time(deadline)?; + let found: Option = self + .connection + .query_row( + "SELECT language, content_hash, size_bytes, mtime_seconds, mtime_nanoseconds, package_assignment, file_facts, file_subgraph FROM candidate_files WHERE candidate_id = ?1 AND path = ?2", + params![candidate.candidate_id.as_slice(), file.path], + |row| { + Ok(CandidateFileRow { + language: row.get(0)?, + content_hash: row.get(1)?, + size_bytes: row.get(2)?, + mtime_seconds: row.get(3)?, + mtime_nanoseconds: row.get(4)?, + package_assignment: row.get(5)?, + file_facts: row.get(6)?, + file_subgraph: row.get(7)?, + }) + }, + ) + .optional() + .map_err(|error| map_sqlite_error(error, deadline))?; + let Some(found) = found else { + return Err(CachePublicationFailure::conflict( + candidate.candidate_id, + PublicationField::FilePresence, + Some(&file.path), + None, + )); + }; + for (differs, field) in [ + (found.language != file.language, PublicationField::Language), + ( + found.content_hash != file.content_hash, + PublicationField::ContentHash, + ), + ( + found.size_bytes != file.size_bytes, + PublicationField::SizeBytes, + ), + ( + found.package_assignment != file.package_assignment, + PublicationField::PackageAssignment, + ), + (found.file_facts != file.facts, PublicationField::FileFacts), + ] { + if differs { + return Err(CachePublicationFailure::conflict( + candidate.candidate_id, + field, + Some(&file.path), + None, + )); + } + } + match (found.file_subgraph, &file.subgraph) { + (None, Some(subgraph)) => { + self.connection.execute( + "UPDATE candidate_files SET file_subgraph = ?1 WHERE candidate_id = ?2 AND path = ?3 AND file_subgraph IS NULL", + params![subgraph, candidate.candidate_id.as_slice(), file.path], + ).map_err(|error| map_sqlite_error(error, deadline))?; + } + (Some(stored), Some(incoming)) if stored != *incoming => { + return Err(CachePublicationFailure::conflict( + candidate.candidate_id, + PublicationField::FileSubgraph, + Some(&file.path), + None, + )); + } + _ => {} + } + if found.mtime_seconds != file.mtime_seconds + || found.mtime_nanoseconds != file.mtime_nanoseconds + { + update_mtime + .execute(params![ + file.mtime_seconds, + file.mtime_nanoseconds, + candidate.candidate_id.as_slice(), + file.path, + ]) + .map_err(|error| map_sqlite_error(error, deadline))?; + } + } + Ok(()) + } +} diff --git a/cli/src/cache/store/graph_publication.rs b/cli/src/cache/store/graph_publication.rs new file mode 100644 index 0000000..964d5d7 --- /dev/null +++ b/cli/src/cache/store/graph_publication.rs @@ -0,0 +1,102 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::{ + CacheError, CachePublicationFailure, CacheStore, Deadline, PreparedGraph, PublicationField, + map_sqlite_error, +}; + +impl CacheStore { + pub(super) fn verify_existing_graph( + &self, + candidate_id: [u8; 32], + snapshot_id: i64, + graph: &PreparedGraph, + deadline: &Deadline, + ) -> Result<(), CachePublicationFailure> { + let stored_symbols = + self.load_graph_payloads(snapshot_id, "graph_symbols", "symbol", deadline)?; + // Edges have no serialized copy to compare, so compare the identity the + // columns carry: `edge_key` is the lossless edge identity and + // `confidence` is the one attribute it deliberately excludes. + let stored_edges = + self.load_graph_payloads(snapshot_id, "graph_edges", "edge_key", deadline)?; + let stored_confidence = + self.load_graph_text(snapshot_id, "graph_edges", "confidence", deadline)?; + for (differs, field) in [ + ( + stored_symbols.len() != graph.symbols.len() + || stored_symbols + .iter() + .zip(&graph.symbols) + .any(|(stored, row)| *stored != row.payload), + PublicationField::GraphSymbols, + ), + ( + stored_edges.len() != graph.edges.len() + || stored_edges + .iter() + .zip(&graph.edges) + .any(|(stored, row)| *stored != row.edge_key), + PublicationField::GraphEdges, + ), + ( + stored_confidence.len() != graph.edges.len() + || stored_confidence + .iter() + .zip(&graph.edges) + .any(|(stored, row)| *stored != row.confidence), + PublicationField::GraphConfidence, + ), + ] { + if differs { + return Err(CachePublicationFailure::conflict( + candidate_id, + field, + None, + Some(graph.tier), + )); + } + } + Ok(()) + } + + fn load_graph_payloads( + &self, + snapshot_id: i64, + table: &str, + column: &str, + deadline: &Deadline, + ) -> Result>, CacheError> { + let sql = + format!("SELECT {column} FROM {table} WHERE snapshot_id = ?1 ORDER BY ordinal ASC"); + let mut statement = self + .connection + .prepare(&sql) + .map_err(|error| map_sqlite_error(error, deadline))?; + statement + .query_map([snapshot_id], |row| row.get::<_, Vec>(0)) + .map_err(|error| map_sqlite_error(error, deadline))? + .collect::, _>>() + .map_err(|error| map_sqlite_error(error, deadline)) + } + + fn load_graph_text( + &self, + snapshot_id: i64, + table: &str, + column: &str, + deadline: &Deadline, + ) -> Result, CacheError> { + let sql = + format!("SELECT {column} FROM {table} WHERE snapshot_id = ?1 ORDER BY ordinal ASC"); + let mut statement = self + .connection + .prepare(&sql) + .map_err(|error| map_sqlite_error(error, deadline))?; + statement + .query_map([snapshot_id], |row| row.get::<_, String>(0)) + .map_err(|error| map_sqlite_error(error, deadline))? + .collect::, _>>() + .map_err(|error| map_sqlite_error(error, deadline)) + } +} diff --git a/cli/src/cache/store/publication_error.rs b/cli/src/cache/store/publication_error.rs new file mode 100644 index 0000000..69d7efe --- /dev/null +++ b/cli/src/cache/store/publication_error.rs @@ -0,0 +1,143 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::{CacheError, CandidateId}; + +#[derive(Debug)] +pub(crate) enum CachePublicationFailure { + Cache(CacheError), + Conflict(PublicationConflict), +} + +#[derive(Debug)] +pub(crate) struct PublicationConflict { + pub candidate_id: CandidateId, + pub field: PublicationField, + pub file: Option, + pub tier: Option<&'static str>, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum PublicationField { + LanguageFingerprint, + PackageFingerprint, + CompatibilityId, + InputDigest, + Completeness, + InventoryFileCount, + InventoryTotalBytes, + FileCount, + Omissions, + FilePresence, + Language, + ContentHash, + SizeBytes, + PackageAssignment, + FileFacts, + FileSubgraph, + GraphSymbols, + GraphEdges, + GraphConfidence, +} + +impl PublicationField { + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::LanguageFingerprint => "language_fingerprint", + Self::PackageFingerprint => "package_fingerprint", + Self::CompatibilityId => "compatibility_id", + Self::InputDigest => "input_digest", + Self::Completeness => "completeness", + Self::InventoryFileCount => "inventory_file_count", + Self::InventoryTotalBytes => "inventory_total_bytes", + Self::FileCount => "file_count", + Self::Omissions => "omissions", + Self::FilePresence => "file_presence", + Self::Language => "language", + Self::ContentHash => "content_hash", + Self::SizeBytes => "size_bytes", + Self::PackageAssignment => "package_assignment", + Self::FileFacts => "file_facts", + Self::FileSubgraph => "file_subgraph", + Self::GraphSymbols => "graph_symbols", + Self::GraphEdges => "graph_edges", + Self::GraphConfidence => "graph_confidence", + } + } +} + +impl CachePublicationFailure { + pub(super) fn conflict( + candidate_id: [u8; 32], + field: PublicationField, + file: Option<&str>, + tier: Option<&'static str>, + ) -> Self { + Self::Conflict(PublicationConflict { + candidate_id: CandidateId::from_bytes(candidate_id), + field, + file: file.map(escaped_publication_identifier), + tier, + }) + } +} + +impl From for CachePublicationFailure { + fn from(error: CacheError) -> Self { + Self::Cache(error) + } +} + +impl From for CacheError { + fn from(error: CachePublicationFailure) -> Self { + match error { + CachePublicationFailure::Cache(error) => error, + CachePublicationFailure::Conflict(_) => Self::CandidateConflict, + } + } +} + +const IDENTIFIER_MAX_BYTES: usize = 256; + +pub(crate) fn escaped_publication_identifier(value: &str) -> String { + const MARKER: &str = "..."; + let mut output = String::with_capacity(IDENTIFIER_MAX_BYTES); + output.push('"'); + let mut chars = value.chars().peekable(); + while let Some(character) = chars.next() { + let escaped = character.escape_debug(); + let width = escaped.clone().map(char::len_utf8).sum::(); + let reserve = 1 + if chars.peek().is_some() { + MARKER.len() + } else { + 0 + }; + if output.len() + width + reserve > IDENTIFIER_MAX_BYTES { + output.push_str(MARKER); + break; + } + output.extend(escaped); + } + output.push('"'); + output +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn conflict_file_allocation_remains_bounded_after_escape_expansion() { + let CachePublicationFailure::Conflict(conflict) = CachePublicationFailure::conflict( + [1; 32], + PublicationField::FilePresence, + Some(&"\u{1b}".repeat(10_000)), + None, + ) else { + panic!("expected conflict"); + }; + let file = conflict.file.expect("file identifier"); + assert!(file.len() <= IDENTIFIER_MAX_BYTES); + assert!(file.ends_with("...\"")); + assert!(!file.chars().any(char::is_control)); + } +} diff --git a/cli/src/cache/store/tests/candidate_conflicts.rs b/cli/src/cache/store/tests/candidate_conflicts.rs new file mode 100644 index 0000000..e7c8178 --- /dev/null +++ b/cli/src/cache/store/tests/candidate_conflicts.rs @@ -0,0 +1,151 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::*; + +#[test] +fn rejects_inconsistent_candidates_and_conflicting_republication() { + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let cache_location = location(&root, temp.path()); + let store = + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); + + let mut unsorted = candidate(CacheCompleteness::Partial, ResolverCacheTier::Name); + unsorted.omissions = vec![ + crate::cache::CacheOmission { + path: "z".into(), + reason: "x".into(), + detail: "detail".into(), + }, + crate::cache::CacheOmission { + path: "a".into(), + reason: "x".into(), + detail: "detail".into(), + }, + ]; + unsorted.candidate_id = CandidateId::new( + unsorted.compatibility.id, + unsorted.input_digest, + unsorted.completeness, + &unsorted.omissions, + ); + assert!(matches!( + store.publish_candidate(&unsorted, &Deadline::new(None)), + Err(CacheError::InvalidCandidate) + )); + + let mut overflow = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + overflow.created_at_ns = u64::MAX; + assert!(matches!( + store.publish_candidate(&overflow, &Deadline::new(None)), + Err(CacheError::InvalidCandidate) + )); + + let original = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + store + .publish_candidate(&original, &Deadline::new(None)) + .expect("publish"); + let mut republished = original.clone(); + republished.created_at_ns += 1; + republished.compatibility.created_at_ns += 1; + store + .publish_candidate(&republished, &Deadline::new(None)) + .expect("timestamps are store-owned and do not conflict"); + assert_eq!( + store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Complete, + original.compatibility.id, + &Deadline::new(None), + ) + .expect("load") + .expect("active") + .created_at_ns, + original.created_at_ns + ); +} + +#[test] +fn scope_publication_requires_and_restores_every_owned_subgraph() { + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let cache_location = location(&root, temp.path()); + let store = + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); + let mut snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Scope); + assert!(matches!( + store.publish_candidate(&snapshot, &Deadline::new(None)), + Err(CacheError::InvalidCandidate) + )); + // A Name snapshot may be published first; a later Scope publication + // for the identical candidate augments its per-file subgraphs. + let mut name = snapshot.clone(); + name.tier_graphs = vec![( + ResolverCacheTier::Name, + CodeGraph { + symbols: Vec::new(), + edges: Vec::new(), + }, + )]; + store + .publish_candidate(&name, &Deadline::new(None)) + .expect("publish name"); + let mut incremental = IncrementalGraph::new(); + incremental.upsert(&snapshot.files[0].facts); + snapshot.files[0].subgraph = incremental.subgraph("src/a.rs").cloned(); + store + .publish_candidate(&snapshot, &Deadline::new(None)) + .expect("augment with scope"); + let restored = store + .hydrate_scope_subgraphs(snapshot.candidate_id, &Deadline::new(None)) + .expect("hydrate"); + assert!(restored.subgraph("src/a.rs").is_some()); +} + +#[test] +fn metadata_refresh_rejects_each_conflicting_immutable_file_payload() { + use super::metadata_refresh::{Fixture, add_scope, hint}; + + let changes = [ + "size_bytes = size_bytes + 1", + "language = 'python'", + "content_hash = zeroblob(32)", + "package_assignment = 'different-package'", + "file_facts = X'00'", + "file_subgraph = X'00'", + ]; + for change in changes { + let fixture = Fixture::new(); + let mut original = candidate(CacheCompleteness::Complete, ResolverCacheTier::Scope); + add_scope(&mut original); + fixture.publish(&original); + fixture + .store + .connection + .execute( + &format!("UPDATE candidate_files SET {change} WHERE candidate_id = ?1"), + [original.candidate_id.as_bytes().as_slice()], + ) + .expect("conflicting payload"); + let mut incoming = original.clone(); + incoming.files[0].mtime = hint(1_000_000_000, 0); + assert!( + matches!( + fixture + .store + .publish_candidate(&incoming, &Deadline::new(None)), + Err(CacheError::CandidateConflict) + ), + "{change}" + ); + let stored: (Option, Option) = fixture.store.connection.query_row( + "SELECT mtime_seconds, mtime_nanoseconds FROM candidate_files WHERE candidate_id = ?1", + [original.candidate_id.as_bytes().as_slice()], + |row| Ok((row.get(0)?, row.get(1)?)), + ).expect("stored hint"); + assert_eq!(stored, (Some(0), Some(4)), "{change}"); + } +} diff --git a/cli/src/cache/store/tests/metadata_refresh.rs b/cli/src/cache/store/tests/metadata_refresh.rs new file mode 100644 index 0000000..f93066a --- /dev/null +++ b/cli/src/cache/store/tests/metadata_refresh.rs @@ -0,0 +1,184 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::*; + +pub(super) struct Fixture { + pub store: CacheStore, + _temp: tempfile::TempDir, +} + +impl Fixture { + pub fn new() -> Self { + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let store = + CacheStore::open_writable(&location(&root, temp.path()), &root, &Deadline::new(None)) + .expect("open"); + Self { store, _temp: temp } + } + + pub fn publish(&self, snapshot: &CandidateSnapshot) { + self.store + .publish_candidate(snapshot, &Deadline::new(None)) + .expect("publish"); + } + + pub fn metadata( + &self, + snapshot: &CandidateSnapshot, + ) -> Vec { + self.store + .candidate_file_metadata(snapshot.candidate_id, &Deadline::new(None)) + .expect("metadata") + } +} + +pub(super) fn hint(seconds: i64, nanoseconds: u32) -> Option { + Some(MtimeHint { + seconds_since_unix_epoch: seconds, + nanoseconds, + }) +} + +pub(super) fn add_scope(snapshot: &mut CandidateSnapshot) { + let mut incremental = IncrementalGraph::new(); + for file in &mut snapshot.files { + incremental.upsert(&file.facts); + file.subgraph = incremental.subgraph(&file.path).cloned(); + } + snapshot.tier_graphs = vec![( + ResolverCacheTier::Scope, + CodeGraph { + symbols: Vec::new(), + edges: Vec::new(), + }, + )]; +} + +pub(super) fn two_files() -> CandidateSnapshot { + let mut snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + let mut second = snapshot.files[0].clone(); + second.path = "src/b.rs".into(); + second.facts = empty_facts(&second.path); + second.package_assignment = "10:assignment8:src/b.rs4:none".into(); + snapshot.files.push(second); + snapshot.input_digest = ProjectInputDigest::from_inputs( + snapshot + .files + .iter() + .map(|file| (&file.path, &file.language, file.content_hash)), + ); + snapshot.candidate_id = CandidateId::new( + snapshot.compatibility.id, + snapshot.input_digest, + snapshot.completeness, + &snapshot.omissions, + ); + snapshot.inventory_file_count = 2; + snapshot.inventory_total_bytes = 2; + snapshot +} + +#[test] +fn republication_updates_mtime_hints_without_replacing_content_identity() { + let transitions = [ + (None, hint(1_000_000_000, 0)), + (hint(1_000_000_000, 0), None), + (hint(1_000_000_001, 0), hint(1_000_000_000, 0)), + (hint(1_000_000_000, 1), hint(1_000_000_000, 2)), + (hint(1, 0), hint(-1, 999_999_999)), + ]; + for (before, after) in transitions { + let fixture = Fixture::new(); + let mut original = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + original.files[0].mtime = before; + fixture.publish(&original); + let original_metadata = fixture.metadata(&original); + let mut incoming = original.clone(); + incoming.files[0].mtime = after; + incoming.created_at_ns += 10; + incoming.compatibility.created_at_ns += 10; + fixture.publish(&incoming); + fixture.publish(&incoming); + let loaded = fixture + .store + .load_candidate(original.candidate_id, &Deadline::new(None)) + .expect("load"); + assert_eq!(loaded.candidate_id, original.candidate_id); + assert_eq!(loaded.input_digest, original.input_digest); + assert_eq!(loaded.created_at_ns, original.created_at_ns); + assert_eq!(loaded.compatibility, original.compatibility); + assert_eq!(loaded.files[0].mtime, after); + let mut expected_metadata = original_metadata; + expected_metadata[0].mtime = after; + assert_eq!(fixture.metadata(&incoming), expected_metadata); + assert_eq!(loaded.tier_graphs.len(), 1); + assert_eq!(loaded.tier_graphs[0].0, original.tier_graphs[0].0); + assert_eq!( + serde_json::to_value(&loaded.tier_graphs[0].1).expect("loaded graph"), + serde_json::to_value(&original.tier_graphs[0].1).expect("original graph") + ); + let counts: (i64, i64) = fixture + .store + .connection + .query_row( + "SELECT (SELECT count(*) FROM candidates), (SELECT count(*) FROM graph_snapshots)", + [], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .expect("row counts"); + assert_eq!(counts, (1, 1)); + } +} + +#[test] +fn name_scope_name_publication_keeps_graphs_and_enriched_subgraphs() { + let fixture = Fixture::new(); + let name = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + fixture.publish(&name); + let mut scope = name.clone(); + add_scope(&mut scope); + scope.files[0].mtime = hint(1_000_000_001, 0); + fixture.publish(&scope); + let mut refreshed_name = name.clone(); + refreshed_name.files[0].mtime = hint(1_000_000_000, 0); + fixture.publish(&refreshed_name); + let loaded = fixture + .store + .load_candidate(name.candidate_id, &Deadline::new(None)) + .expect("load"); + assert_eq!(loaded.files[0].mtime, refreshed_name.files[0].mtime); + assert!(loaded.files[0].subgraph.is_some()); + assert_eq!(loaded.tier_graphs.len(), 2); + for ((tier, graph), (expected_tier, expected_graph)) in loaded + .tier_graphs + .iter() + .zip([&name.tier_graphs[0], &scope.tier_graphs[0]]) + { + assert_eq!(tier, expected_tier); + assert_eq!( + serde_json::to_value(graph).expect("loaded graph"), + serde_json::to_value(expected_graph).expect("expected graph") + ); + } + let hydrated = fixture + .store + .hydrate_scope_subgraphs(name.candidate_id, &Deadline::new(None)) + .expect("hydrate"); + assert!(hydrated.subgraph("src/a.rs").is_some()); + for tier in [ResolverCacheTier::Name, ResolverCacheTier::Scope] { + let active = fixture + .store + .load_active( + tier, + CacheCompleteness::Complete, + name.compatibility.id, + &Deadline::new(None), + ) + .expect("active load") + .expect("active"); + assert_eq!(active.candidate_id, name.candidate_id); + assert_eq!(active.files[0].mtime, refreshed_name.files[0].mtime); + } +} diff --git a/cli/src/cache/store/tests/publication_diagnostics.rs b/cli/src/cache/store/tests/publication_diagnostics.rs new file mode 100644 index 0000000..d794863 --- /dev/null +++ b/cli/src/cache/store/tests/publication_diagnostics.rs @@ -0,0 +1,276 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::metadata_refresh::{Fixture, add_scope, hint, two_files}; +use super::*; + +fn conflict(store: &CacheStore, snapshot: &CandidateSnapshot) -> PublicationConflict { + match store.publish_candidate_detailed(snapshot, &Deadline::new(None)) { + Err(CachePublicationFailure::Conflict(conflict)) => conflict, + result => panic!("expected structured conflict, got {result:?}"), + } +} + +#[test] +fn stored_header_and_compatibility_differences_identify_named_fields() { + let changes = [ + ( + "compatibility", + "language_fingerprint = zeroblob(32)", + PublicationField::LanguageFingerprint, + ), + ( + "compatibility", + "package_fingerprint = zeroblob(32)", + PublicationField::PackageFingerprint, + ), + ( + "candidates", + "compatibility_id = zeroblob(32)", + PublicationField::CompatibilityId, + ), + ( + "candidates", + "input_digest = zeroblob(32)", + PublicationField::InputDigest, + ), + ( + "candidates", + "completeness = 1", + PublicationField::Completeness, + ), + ( + "candidates", + "inventory_file_count = 2", + PublicationField::InventoryFileCount, + ), + ( + "candidates", + "inventory_total_bytes = 2", + PublicationField::InventoryTotalBytes, + ), + ]; + for (table, change, field) in changes { + let fixture = Fixture::new(); + let snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + fixture.publish(&snapshot); + fixture + .store + .connection + .execute_batch("PRAGMA foreign_keys = OFF") + .expect("controlled mutation"); + fixture + .store + .connection + .execute(&format!("UPDATE {table} SET {change}"), []) + .expect("stored field"); + let diagnostic = conflict(&fixture.store, &snapshot); + assert_eq!(diagnostic.candidate_id, snapshot.candidate_id); + assert_eq!(diagnostic.field, field); + assert!(diagnostic.file.is_none()); + assert!(diagnostic.tier.is_none()); + assert!(matches!( + fixture + .store + .publish_candidate(&snapshot, &Deadline::new(None)), + Err(CacheError::CandidateConflict) + )); + } +} + +#[test] +fn immutable_file_differences_identify_fields_and_roll_back_prior_mtime_updates() { + let changes = [ + ("language = 'python'", PublicationField::Language), + ("content_hash = zeroblob(32)", PublicationField::ContentHash), + ("size_bytes = size_bytes + 1", PublicationField::SizeBytes), + ( + "package_assignment = 'secret-package'", + PublicationField::PackageAssignment, + ), + ("file_facts = X'736563726574'", PublicationField::FileFacts), + ( + "file_subgraph = X'736563726574'", + PublicationField::FileSubgraph, + ), + ]; + for (change, field) in changes { + let fixture = Fixture::new(); + let mut snapshot = two_files(); + add_scope(&mut snapshot); + fixture.publish(&snapshot); + fixture + .store + .connection + .execute( + &format!("UPDATE candidate_files SET {change} WHERE path = 'src/b.rs'"), + [], + ) + .expect("stored file field"); + let mut incoming = snapshot.clone(); + incoming.files[0].mtime = hint(1_000_000_000, 0); + let diagnostic = conflict(&fixture.store, &incoming); + assert_eq!(diagnostic.field, field); + assert_eq!(diagnostic.file.as_deref(), Some("\"src/b.rs\"")); + assert!(diagnostic.tier.is_none()); + assert_eq!( + fixture.metadata(&snapshot)[0].mtime, + snapshot.files[0].mtime + ); + } +} + +#[test] +fn file_presence_preserves_row_count_and_names_the_incoming_file() { + let fixture = Fixture::new(); + let snapshot = two_files(); + fixture.publish(&snapshot); + fixture + .store + .connection + .execute( + "UPDATE candidate_files SET path = 'other.rs' WHERE path = 'src/b.rs'", + [], + ) + .expect("stored path"); + let diagnostic = conflict(&fixture.store, &snapshot); + assert_eq!(diagnostic.field, PublicationField::FilePresence); + assert_eq!(diagnostic.file.as_deref(), Some("\"src/b.rs\"")); + let count: i64 = fixture + .store + .connection + .query_row("SELECT count(*) FROM candidate_files", [], |row| row.get(0)) + .expect("row count"); + assert_eq!(count, 2); +} + +#[test] +fn file_count_and_omissions_have_payload_free_categories() { + for (sql, field) in [ + ("DELETE FROM candidate_files", PublicationField::FileCount), + ( + "INSERT INTO candidate_omissions SELECT candidate_id, 'secret.rs', 'secret-reason', 'secret-detail' FROM candidates", + PublicationField::Omissions, + ), + ] { + let fixture = Fixture::new(); + let snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + fixture.publish(&snapshot); + fixture + .store + .connection + .execute(sql, []) + .expect("stored collection"); + let diagnostic = conflict(&fixture.store, &snapshot); + assert_eq!(diagnostic.field, field); + assert!(diagnostic.file.is_none()); + assert!(diagnostic.tier.is_none()); + } +} + +#[test] +fn graph_differences_name_the_tier_and_roll_back_metadata_updates() { + let changes = [ + ( + "INSERT INTO graph_symbols SELECT snapshot_id, 0, X'01', 'id', 'secret-name', 'secret.rs', 0, 0, 'function', X'736563726574' FROM graph_snapshots", + PublicationField::GraphSymbols, + ), + ( + "INSERT INTO graph_edges SELECT snapshot_id, 0, zeroblob(32), 0, 0, 'call', 'exact', 3, 'scope', 'secret.rs', 0, 0, 0 FROM graph_snapshots", + PublicationField::GraphEdges, + ), + ]; + for (sql, field) in changes { + let fixture = Fixture::new(); + let snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + fixture.publish(&snapshot); + fixture + .store + .connection + .execute(sql, []) + .expect("stored graph rows"); + let mut incoming = snapshot.clone(); + incoming.files[0].mtime = hint(1_000_000_000, 0); + let diagnostic = conflict(&fixture.store, &incoming); + assert_eq!(diagnostic.field, field); + assert_eq!(diagnostic.tier, Some("name")); + assert!(diagnostic.file.is_none()); + assert_eq!( + fixture.metadata(&snapshot)[0].mtime, + snapshot.files[0].mtime + ); + } +} + +#[test] +fn ordinary_publication_errors_keep_the_original_category() { + let fixture = Fixture::new(); + let snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + assert!(matches!( + fixture + .store + .publish_candidate_detailed(&snapshot, &Deadline::new(Some(Duration::ZERO))), + Err(CachePublicationFailure::Cache(CacheError::Timeout)) + )); + fixture.store.connection.execute_batch("CREATE TEMP TRIGGER reject_candidate BEFORE INSERT ON candidates BEGIN SELECT RAISE(ABORT, 'secret SQL detail'); END").expect("SQL trigger"); + assert!(matches!( + fixture + .store + .publish_candidate_detailed(&snapshot, &Deadline::new(None)), + Err(CachePublicationFailure::Cache(CacheError::Access)) + )); + fixture + .store + .connection + .execute_batch("DROP TRIGGER reject_candidate") + .expect("remove trigger"); + let mut readonly = fixture.store; + readonly.writable = false; + assert!(matches!( + readonly.publish_candidate_detailed(&snapshot, &Deadline::new(None)), + Err(CachePublicationFailure::Cache(CacheError::ReadOnly)) + )); +} + +#[test] +fn graph_confidence_difference_keeps_edge_identity_and_names_the_tier() { + let fixture = Fixture::new(); + let mut snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + snapshot.tier_graphs[0].1.edges.push(Edge { + from: SymbolId::global("rust", vec![Descriptor::Term("from".into())]), + to: SymbolId::global("rust", vec![Descriptor::Term("to".into())]), + role: RefRole::Call, + confidence: Confidence::Scoped, + provenance: Provenance::ScopeGraph, + occ: Occurrence { + file: "src/a.rs".into(), + byte: 0, + line: 1, + col: 0, + }, + }); + fixture.publish(&snapshot); + fixture + .store + .connection + .execute( + "UPDATE graph_edges SET confidence = 'secret-confidence'", + [], + ) + .expect("stored confidence"); + let diagnostic = conflict(&fixture.store, &snapshot); + assert_eq!(diagnostic.field, PublicationField::GraphConfidence); + assert_eq!(diagnostic.tier, Some("name")); + assert!(diagnostic.file.is_none()); +} + +#[test] +fn multiple_stored_differences_choose_the_first_named_field() { + let fixture = Fixture::new(); + let snapshot = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + fixture.publish(&snapshot); + fixture.store.connection.execute("UPDATE candidate_files SET language = 'python', content_hash = zeroblob(32), file_facts = X'00'", []).expect("stored differences"); + assert_eq!( + conflict(&fixture.store, &snapshot).field, + PublicationField::Language + ); +} diff --git a/cli/src/cache/store/tests/publication_slots.rs b/cli/src/cache/store/tests/publication_slots.rs new file mode 100644 index 0000000..ee22d90 --- /dev/null +++ b/cli/src/cache/store/tests/publication_slots.rs @@ -0,0 +1,206 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::*; + +#[test] +fn candidate_publication_keeps_complete_and_partial_slots_independent() { + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let cache_location = location(&root, temp.path()); + let store = + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); + let complete = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + let partial = candidate(CacheCompleteness::Partial, ResolverCacheTier::Name); + store + .publish_candidate(&complete, &Deadline::new(None)) + .expect("publish complete"); + store + .publish_candidate(&partial, &Deadline::new(None)) + .expect("publish partial"); + assert_eq!( + store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Complete, + complete.compatibility.id, + &Deadline::new(None) + ) + .expect("load") + .expect("active") + .candidate_id, + complete.candidate_id + ); + assert_eq!( + store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Partial, + partial.compatibility.id, + &Deadline::new(None) + ) + .expect("load") + .expect("active") + .candidate_id, + partial.candidate_id + ); + let incompatible = CompatibilityFingerprint::new( + crate::cache::LanguageFeatureFingerprint::current(), + crate::cache::PackageFingerprint::from_normalized(["different-package"]), + ); + let loaded_complete = store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Complete, + complete.compatibility.id, + &Deadline::new(None), + ) + .expect("load") + .expect("active"); + assert_eq!( + loaded_complete.compatibility.language_fingerprint, + complete.compatibility.language_fingerprint + ); + assert_eq!( + loaded_complete.compatibility.package_fingerprint, + complete.compatibility.package_fingerprint + ); + assert!( + store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Complete, + incompatible, + &Deadline::new(None), + ) + .expect("compatibility miss") + .is_none() + ); + store + .publish_candidate(&complete, &Deadline::new(None)) + .expect("idempotent publish"); +} + +#[test] +fn superseding_a_slot_garbage_collects_the_prior_snapshot_and_candidate() { + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let cache_location = location(&root, temp.path()); + let store = + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); + // Two distinct candidates (different input digests) target the same + // (tier, completeness) slot; publishing B flips active away from A. + let a = candidate_with_hash( + CacheCompleteness::Complete, + ResolverCacheTier::Name, + [3; 32], + ); + let b = candidate_with_hash( + CacheCompleteness::Complete, + ResolverCacheTier::Name, + [7; 32], + ); + assert_ne!(a.candidate_id, b.candidate_id); + store + .publish_candidate(&a, &Deadline::new(None)) + .expect("publish a"); + store + .publish_candidate(&b, &Deadline::new(None)) + .expect("publish b"); + + // Only B's snapshot survives; A's snapshot and candidate rows are gone. + let snapshot_count: i64 = store + .connection + .query_row("SELECT count(*) FROM graph_snapshots", [], |row| row.get(0)) + .expect("snapshot count"); + assert_eq!(snapshot_count, 1); + let surviving_candidate: Vec = store + .connection + .query_row("SELECT candidate_id FROM graph_snapshots", [], |row| { + row.get(0) + }) + .expect("surviving candidate"); + assert_eq!( + surviving_candidate.as_slice(), + b.candidate_id.as_bytes().as_slice() + ); + let a_candidate_count: i64 = store + .connection + .query_row( + "SELECT count(*) FROM candidates WHERE candidate_id = ?1", + [a.candidate_id.as_bytes().as_slice()], + |row| row.get(0), + ) + .expect("a candidate count"); + assert_eq!(a_candidate_count, 0); + + // B remains the queryable active snapshot for the slot. + assert_eq!( + store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Complete, + b.compatibility.id, + &Deadline::new(None), + ) + .expect("load") + .expect("active") + .candidate_id, + b.candidate_id + ); +} + +#[test] +fn concurrent_publishers_commit_whole_candidates() { + use std::sync::{Arc, Barrier}; + + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let cache_location = location(&root, temp.path()); + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("initialize"); + let barrier = Arc::new(Barrier::new(2)); + let handles: Vec<_> = [CacheCompleteness::Complete, CacheCompleteness::Partial] + .into_iter() + .map(|completeness| { + let barrier = Arc::clone(&barrier); + let root = root.clone(); + let cache_location = cache_location.clone(); + std::thread::spawn(move || { + let store = + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None))?; + let candidate = candidate(completeness, ResolverCacheTier::Name); + barrier.wait(); + store.publish_candidate(&candidate, &Deadline::new(None))?; + Ok::<_, CacheError>(candidate.candidate_id) + }) + }) + .collect(); + let ids: Vec<_> = handles + .into_iter() + .map(|handle| handle.join().expect("publisher thread").expect("publish")) + .collect(); + let store = + CacheStore::open_frozen(&cache_location, &root, &Deadline::new(None)).expect("frozen"); + for (completeness, id) in [CacheCompleteness::Complete, CacheCompleteness::Partial] + .into_iter() + .zip(ids) + { + assert_eq!( + store + .load_active( + ResolverCacheTier::Name, + completeness, + candidate(completeness, ResolverCacheTier::Name) + .compatibility + .id, + &Deadline::new(None), + ) + .expect("load") + .expect("active") + .candidate_id, + id + ); + } +} diff --git a/cli/src/cache/store/tests/publication_transactions.rs b/cli/src/cache/store/tests/publication_transactions.rs new file mode 100644 index 0000000..3ba41e7 --- /dev/null +++ b/cli/src/cache/store/tests/publication_transactions.rs @@ -0,0 +1,147 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::*; + +#[test] +fn failed_graph_write_rolls_back_candidate_and_active_publication() { + let temp = tempdir().expect("tempdir"); + let root = temp.path().join("project"); + fs::create_dir(&root).expect("project"); + let cache_location = location(&root, temp.path()); + let store = + CacheStore::open_writable(&cache_location, &root, &Deadline::new(None)).expect("open"); + let candidate = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + store.connection.execute_batch( + "CREATE TEMP TRIGGER fail_graph BEFORE INSERT ON graph_snapshots BEGIN SELECT RAISE(ABORT, 'injected graph failure'); END", + ).expect("failure trigger"); + assert!(matches!( + store.publish_candidate(&candidate, &Deadline::new(None)), + Err(CacheError::Access) + )); + let candidate_count: i64 = store + .connection + .query_row( + "SELECT count(*) FROM candidates WHERE candidate_id = ?1", + [candidate.candidate_id.as_bytes().as_slice()], + |row| row.get(0), + ) + .expect("candidate count"); + let active_count: i64 = store + .connection + .query_row("SELECT count(*) FROM active_snapshots", [], |row| { + row.get(0) + }) + .expect("active count"); + assert_eq!((candidate_count, active_count), (0, 0)); + store + .connection + .execute_batch("DROP TRIGGER fail_graph") + .expect("drop trigger"); + store + .publish_candidate(&candidate, &Deadline::new(None)) + .expect("retry"); +} + +#[test] +fn later_file_conflict_rolls_back_earlier_metadata_updates() { + use super::metadata_refresh::{Fixture, hint, two_files}; + + let fixture = Fixture::new(); + let original = two_files(); + fixture.publish(&original); + let mut incoming = original.clone(); + incoming.files[0].mtime = hint(1_000_000_000, 0); + incoming.files[1].mtime = hint(1_000_000_001, 0); + fixture.store.connection.execute( + "UPDATE candidate_files SET size_bytes = size_bytes + 1 WHERE candidate_id = ?1 AND path = 'src/b.rs'", + [original.candidate_id.as_bytes().as_slice()], + ).expect("conflicting second file size"); + assert!(matches!( + fixture + .store + .publish_candidate(&incoming, &Deadline::new(None)), + Err(CacheError::CandidateConflict) + )); + let metadata = fixture.metadata(&original); + assert_eq!(metadata[0].mtime, original.files[0].mtime); + assert_eq!(metadata[1].mtime, original.files[1].mtime); + fixture.store.connection.execute( + "UPDATE candidate_files SET size_bytes = ?1 WHERE candidate_id = ?2 AND path = 'src/b.rs'", + params![i64::try_from(original.files[1].size_bytes).expect("SQLite size"), original.candidate_id.as_bytes().as_slice()], + ).expect("restore second file size"); + fixture.publish(&original); +} + +#[test] +fn graph_rejection_rolls_back_metadata_and_subgraphs_and_preserves_active_slots() { + use super::metadata_refresh::{Fixture, add_scope, hint}; + + let fixture = Fixture::new(); + let original = candidate(CacheCompleteness::Complete, ResolverCacheTier::Name); + fixture.publish(&original); + let mut scope = original.clone(); + add_scope(&mut scope); + scope.files[0].mtime = hint(1_000_000_000, 0); + fixture.store.connection.execute_batch( + "CREATE TEMP TRIGGER reject_scope BEFORE INSERT ON graph_snapshots BEGIN SELECT RAISE(ABORT, 'rejected graph insert'); END", + ).expect("graph trigger"); + assert!(matches!( + fixture + .store + .publish_candidate(&scope, &Deadline::new(None)), + Err(CacheError::Access) + )); + let metadata = fixture.metadata(&original); + assert_eq!(metadata[0].mtime, original.files[0].mtime); + assert!(!metadata[0].has_subgraph); + let active = fixture + .store + .load_active( + ResolverCacheTier::Name, + CacheCompleteness::Complete, + original.compatibility.id, + &Deadline::new(None), + ) + .expect("name load") + .expect("name active"); + assert_eq!(active.candidate_id, original.candidate_id); + assert!( + fixture + .store + .load_active( + ResolverCacheTier::Scope, + CacheCompleteness::Complete, + original.compatibility.id, + &Deadline::new(None) + ) + .expect("scope load") + .is_none() + ); + let graph_count: i64 = fixture + .store + .connection + .query_row("SELECT count(*) FROM graph_snapshots", [], |row| row.get(0)) + .expect("graph count"); + assert_eq!(graph_count, 1); + fixture + .store + .connection + .execute_batch("DROP TRIGGER reject_scope") + .expect("drop trigger"); + fixture.publish(&scope); + let metadata = fixture.metadata(&scope); + assert_eq!(metadata[0].mtime, scope.files[0].mtime); + assert!(metadata[0].has_subgraph); + assert!( + fixture + .store + .load_active( + ResolverCacheTier::Scope, + CacheCompleteness::Complete, + scope.compatibility.id, + &Deadline::new(None) + ) + .expect("scope load") + .is_some() + ); +} diff --git a/cli/src/refresh/mod.rs b/cli/src/refresh/mod.rs index b26d731..8b8ac35 100644 --- a/cli/src/refresh/mod.rs +++ b/cli/src/refresh/mod.rs @@ -4,6 +4,7 @@ mod plan; mod prepare; +mod publication_error; mod publish; mod resolve; mod types; diff --git a/cli/src/refresh/publication_error.rs b/cli/src/refresh/publication_error.rs new file mode 100644 index 0000000..e0c7be4 --- /dev/null +++ b/cli/src/refresh/publication_error.rs @@ -0,0 +1,83 @@ +// SPDX-License-Identifier: Apache-2.0 + +use std::path::Path; + +use crate::CliError; +use crate::cache::{CachePublicationFailure, PublicationConflict, escaped_publication_identifier}; + +pub(super) fn publication_error(error: CachePublicationFailure, root: &Path) -> CliError { + match error { + CachePublicationFailure::Cache(error) => error.into(), + CachePublicationFailure::Conflict(conflict) => { + CliError::Cache(render_conflict(&conflict, root)) + } + } +} + +fn render_conflict(conflict: &PublicationConflict, root: &Path) -> String { + let mut message = format!( + "cache candidate {} conflicts in {}", + conflict.candidate_id, + conflict.field.as_str(), + ); + if let Some(file) = &conflict.file { + message.push_str(" for file "); + message.push_str(file); + } + if let Some(tier) = conflict.tier { + message.push_str(" for tier "); + message.push_str(tier); + } + message.push_str(". Selected project root: "); + message.push_str(&escaped_publication_identifier(&root.to_string_lossy())); + message.push_str( + ". Run `c2g cache clear --root ` with that root, then retry.", + ); + message +} + +#[cfg(test)] +mod tests { + use super::*; + const IDENTIFIER_MAX_BYTES: usize = 256; + + #[test] + fn identifiers_escape_terminal_controls_quotes_and_backslashes() { + let rendered = escaped_publication_identifier("file\n\r\t\u{1b}\"\\"); + assert_eq!(rendered, "\"file\\n\\r\\t\\u{1b}\\\"\\\\\""); + assert!(!rendered.chars().any(char::is_control)); + } + + #[test] + fn escaped_expansion_and_unicode_remain_bounded() { + for value in [ + "\u{1b}".repeat(10_000), + "界".repeat(10_000), + "x".repeat(256), + ] { + let rendered = escaped_publication_identifier(&value); + assert!(rendered.len() <= IDENTIFIER_MAX_BYTES); + assert!(rendered.ends_with("...\"")); + assert!(!rendered.chars().any(char::is_control)); + } + assert_eq!(escaped_publication_identifier("short"), "\"short\""); + } + + #[test] + fn ordinary_errors_retain_cli_text_and_exit_code() { + for error in [ + crate::cache::CacheError::Access, + crate::cache::CacheError::Timeout, + crate::cache::CacheError::ReadOnly, + ] { + let expected_text = error.to_string(); + let actual = + publication_error(CachePublicationFailure::Cache(error), Path::new("/root")); + assert_eq!( + actual.to_string(), + CliError::Cache(expected_text).to_string() + ); + assert_eq!(actual.exit_code(), crate::ExitCode::Operational); + } + } +} diff --git a/cli/src/refresh/publish.rs b/cli/src/refresh/publish.rs index 45f8590..c0a95ff 100644 --- a/cli/src/refresh/publish.rs +++ b/cli/src/refresh/publish.rs @@ -117,7 +117,11 @@ fn prepare_and_publish_inner( // CacheStore begins its SQLite write transaction only here. All // filesystem discovery, bounded reads, extraction, and revalidation // above deliberately happen outside that transaction. - store.publish_candidate(&prepared.snapshot, inputs.deadline)?; + store + .publish_candidate_detailed(&prepared.snapshot, inputs.deadline) + .map_err(|error| { + super::publication_error::publication_error(error, &inputs.selection.canonical_root) + })?; inputs.deadline.check(inputs.cancellation)?; // The snapshot was just published verbatim; build the loaded view from // the in-memory candidate instead of re-decoding it back out of SQLite. @@ -357,6 +361,8 @@ mod tests { use crate::worker::RequestId; use crate::{Cancellation, Deadline, NeverCancelled}; + mod publication_diagnostics; + struct Extractor; impl FactsExtractor for Extractor { type Session = ExtractorSession; diff --git a/cli/src/refresh/publish/tests/publication_diagnostics.rs b/cli/src/refresh/publish/tests/publication_diagnostics.rs new file mode 100644 index 0000000..7a2e543 --- /dev/null +++ b/cli/src/refresh/publish/tests/publication_diagnostics.rs @@ -0,0 +1,61 @@ +// SPDX-License-Identifier: Apache-2.0 + +use super::*; + +struct CorruptPublishedFacts { + database: std::path::PathBuf, +} + +impl PublicationHook for CorruptPublishedFacts { + fn before_publish(&self) -> Result<()> { + let connection = rusqlite::Connection::open(&self.database) + .map_err(|error| CliError::Index(error.to_string()))?; + connection + .execute( + "UPDATE candidate_files SET file_facts = X'736563726574'", + [], + ) + .map_err(|error| CliError::Index(error.to_string()))?; + Ok(()) + } +} + +#[test] +fn publication_conflict_reaches_refresh_with_selected_root_and_recovery_command() { + let (temp, selection) = project("fn secret_source() {}\n"); + let limits = ResourceLimits::default(); + let deadline = Deadline::new(None); + let store = store(&temp, &selection); + let first = prepare_and_publish_with( + &Extractor, + &store, + inputs(&selection, &limits, &deadline), + false, + ) + .expect("initial publication"); + let location = CacheLocation::for_project(Some(temp.path()), &selection.canonical_root) + .expect("cache location"); + let result = prepare_and_publish_with_hook( + &Extractor, + &store, + inputs(&selection, &limits, &deadline), + false, + &CorruptPublishedFacts { + database: location.database_path, + }, + ); + let Err(CliError::Cache(message)) = result else { + panic!("expected refresh publication conflict"); + }; + assert!(message.contains(&first.loaded.candidate_id.to_string())); + assert!(message.contains("file_facts for file \"a.rs\"")); + assert!( + message.contains(&crate::cache::escaped_publication_identifier( + &selection.canonical_root.to_string_lossy(), + )) + ); + assert!(message.contains("c2g cache clear --root ")); + assert!(!message.contains("secret")); + assert!(!message.contains("UPDATE")); + assert!(!message.contains("--all")); +} diff --git a/cli/tests/cache_metadata_refresh.rs b/cli/tests/cache_metadata_refresh.rs new file mode 100644 index 0000000..0c1b1f6 --- /dev/null +++ b/cli/tests/cache_metadata_refresh.rs @@ -0,0 +1,231 @@ +// SPDX-License-Identifier: Apache-2.0 + +use std::fs::{self, File, FileTimes}; +use std::path::PathBuf; +use std::process::Command; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use rusqlite::Connection; +use serde_json::Value; + +struct Fixture { + _temp: tempfile::TempDir, + root: PathBuf, + _home: PathBuf, + cache_directory: PathBuf, + database: PathBuf, +} + +impl Fixture { + fn new() -> Self { + let temp = tempfile::tempdir().expect("tempdir"); + let root = temp.path().join("project"); + let home = temp.path().join("home"); + fs::create_dir(&root).expect("project directory"); + fs::create_dir(&home).expect("home directory"); + for (path, source) in [ + ("a.rs", "pub fn alpha() {}\n"), + ("b.rs", "pub fn bravo() { alpha(); }\n"), + ("c.rs", "pub fn charlie() { bravo(); }\n"), + ] { + let path = root.join(path); + fs::write(&path, source).expect("source"); + set_mtime(&path, 1_000_000_000); + } + let mut fixture = Self { + _temp: temp, + root: root.canonicalize().expect("canonical project"), + _home: home, + cache_directory: PathBuf::new(), + database: PathBuf::new(), + }; + let paths = fixture.run(&["cache", "path"]); + fixture.cache_directory = + PathBuf::from(paths["cache_dir"].as_str().expect("cache directory")); + fixture.database = PathBuf::from(paths["database_path"].as_str().expect("database path")); + #[cfg(unix)] + assert!(fixture.cache_directory.starts_with(&fixture._home)); + fixture + } + + fn command(&self) -> Command { + let mut command = Command::new(env!("CARGO_BIN_EXE_c2g")); + command.current_dir(&self.root); + #[cfg(unix)] + command + .env("HOME", &self._home) + .env("XDG_CACHE_HOME", self._home.join("cache")); + // Windows known folders use unique canonical roots to isolate project keys. + // Drop removes only this partition. + command.env("CODE2GRAPH_AUTO_PRUNE", "off"); + command + } + + fn run(&self, args: &[&str]) -> Value { + let output = self + .command() + .args(args) + .arg("--json") + .output() + .expect("run c2g"); + assert!( + output.status.success(), + "{args:?}: {}", + String::from_utf8_lossy(&output.stderr) + ); + assert!( + output.stderr.is_empty(), + "{args:?}: {}", + String::from_utf8_lossy(&output.stderr) + ); + serde_json::from_slice(&output.stdout).expect("JSON output") + } + + fn touch(&self, seconds: u64) -> SystemTime { + let path = self.root.join("a.rs"); + let before = fs::metadata(&path) + .expect("metadata") + .modified() + .expect("mtime"); + let after = set_mtime(&path, seconds); + assert_ne!(before, after); + after + } + + fn stored_mtime(&self) -> SystemTime { + let connection = Connection::open(&self.database).expect("cache database"); + let (seconds, nanoseconds): (i64, u32) = connection + .query_row( + "SELECT mtime_seconds, mtime_nanoseconds FROM candidate_files WHERE path = 'a.rs'", + [], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .expect("stored mtime"); + UNIX_EPOCH + + Duration::new( + u64::try_from(seconds).expect("nonnegative seconds"), + nanoseconds, + ) + } +} + +impl Drop for Fixture { + fn drop(&mut self) { + if !self.cache_directory.as_os_str().is_empty() + && self.cache_directory.exists() + && let Err(error) = fs::remove_dir_all(&self.cache_directory) + { + if std::thread::panicking() { + eprintln!("remove fixture cache partition: {error}"); + } else { + panic!("remove fixture cache partition: {error}"); + } + } + } +} + +fn set_mtime(path: &std::path::Path, seconds: u64) -> SystemTime { + let desired = UNIX_EPOCH + Duration::from_secs(seconds); + File::options() + .write(true) + .open(path) + .expect("source file") + .set_times(FileTimes::new().set_modified(desired)) + .expect("set mtime"); + let actual = fs::metadata(path) + .expect("metadata") + .modified() + .expect("mtime"); + assert_eq!(actual, desired); + actual +} + +fn assert_identity(value: &Value, snapshot: &Value) { + assert_eq!(value["status"], "ok"); + assert_eq!(value["project"]["snapshot"], *snapshot); + assert_eq!(value["project"]["freshness"], "fresh"); +} + +#[test] +fn timestamp_only_changes_refresh_index_symbols_and_status_then_reuse_cache() { + let fixture = Fixture::new(); + let initial = fixture.run(&["index"]); + assert_eq!(initial["results"]["inventory_file_count"], 3); + let snapshot = &initial["project"]["snapshot"]; + for (index, args) in [&["index"][..], &["symbols", "alpha"][..], &["status"][..]] + .into_iter() + .enumerate() + { + let actual = fixture.touch(1_000_000_010 + index as u64 * 10); + let refreshed = fixture.run(args); + assert_identity(&refreshed, snapshot); + assert_eq!(fixture.stored_mtime(), actual); + let repeated = fixture.run(args); + assert_identity(&repeated, snapshot); + assert_eq!(repeated["project"]["cache"], "hit"); + assert_eq!(fixture.stored_mtime(), actual); + if args[0] == "index" { + assert_eq!(refreshed["results"]["changed"], 0); + assert_eq!(repeated["results"]["changed"], 0); + } + if args[0] == "symbols" { + assert_eq!(refreshed["total"], 1); + assert_eq!(repeated["results"], refreshed["results"]); + } + } +} + +#[test] +fn trust_mtime_refreshes_changed_hints_and_reuses_unchanged_hints() { + let fixture = Fixture::new(); + let initial = fixture.run(&["index", "--trust-mtime"]); + let snapshot = &initial["project"]["snapshot"]; + let actual = fixture.touch(1_000_000_100); + let refreshed = fixture.run(&["index", "--trust-mtime"]); + assert_identity(&refreshed, snapshot); + assert_eq!(fixture.stored_mtime(), actual); + assert_eq!(refreshed["results"]["changed"], 0); + assert_eq!(refreshed["results"]["plan_decisions"]["reuse_facts"], 3); + assert_eq!(refreshed["results"]["plan_decisions"]["extract"], 0); + let repeated = fixture.run(&["index", "--trust-mtime"]); + assert_identity(&repeated, snapshot); + assert_eq!(repeated["project"]["cache"], "hit"); + assert_eq!(repeated["results"]["changed"], 0); + assert_eq!(repeated["results"]["attempts"], 0); + assert_eq!( + repeated["results"]["plan_decisions"], + serde_json::json!({ + "need_hash": 0, "reuse_facts": 0, "extract": 0, + "remove": 0, "omit": 0, + }) + ); + assert_eq!(fixture.stored_mtime(), actual); +} + +#[test] +fn default_refresh_detects_same_size_content_changes_with_unchanged_mtime() { + let fixture = Fixture::new(); + let initial = fixture.run(&["index"]); + let path = fixture.root.join("a.rs"); + let before = fs::metadata(&path).expect("metadata"); + fs::write(&path, "pub fn delta() {}\n").expect("changed source"); + assert_eq!(fs::metadata(&path).expect("metadata").len(), before.len()); + let restored = set_mtime(&path, 1_000_000_000); + assert_eq!(restored, before.modified().expect("mtime")); + let refreshed = fixture.run(&["index"]); + assert_eq!(refreshed["status"], "ok"); + assert_ne!( + refreshed["project"]["snapshot"], + initial["project"]["snapshot"] + ); + assert_eq!(refreshed["results"]["changed"], 1); + let symbols = fixture.run(&["symbols", "delta"]); + assert_eq!(symbols["total"], 1); + assert_eq!( + symbols["project"]["snapshot"], + refreshed["project"]["snapshot"] + ); + let repeated = fixture.run(&["index"]); + assert_eq!(repeated["results"]["changed"], 0); + assert_eq!(repeated["project"]["cache"], "hit"); +} diff --git a/cli/tests/contracts.rs b/cli/tests/contracts.rs index 6993a8d..3ea1ba1 100644 --- a/cli/tests/contracts.rs +++ b/cli/tests/contracts.rs @@ -19,6 +19,8 @@ use code2graph_cli::{ fn legacy_public_error_enum_shapes_compile() { let cache_error = CacheError::InvalidFacts; assert!(matches!(cache_error, CacheError::InvalidFacts)); + let conflict = CacheError::CandidateConflict; + assert!(matches!(conflict, CacheError::CandidateConflict)); let worker_error = WorkerErrorCode::Extraction; assert_eq!(worker_error as u16, 1); assert_eq!(WorkerErrorCode::InvalidRequest as u16, 2);