diff --git a/crates/graphforge-api/src/resumable_construction.rs b/crates/graphforge-api/src/resumable_construction.rs index e9db90082..7941cb684 100644 --- a/crates/graphforge-api/src/resumable_construction.rs +++ b/crates/graphforge-api/src/resumable_construction.rs @@ -498,6 +498,99 @@ mod tests { .unwrap() } + /// A published construction ships its adjacency CSR (#1388): a fresh + /// process finds it current in the hydrated workspace without rebuilding, + /// hop queries answer from it, and a corrupted published shard is refused + /// by the open-time digest sweep like every other graph object. + #[test] + fn published_construction_serves_its_adjacency_index_and_refuses_corruption() { + use graphforge_storage::adjacency::AdjacencyFreshnessState; + + let directory = tempfile::TempDir::new().unwrap(); + let path = directory.path().to_str().unwrap(); + let graph = GraphForge::new(Some(path)).unwrap(); + let node_ids: Vec = (0..4).map(|_| Uuid::now_v7()).collect(); + let edge_ids: Vec = (0..3).map(|_| Uuid::now_v7()).collect(); + let mut session = graph.begin_graph_construction(Default::default()).unwrap(); + session.append_nodes("nodes", &nodes(&node_ids)).unwrap(); + session + .append_edges( + "edges", + &edges( + &edge_ids, + &[ + (node_ids[0], node_ids[1]), + (node_ids[1], node_ids[2]), + (node_ids[2], node_ids[3]), + ], + ), + ) + .unwrap(); + session.seal_and_publish().unwrap(); + drop(session); + let inspection = graph.inspect_adjacency().unwrap(); + assert_eq!(inspection.state, AdjacencyFreshnessState::Current); + assert_eq!(inspection.artifact_source_generation, Some(1)); + drop(graph); + + let reopened = GraphForge::new(Some(path)).unwrap(); + assert!(graphforge_storage::adjacency::manifest_path(&reopened.dir()).is_file()); + let inspection = reopened.inspect_adjacency().unwrap(); + assert_eq!(inspection.state, AdjacencyFreshnessState::Current); + assert_eq!(inspection.artifact_source_generation, Some(1)); + assert_construction_relationships( + &reopened, + &[ + [node_ids[0], edge_ids[0], node_ids[1]], + [node_ids[1], edge_ids[1], node_ids[2]], + [node_ids[2], edge_ids[2], node_ids[3]], + ], + ); + let hops = reopened + .execute("MATCH (a)-[r]->(b)-[s]->(c) RETURN count(*) AS n") + .unwrap(); + assert_eq!( + hops.batches[0] + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + 2 + ); + drop(reopened); + + let generation = graphforge_storage::resolve_project_generation(directory.path()).unwrap(); + let inventory = generation.graph_files_inventory().unwrap().unwrap(); + let shard = inventory + .files + .iter() + .find(|entry| { + entry.relative_path.starts_with("indexes/adjacency/") + && entry.relative_path.ends_with(".csr") + }) + .expect("published CSR shard"); + let object = graphforge_storage::graph_object_path( + generation.container_root(), + &shard.content_sha256, + ) + .unwrap(); + let mut permissions = std::fs::metadata(&object).unwrap().permissions(); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt as _; + permissions.set_mode(0o600); + } + #[cfg(not(unix))] + permissions.set_readonly(false); + std::fs::set_permissions(&object, permissions).unwrap(); + let mut bytes = std::fs::read(&object).unwrap(); + bytes[0] ^= 1; + std::fs::write(&object, bytes).unwrap(); + let error = GraphForge::new(Some(path)).unwrap_err(); + assert!(error.to_string().contains("digest"), "{error}"); + } + #[test] fn cancelled_public_construction_preserves_current() { let graph = GraphForge::new(None).unwrap(); diff --git a/crates/graphforge-api/tests/permanent_storage_budgets.rs b/crates/graphforge-api/tests/permanent_storage_budgets.rs index 6517669ac..59eb25620 100644 --- a/crates/graphforge-api/tests/permanent_storage_budgets.rs +++ b/crates/graphforge-api/tests/permanent_storage_budgets.rs @@ -955,17 +955,22 @@ fn assess(f: Fixture) { construct(&source, f, &nodes, &edges); let construction_ns = started.elapsed().as_nanos(); let graph = GraphForge::new(source.to_str()).unwrap(); - let unindexed = graph.storage_attribution().unwrap(); - assert_eq!( - unindexed.categories[&graphforge_storage::ArtifactCategory::Adjacency].allocated_bytes, - 0 + // ADR 0037: construction publishes the adjacency CSR with the generation, so + // a freshly constructed project is already indexed. This previously asserted + // the category was empty, which was true only while adjacency did not exist + // until some reader built it. + let constructed = graph.storage_attribution().unwrap(); + assert!( + constructed.categories[&graphforge_storage::ArtifactCategory::Adjacency].allocated_bytes + > 0, + "construction must publish adjacency with the generation (ADR 0037)" ); validate_storage_budgets( Fixture { - adjacency: false, + adjacency: true, ..f }, - &unindexed, + &constructed, ); // Candidate codecs inspect construction's published payload generation. // Adjacency publication may switch to a generation-owned graph tree. @@ -983,7 +988,12 @@ fn assess(f: Fixture) { let identity_experiment = identity_padding_experiment(&source); let manifest_experiment = bucket_manifest_experiment(&source); if f.heterogeneous { - assert_eq!(manifest_experiment["authenticated_lookups_checked"], 352); + // 352 before ADR 0037. Construction now publishes the adjacency CSR with + // the generation, so `indexes/adjacency/**` are graph files, the manifest + // Merkle tree covers them, and the authenticated walk visits their leaves + // and the interior nodes above them. The count is a pinned observation of + // that tree's shape, not an invariant, so it moves when the file set does. + assert_eq!(manifest_experiment["authenticated_lookups_checked"], 373); assert!( manifest_experiment["source_manifest_allocated_bytes"] .as_u64() @@ -998,18 +1008,27 @@ fn assess(f: Fixture) { let adjacency_experiment = f.adjacency.then(|| adjacency_codec_experiment(&source)); let storage = graph.storage_attribution().unwrap(); storage.validate_for_qualification().unwrap(); - validate_storage_budgets(f, &storage); - assert_eq!( + validate_storage_budgets( + Fixture { + adjacency: true, + ..f + }, + &storage, + ); + // Before ADR 0037 this asserted presence == `f.adjacency`, because an + // unbuilt project carried no adjacency. Construction now always publishes + // it, so the remaining property is that an explicit rebuild does not + // remove it and does not push the category past its budget. + assert!( storage.categories[&graphforge_storage::ArtifactCategory::Adjacency].allocated_bytes > 0, - f.adjacency, - "unbuilt and built adjacency are distinct measured capabilities" + "a constructed project carries adjacency whether or not it was rebuilt" ); drop(graph); let whole_project = whole_project_file_census(&source); let fingerprint = round_trip(root.path(), &source, f, &nodes, &edges); println!( "PERMANENT_STORAGE_ASSESSMENT {}", - json!({"fixture":f.name,"nodes":f.nodes,"edges":f.edges,"routes":f.routes,"random_ids":matches!(f.identifiers, Identifiers::Random),"properties":f.properties,"heterogeneous_schemas":f.heterogeneous,"adjacency_built":f.adjacency,"unindexed_allocated_bytes":unindexed.allocated_bytes,"unindexed_categories":unindexed.categories, + json!({"fixture":f.name,"nodes":f.nodes,"edges":f.edges,"routes":f.routes,"random_ids":matches!(f.identifiers, Identifiers::Random),"properties":f.properties,"heterogeneous_schemas":f.heterogeneous,"adjacency_built":f.adjacency,"constructed_allocated_bytes":constructed.allocated_bytes,"constructed_categories":constructed.categories, "construction_elapsed_ns":construction_ns,"semantic_fingerprint":fingerprint,"whole_project":whole_project, "permanent": {"logical_bytes":storage.logical_bytes,"physical_logical_bytes":storage.physical_logical_bytes,"allocated_bytes":storage.allocated_bytes,"physical_objects":storage.physical_objects,"categories":storage.categories}, "adjacency_codec_experiment":adjacency_experiment,"parquet_experiment":experiment,"identity_padding_experiment":identity_experiment,"manifest_bucket_experiment":manifest_experiment}) diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index 3eca86ad8..5eb96496e 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -37,6 +37,7 @@ mod builder; mod codec; +pub(crate) use builder::build_adjacency_index_for_edge_files; pub use builder::{ ADJACENCY_SPILL_DIR_NAME, AdjacencyBuildMetrics, AdjacencyBuildOptions, DEFAULT_ADJACENCY_CHUNK_ROWS, DEFAULT_ADJACENCY_MERGE_FAN_IN, build_adjacency_index, @@ -1170,21 +1171,40 @@ fn capture_adjacency_inventory( crate::AuthenticatedPropertyInventory::from_inventory_at_root(root, inventory, None) } +/// The `(relation, path)` edge tables an adjacency build streams, resolved +/// through the admitted inventory when one is supplied, else the legacy raw +/// `topology/edges/` layout. +pub(crate) fn resolve_adjacency_edge_files( + project_dir: &Path, + inventory: Option<&crate::AuthenticatedPropertyInventory>, +) -> Result, GfError> { + match inventory { + Some(inventory) => Ok(inventory.edge_files(None)), + None => crate::mutator::edge_parquet_files(project_dir, None), + } +} + fn for_each_adjacency_edge_file( project_dir: &Path, inventory: Option<&crate::AuthenticatedPropertyInventory>, batch_size: usize, on_batch: &mut dyn FnMut(&str, bool, &RecordBatch) -> Result<(), GfError>, ) -> Result<(), GfError> { - let paths = match inventory { - Some(inventory) => inventory.edge_files(None), - None => crate::mutator::edge_parquet_files(project_dir, None)?, - }; + let paths = resolve_adjacency_edge_files(project_dir, inventory)?; + for_each_adjacency_edge_path(&paths, batch_size, on_batch) +} + +/// Stream every `(relation, path)` edge table in `paths` as projected batches. +pub(crate) fn for_each_adjacency_edge_path( + paths: &[(String, PathBuf)], + batch_size: usize, + on_batch: &mut dyn FnMut(&str, bool, &RecordBatch) -> Result<(), GfError>, +) -> Result<(), GfError> { for (stem, path) in paths { // An unreadable edge file must FAIL the build, not be skipped: a // manifest written without it would stamp the current generation and // make an index missing a relation's edges look fresh. - let _schema = match crate::catalog::discover_parquet_schema_detailed(&path) { + let _schema = match crate::catalog::discover_parquet_schema_detailed(path) { Ok(schema) => schema, Err(detail) => { return Err(GfError::Storage(format!( @@ -1199,8 +1219,8 @@ fn for_each_adjacency_edge_file( } else { &["edge_id", "src_id", "dst_id"] }; - stream_projected_parquet_batches(&path, columns, batch_size, &mut |batch| { - on_batch(&stem, exploratory, &batch) + stream_projected_parquet_batches(path, columns, batch_size, &mut |batch| { + on_batch(stem, exploratory, &batch) })?; } Ok(()) diff --git a/crates/graphforge-storage/src/adjacency/builder.rs b/crates/graphforge-storage/src/adjacency/builder.rs index 8a4e515c7..3494f6949 100644 --- a/crates/graphforge-storage/src/adjacency/builder.rs +++ b/crates/graphforge-storage/src/adjacency/builder.rs @@ -3,7 +3,7 @@ use super::{ ALL_RELATIONS_STEM, AdjacencyManifestRow, BuildEntry, DEFAULT_ADJACENCY_BATCH_SIZE, DEFAULT_CSR_SHARD_EDGES, DEFAULT_CSR_SHARD_NODES, Direction, ShardedCsrWriter, adjacency_dir, - csr_path, for_each_adjacency_edge_file, named_column, storage_err, string_column, + csr_path, for_each_adjacency_edge_path, named_column, storage_err, string_column, uint64_column, usable_stem, write_manifest, }; use graphforge_core::GfError; @@ -242,6 +242,34 @@ pub fn build_adjacency_index_from_inventory( checkpoint()?; // Generation BEFORE the scan — see the race note in the doc comment. let generation = crate::generation::read_topology_generation(source_project_dir)?; + let edge_files = super::resolve_adjacency_edge_files(source_project_dir, inventory)?; + build_adjacency_index_for_edge_files( + artifact_project_dir, + &edge_files, + generation, + built_at_micros, + options, + checkpoint, + ) +} + +/// Build the index from an explicit `(relation, path)` edge-table list already +/// admitted by the caller, stamping `topology_generation` into the manifest. +/// +/// This is the publish-side entry (#1388): construction encoding and portable +/// import hand over the exact edge tables of the generation they are about to +/// publish, so the CSR ships inside that generation instead of being rebuilt +/// by every query process. The relation names are the semantic routes the +/// caller resolved; no inventory or counter file is consulted here. +pub(crate) fn build_adjacency_index_for_edge_files( + artifact_project_dir: &Path, + edge_files: &[(String, PathBuf)], + topology_generation: u64, + built_at_micros: i64, + options: &AdjacencyBuildOptions, + mut checkpoint: impl FnMut() -> Result<(), GfError>, +) -> Result<(Vec, AdjacencyBuildMetrics), GfError> { + let generation = topology_generation; let options = options.effective(); let adjacency = adjacency_dir(artifact_project_dir); @@ -261,8 +289,7 @@ pub fn build_adjacency_index_from_inventory( let build_result = (|| { let mut groups = stream_build_groups( - source_project_dir, - inventory, + edge_files, &options, &mut spill, &mut metrics, @@ -782,8 +809,7 @@ fn merge_keyed_runs( } fn stream_build_groups( - project_dir: &Path, - inventory: Option<&crate::AuthenticatedPropertyInventory>, + edge_files: &[(String, PathBuf)], options: &AdjacencyBuildOptions, spill: &mut SpillSession, metrics: &mut AdjacencyBuildMetrics, @@ -797,9 +823,8 @@ fn stream_build_groups( EntryGroup::with_label(ALL_RELATIONS_STEM), ); - for_each_adjacency_edge_file( - project_dir, - inventory, + for_each_adjacency_edge_path( + edge_files, options.batch_size, &mut |stem, exploratory, batch| { checkpoint()?; diff --git a/crates/graphforge-storage/src/construction_determinism_tests.rs b/crates/graphforge-storage/src/construction_determinism_tests.rs index 94f45efb6..e0e773b44 100644 --- a/crates/graphforge-storage/src/construction_determinism_tests.rs +++ b/crates/graphforge-storage/src/construction_determinism_tests.rs @@ -366,6 +366,96 @@ mod determinism { ); } + /// Initial construction publishes the derived adjacency CSR as ordinary + /// SHA-256-declared artifacts of the generation (#1388): a query process + /// hydrates and opens it instead of rebuilding it into a temporary + /// directory. Byte reproducibility of those artifacts is covered by + /// `same_input_twice_produces_identical_digests`, which compares every + /// encoded artifact; this test pins presence, generation stamp, and that + /// the published index validates against the published topology. + #[test] + fn initial_construction_publishes_a_current_adjacency_index() { + use crate::adjacency::{Direction, ShardedCsrIndex, csr_path}; + + let nodes = node_ids(256); + let edges = edge_ids(1_024); + let root = TempDir::new().unwrap(); + let mut session = pinned_session(&root, 4); + append_all(&mut session, &nodes, &edges, 128); + session.seal().unwrap(); + let shape = session.shape_canonical_with_cancellation(|| false).unwrap(); + let encoding = session.encode_canonical(&shape, 1).unwrap(); + + let published_index = encoding + .artifacts + .iter() + .map(|artifact| artifact.path.as_str()) + .filter(|path| path.starts_with("indexes/adjacency/")) + .collect::>(); + let shard_manifest = |stem: &str, direction: Direction| { + csr_path(std::path::Path::new(""), stem, direction) + .with_extension("csr.json") + .to_string_lossy() + .into_owned() + }; + let expected_manifests = [ + shard_manifest(crate::adjacency::ALL_RELATIONS_STEM, Direction::Out), + shard_manifest(crate::adjacency::ALL_RELATIONS_STEM, Direction::In), + shard_manifest("R", Direction::Out), + shard_manifest("R", Direction::In), + ]; + for expected in std::iter::once("indexes/adjacency/index_manifest.parquet") + .chain(expected_manifests.iter().map(String::as_str)) + { + assert!( + published_index.contains(&expected), + "{expected} is not among {published_index:?}" + ); + } + assert!(encoding.evidence.adjacency.write_bytes > 0); + assert_eq!(encoding.evidence.adjacency.source_rows, edges.len() as u64); + assert!(encoding.evidence.adjacency.csr_shards >= 4); + + let published = session + .publish_canonical(&encoding, Uuid::from_u128(0x71), Uuid::from_u128(0x72)) + .unwrap(); + let generation = crate::resolve_project_generation(root.path()).unwrap(); + assert_eq!(generation.generation_uuid(), published.generation_uuid); + let inventory = generation.graph_files_inventory().unwrap().unwrap(); + assert_eq!( + inventory + .files + .iter() + .filter(|entry| entry.role == crate::GraphFileRole::Index) + .count(), + published_index.len() + ); + + // Hydrate exactly as a query process does, then open presence-only. + let workspace = TempDir::new().unwrap(); + crate::materialize_graph_objects( + generation.container_root(), + &inventory, + workspace.path(), + ) + .unwrap(); + let rows = crate::adjacency::read_manifest(workspace.path()).unwrap(); + assert!(!rows.is_empty()); + assert!(rows.iter().all(|row| row.topology_generation == 1)); + assert_eq!( + crate::read_topology_generation(workspace.path()).unwrap(), + 1 + ); + let union = ShardedCsrIndex::open(&csr_path(workspace.path(), "_all", Direction::Out)) + .unwrap(); + assert_eq!(union.edge_count(), edges.len() as u64); + assert!( + crate::adjacency::validate_adjacency_index(workspace.path()) + .unwrap() + .is_empty() + ); + } + /// A graph too small to fill one partition shapes into exactly one, and /// that is intended rather than a degenerate case to refuse. /// diff --git a/crates/graphforge-storage/src/graph_construction/encoding_publication/tests.rs b/crates/graphforge-storage/src/graph_construction/encoding_publication/tests.rs index f431288f3..f0ddd2b3a 100644 --- a/crates/graphforge-storage/src/graph_construction/encoding_publication/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/encoding_publication/tests.rs @@ -351,16 +351,13 @@ fn canonical_encoder_outputs_feed_ordinary_readers_index_and_adjacency() { let index = crate::UuidMembershipIndex::open(&graph).unwrap(); assert_eq!(index.count(crate::UuidIndexKind::Node), 3); assert_eq!(index.count(crate::UuidIndexKind::Edge), 2); - let adjacency = crate::adjacency::build_adjacency_index_from_inventory( - &graph, - &graph, - Some(&admitted), - shape.runtime_catalog_now_micros, - &crate::adjacency::AdjacencyBuildOptions::default(), - || Ok(()), - ) - .unwrap(); - assert!(!adjacency.0.is_empty()); + // ADR 0037: the encoder publishes the adjacency CSR with the generation, so + // an ordinary reader reads that index instead of building its own. Rebuilding + // here would replace artifacts the checkpoint's identity ledger has already + // pinned, and the resume below would then refuse them as changed -- which is + // the ledger doing its job, not a defect. + let adjacency = crate::adjacency::read_manifest(&graph).unwrap(); + assert!(!adjacency.is_empty()); drop(session); let mut resumed_session = GraphConstructionSession::open_with_semantic_authority( diff --git a/crates/graphforge-storage/src/graph_construction_encoding.rs b/crates/graphforge-storage/src/graph_construction_encoding.rs index f7ec818f8..ebc69b679 100644 --- a/crates/graphforge-storage/src/graph_construction_encoding.rs +++ b/crates/graphforge-storage/src/graph_construction_encoding.rs @@ -51,6 +51,8 @@ use crate::uuid_membership::{ }; use crate::{SemanticRouteKind, SemanticStorageBindings}; +mod adjacency; + const ENCODED_ROOT: &str = "encoded-v1"; const INVENTORY: &str = "inventory.json"; const ENCODING_INTENT: &str = "encoding-intent.json"; @@ -446,6 +448,9 @@ pub struct GraphConstructionEncodingEvidence { /// v4 receipt, manifest, and authority-lock writes at publication. #[serde(default)] pub ordinal_publication_write_operations: u64, + /// Adjacency CSR artifacts published with the inventory (#1388). + #[serde(default)] + pub adjacency: adjacency::AdjacencyEncodingEvidence, } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] @@ -810,6 +815,15 @@ pub(crate) fn encode( &mut evidence, )?; } + adjacency::encode_adjacency( + &output, + shape, + generation, + &routes, + cancelled, + &mut artifacts, + &mut evidence, + )?; evidence.membership_records = index.input_records; evidence.membership_read_bytes = index.read_bytes; @@ -2947,49 +2961,4 @@ fn hex(bytes: &[u8]) -> String { } #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn runtime_label_decode_rejects_cross_batch_duplicate_identity() { - let mut catalog = graphforge_value::RuntimeCatalogData::new(); - catalog.intern_label_at("First", 1).unwrap(); - catalog.intern_label_at("Second", 1).unwrap(); - let batch = catalog.to_record_batch(); - let mut columns = batch.columns().to_vec(); - columns[2] = Arc::new(UInt32Array::from(vec![0, 0])); - let invalid = RecordBatch::try_new(batch.schema(), columns).unwrap(); - let directory = tempfile::TempDir::new().unwrap(); - let path = directory.path().join("catalog.parquet"); - let mut writer = parquet::arrow::ArrowWriter::try_new( - File::create(&path).unwrap(), - invalid.schema(), - None, - ) - .unwrap(); - writer.write(&invalid).unwrap(); - writer.close().unwrap(); - let budgets = GraphConstructionBudgets { - max_batch_rows: 1, - ..Default::default() - }; - let error = read_runtime_label_ids( - File::open(path).unwrap(), - budgets, - &mut GraphConstructionEncodingEvidence::default(), - ) - .unwrap_err(); - assert!( - error.to_string().contains("unique and contiguous"), - "{error}" - ); - } - - #[test] - fn proof_counter_addition_rejects_overflow_without_clamping() { - let mut counter = u64::MAX; - let error = add_evidence_counter(&mut counter, 1, "mutation").unwrap_err(); - assert!(error.to_string().contains("mutation overflow")); - assert_eq!(counter, u64::MAX); - } -} +mod tests; diff --git a/crates/graphforge-storage/src/graph_construction_encoding/adjacency.rs b/crates/graphforge-storage/src/graph_construction_encoding/adjacency.rs new file mode 100644 index 000000000..cc871ab89 --- /dev/null +++ b/crates/graphforge-storage/src/graph_construction_encoding/adjacency.rs @@ -0,0 +1,190 @@ +//! Publish the derived adjacency CSR with the canonical inventory (#1388). + +use std::ffi::OsStr; +use std::path::{Component, Path}; + +use graphforge_core::GfError; +use serde::{Deserialize, Serialize}; + +use super::{ + ConstructionEncodedArtifact, GraphConstructionEncodingEvidence, StableDirectory, + account_cache_release, add_evidence_counter, authenticate_file_cancellable, directory_for, + storage, +}; +use crate::graph_construction::ConstructionShape; + +/// Measured work of the adjacency build inside canonical encoding. +#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)] +pub struct AdjacencyEncodingEvidence { + /// Exact bytes of the published CSR artifacts. Also counted in the + /// encoding's `output_write_bytes`; the bounded builder's individual write + /// and fsync calls are not attributed. + #[serde(default)] + pub write_bytes: u64, + /// Projected edge rows streamed from the encoded edge tables. + #[serde(default)] + pub source_rows: u64, + /// Sorted spill runs written and merged. + #[serde(default)] + pub spill_runs: u64, + /// Peak bytes charged to the spill session. + #[serde(default)] + pub spill_peak_bytes: u64, + /// CSR shards published across every relation/direction pair. + #[serde(default)] + pub csr_shards: u64, +} + +/// Private spill root for the adjacency build, beside (never inside) the +/// encoded graph tree so a crash cannot leave spill runs among artifacts. +const ADJACENCY_SPILL_ROOT: &str = ".adjacency-spill"; + +/// Publish the derived adjacency CSR with the canonical inventory (#1388). +/// +/// Every query process used to rebuild the whole CSR into a private temporary +/// directory and delete it at exit. Building it once here, from the exact edge +/// tables just encoded, makes it an ordinary SHA-256-declared artifact under +/// `indexes/adjacency/` that the publisher installs into the CAS like every +/// other file: hydration verifies its digest at open, and the persistent +/// adjacency provider then opens it presence-only instead of rebuilding. +/// +/// Only an initial construction (`parent_topology_generation == 0`) is +/// covered: the encoder holds every edge table of that generation. An append +/// carries the parent's files forward structurally without re-reading them, +/// so its index would need the parent's edge tables too; until that lands an +/// append keeps today's lazy rebuild (the provider reads the carried-forward +/// manifest as stale and never serves it). +/// +/// Deterministic by construction: CSR bytes derive from `topology/` alone +/// (R-ADJ-2), the shard directory takes its content digest as its name, and +/// the manifest's build time is the session's recorded clock, so the encoded +/// inventory authority stays reproducible across sessions and resume. +pub(super) fn encode_adjacency( + output: &StableDirectory, + shape: &ConstructionShape, + generation: u64, + routes: &crate::route_component::RouteTable, + cancelled: &mut impl FnMut() -> bool, + artifacts: &mut Vec, + evidence: &mut GraphConstructionEncodingEvidence, +) -> Result<(), GfError> { + if shape.parent_topology_generation != 0 { + return Ok(()); + } + let graph_root = output.path().join("graph"); + let mut edge_files = Vec::new(); + for artifact in artifacts.iter() { + if !artifact.path.starts_with("topology/edges/") { + continue; + } + let semantic = routes.semantic_relative_path(&artifact.path)?; + let relation = crate::route_component::route_position(&semantic)? + .ok_or_else(|| storage("encoded edge table lacks a relation route"))?; + edge_files.push((relation.to_owned(), graph_root.join(&artifact.path))); + } + edge_files.sort(); + + // A crashed earlier attempt may have left a torn index or spill behind; + // the artifact set must be exactly this build's. + let adjacency = crate::adjacency::adjacency_dir(&graph_root); + if adjacency.exists() { + std::fs::remove_dir_all(&adjacency).map_err(storage)?; + } + let spill_root = output.path().join(ADJACENCY_SPILL_ROOT); + if spill_root.exists() { + std::fs::remove_dir_all(&spill_root).map_err(storage)?; + } + let options = crate::adjacency::AdjacencyBuildOptions { + spill_dir: Some(spill_root.clone()), + ..crate::adjacency::AdjacencyBuildOptions::default() + }; + let (rows, metrics) = crate::adjacency::build_adjacency_index_for_edge_files( + &graph_root, + &edge_files, + generation, + shape.runtime_catalog_now_micros, + &options, + || crate::graph_construction::reject_cancelled(cancelled), + )?; + let _ = std::fs::remove_dir_all(&spill_root); + if rows.is_empty() { + return Err(storage("adjacency build published no manifest rows")); + } + + let mut relative_paths = Vec::new(); + collect_relative_files(&graph_root, &adjacency, &mut relative_paths)?; + relative_paths.sort(); + for relative in relative_paths { + let (directory, name) = directory_for(output, &relative)?; + let file = directory + .open_child_file(OsStr::new(&name)) + .map_err(storage)?; + // The CSR is a published artifact like any other, so the allocation + // tracker has to own it. Registering here rather than inside the + // builder keeps the bounded builder's own writes unattributed, which + // is deliberate, while still accounting for the files it leaves + // behind -- `matches_file_inventory` admits no untracked file. + if let Some(allocation) = output.allocation() { + allocation.replace_file_at(&graph_root.join(&relative), &file)?; + } + let (artifact, released, _) = authenticate_file_cancellable(&relative, file, cancelled)?; + add_evidence_counter( + &mut evidence.adjacency.write_bytes, + artifact.bytes, + "adjacency write bytes", + )?; + add_evidence_counter( + &mut evidence.output_write_bytes, + artifact.bytes, + "output write bytes", + )?; + account_cache_release(released, evidence)?; + artifacts.push(artifact); + } + evidence.adjacency.source_rows = metrics.source_rows; + evidence.adjacency.spill_runs = metrics.spill_runs; + evidence.adjacency.spill_peak_bytes = metrics.spill_bytes; + evidence.adjacency.csr_shards = metrics.csr_shards; + Ok(()) +} + +/// Every regular file below `directory`, as normalized paths relative to +/// `root`, refusing links and other non-regular entries. +fn collect_relative_files( + root: &Path, + directory: &Path, + output: &mut Vec, +) -> Result<(), GfError> { + for entry in std::fs::read_dir(directory).map_err(storage)? { + let entry = entry.map_err(storage)?; + let path = entry.path(); + let kind = entry.file_type().map_err(storage)?; + if kind.is_dir() { + collect_relative_files(root, &path, output)?; + continue; + } + if !kind.is_file() { + return Err(storage( + "adjacency artifact tree contains a non-regular file", + )); + } + let relative = path + .strip_prefix(root) + .map_err(|_| storage("adjacency artifact escaped the encoded graph tree"))?; + let mut text = String::new(); + for component in relative.components() { + let Component::Normal(name) = component else { + return Err(storage("adjacency artifact path is not normalized")); + }; + if !text.is_empty() { + text.push('/'); + } + text.push_str( + name.to_str() + .ok_or_else(|| storage("adjacency artifact name is not UTF-8"))?, + ); + } + output.push(text); + } + Ok(()) +} diff --git a/crates/graphforge-storage/src/graph_construction_encoding/tests.rs b/crates/graphforge-storage/src/graph_construction_encoding/tests.rs new file mode 100644 index 000000000..2308dbee0 --- /dev/null +++ b/crates/graphforge-storage/src/graph_construction_encoding/tests.rs @@ -0,0 +1,41 @@ +use super::*; + +#[test] +fn runtime_label_decode_rejects_cross_batch_duplicate_identity() { + let mut catalog = graphforge_value::RuntimeCatalogData::new(); + catalog.intern_label_at("First", 1).unwrap(); + catalog.intern_label_at("Second", 1).unwrap(); + let batch = catalog.to_record_batch(); + let mut columns = batch.columns().to_vec(); + columns[2] = Arc::new(UInt32Array::from(vec![0, 0])); + let invalid = RecordBatch::try_new(batch.schema(), columns).unwrap(); + let directory = tempfile::TempDir::new().unwrap(); + let path = directory.path().join("catalog.parquet"); + let mut writer = + parquet::arrow::ArrowWriter::try_new(File::create(&path).unwrap(), invalid.schema(), None) + .unwrap(); + writer.write(&invalid).unwrap(); + writer.close().unwrap(); + let budgets = GraphConstructionBudgets { + max_batch_rows: 1, + ..Default::default() + }; + let error = read_runtime_label_ids( + File::open(path).unwrap(), + budgets, + &mut GraphConstructionEncodingEvidence::default(), + ) + .unwrap_err(); + assert!( + error.to_string().contains("unique and contiguous"), + "{error}" + ); +} + +#[test] +fn proof_counter_addition_rejects_overflow_without_clamping() { + let mut counter = u64::MAX; + let error = add_evidence_counter(&mut counter, 1, "mutation").unwrap_err(); + assert!(error.to_string().contains("mutation overflow")); + assert_eq!(counter, u64::MAX); +} diff --git a/crates/graphforge-storage/src/project_portable_v2_import.rs b/crates/graphforge-storage/src/project_portable_v2_import.rs index 4469d6abb..014dbfd1a 100644 --- a/crates/graphforge-storage/src/project_portable_v2_import.rs +++ b/crates/graphforge-storage/src/project_portable_v2_import.rs @@ -18,6 +18,8 @@ use crate::project_portable_v2::{ PortableV2Error, PortableV2ErrorCode, PortableV2Limits, PortableV2PackageClass, PortableV2Report, RUNTIME_MAP_PATH, decode_runtime_map, materialize_verified_portable_v2, }; +mod adjacency; + use crate::project_publication::{ ProjectCapability, ProjectFileParticipant, ProjectGenerationRequest, ProjectParticipant, ProjectParticipantEncoding, ProjectPublicationReceipt, ProjectStageOutcome, @@ -288,6 +290,9 @@ pub fn import_complete_portable_v2_with_allocation( package_digest: Some(report.package_digest.clone()), }); let allocation_on_error = materialized_identity_allocated_bytes.clone(); + // Files the derived adjacency build adds to the stage (#1388); the stage + // identity walks below are bounded by the entry count and must include them. + let mut added_stage_entries = 0_usize; let result = import_materialized( &stage, target, @@ -299,6 +304,7 @@ pub fn import_complete_portable_v2_with_allocation( &report, owned_retry, allocation, + &mut added_stage_entries, ) .map(|mut receipt| { receipt.materialized_identity_allocated_bytes = materialized_identity_allocated_bytes; @@ -311,7 +317,7 @@ pub fn import_complete_portable_v2_with_allocation( &owner, materialized_stage_identity, &mut owned_identities, - entry_count, + entry_count.saturating_add(added_stage_entries), allocation, ) { return cleanup_error.with_allocation_identities(owned_identities); @@ -331,7 +337,7 @@ pub fn import_complete_portable_v2_with_allocation( &owner, materialized_stage_identity, &mut receipt, - entry_count, + entry_count.saturating_add(added_stage_entries), allocation, )?; Ok(receipt) @@ -862,6 +868,7 @@ fn import_materialized( report: &PortableV2Report, owned_retry: bool, allocation: Option<&crate::StorageAllocationOperation>, + added_stage_entries: &mut usize, ) -> Result { if report.package_class != PortableV2PackageClass::Complete { return Err(PortableV2Error::new( @@ -1009,7 +1016,17 @@ fn import_materialized( if let Some(graph_tree) = &package_graph_tree { validate_import_graph_identities(graph_tree, cancelled)?; validate_import_property_values(graph_tree, report.entry_count, cancelled)?; + *added_stage_entries = + adjacency::persist_import_adjacency(stage, graph_tree, &participants, cancelled)?; } + let stage_entry_count = usize::try_from(report.entry_count) + .map_err(|_| { + PortableV2Error::new( + PortableV2ErrorCode::LimitExceeded, + "import entry count exceeds platform capacity", + ) + })? + .saturating_add(*added_stage_entries); let admission = crate::filesystem_admission::admit_project_lifecycle( target, crate::filesystem_admission::ProjectLifecycleMode::Durable, @@ -1049,12 +1066,7 @@ fn import_materialized( admission.root(), package_graph_tree.as_deref(), &mut participants, - usize::try_from(report.entry_count).map_err(|_| { - PortableV2Error::new( - PortableV2ErrorCode::LimitExceeded, - "import entry count exceeds platform capacity", - ) - })?, + stage_entry_count, allocation, )?; let request = ProjectGenerationRequest { @@ -1966,7 +1978,7 @@ mod tests { assert!(receipt.parent_sync_confirmed); } - fn supported() -> Vec { + pub(super) fn supported() -> Vec { vec![ ProjectCapability { capability_id: "graph".into(), diff --git a/crates/graphforge-storage/src/project_portable_v2_import/adjacency.rs b/crates/graphforge-storage/src/project_portable_v2_import/adjacency.rs new file mode 100644 index 000000000..5a1d0a165 --- /dev/null +++ b/crates/graphforge-storage/src/project_portable_v2_import/adjacency.rs @@ -0,0 +1,375 @@ +//! Import-time publication of the derived adjacency CSR (#1388). + +use std::fs; +use std::path::{Path, PathBuf}; +use std::sync::atomic::AtomicBool; + +use graphforge_core::GfError; + +use super::{ + PortableV2Error, PortableV2ErrorCode, ProjectFileParticipant, storage, storage_or_cancel, +}; + +/// Private spill root for the import-time adjacency build, inside the owned +/// stage so failure cleanup reclaims it with everything else. +const IMPORT_ADJACENCY_SPILL_ROOT: &str = ".adjacency-spill"; + +/// Build the derived adjacency CSR into the verified package graph tree when +/// the package carries none, so the imported generation ships it (#1388) +/// exactly as a construction-published generation does. A complete package of +/// such a generation already carries the index and is installed unchanged; +/// subset exports and packages made before #1388 lack it, and a clean import +/// of those would otherwise rebuild the whole CSR in every query process. The +/// tree is still private here: the compact CAS append that follows hashes and +/// installs the new files like any other package payload. +/// +/// Only the compact (v2 / mapped-root) participant contract is covered. A v1 +/// inventory participant is verified file-for-file against the package tree +/// at staging, so derived files cannot be added there; that path keeps the +/// lazy rebuild. +/// +/// Returns the number of files added to the stage. +pub(super) fn persist_import_adjacency( + stage: &Path, + graph_tree: &Path, + participants: &[ProjectFileParticipant], + cancelled: Option<&AtomicBool>, +) -> Result { + let Some(participant) = participants.iter().find(|participant| { + participant.participant.capability_id == crate::GRAPH_CAPABILITY_ID + && participant.participant.record_family_id == crate::GRAPH_FILES_FAMILY + }) else { + return Ok(0); + }; + if !matches!( + participant.participant.record_version, + crate::GRAPH_FILES_V2_RECORD_VERSION + | crate::graph_files::GRAPH_FILES_MAPPED_ROOT_RECORD_VERSION + ) { + return Ok(0); + } + let adjacency = crate::adjacency::adjacency_dir(graph_tree); + if adjacency.exists() { + return Ok(0); + } + let generation = + crate::generation::read_topology_generation(graph_tree).map_err(|error| storage(&error))?; + let edge_files = import_edge_files(graph_tree)?; + let spill_root = stage.join(IMPORT_ADJACENCY_SPILL_ROOT); + let options = crate::adjacency::AdjacencyBuildOptions { + spill_dir: Some(spill_root.clone()), + ..crate::adjacency::AdjacencyBuildOptions::default() + }; + let built_at_micros = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map_or(0, |elapsed| { + i64::try_from(elapsed.as_micros()).unwrap_or(i64::MAX) + }); + let outcome = crate::adjacency::build_adjacency_index_for_edge_files( + graph_tree, + &edge_files, + generation, + built_at_micros, + &options, + || { + if cancelled.is_some_and(|flag| flag.load(std::sync::atomic::Ordering::Relaxed)) { + Err(GfError::Storage("import cancelled".into())) + } else { + Ok(()) + } + }, + ); + let _ = fs::remove_dir_all(&spill_root); + outcome.map_err(|error| storage_or_cancel(&error, cancelled))?; + let mut added = 0_usize; + let mut pending = vec![adjacency]; + while let Some(directory) = pending.pop() { + for entry in fs::read_dir(&directory).map_err(|_| { + PortableV2Error::new(PortableV2ErrorCode::Io, "cannot inventory staged adjacency") + })? { + let entry = entry.map_err(|_| { + PortableV2Error::new(PortableV2ErrorCode::Io, "cannot inventory staged adjacency") + })?; + if entry.path().is_dir() { + pending.push(entry.path()); + } else { + added = added.saturating_add(1); + } + } + } + Ok(added) +} + +/// The `(relation, path)` edge tables of a verified package graph tree, +/// resolved through its owned route authority when it carries one. +fn import_edge_files(graph_tree: &Path) -> Result, PortableV2Error> { + let directory = graphforge_filesystem::StableDirectory::open(graph_tree).map_err(|_| { + PortableV2Error::new( + PortableV2ErrorCode::Io, + "cannot authenticate portable graph tree", + ) + })?; + let routes = crate::route_component::owned::read_owned_layout_table(&directory) + .map_err(|error| storage(&error))?; + let edges = graph_tree.join("topology").join("edges"); + let mut files = Vec::new(); + if !edges.is_dir() { + return Ok(files); + } + let mut pending = vec![edges]; + while let Some(current) = pending.pop() { + for entry in fs::read_dir(¤t).map_err(|_| { + PortableV2Error::new(PortableV2ErrorCode::Io, "cannot read portable edge tables") + })? { + let path = entry + .map_err(|_| { + PortableV2Error::new( + PortableV2ErrorCode::Io, + "cannot read portable edge tables", + ) + })? + .path(); + if path.is_dir() { + pending.push(path); + continue; + } + if path.extension().and_then(|extension| extension.to_str()) != Some("parquet") { + continue; + } + let relative = path + .strip_prefix(graph_tree) + .map_err(|_| { + PortableV2Error::new( + PortableV2ErrorCode::InvalidStructure, + "portable edge table escaped the graph tree", + ) + })? + .components() + .map(|component| match component { + std::path::Component::Normal(value) => value.to_str().ok_or_else(|| { + PortableV2Error::new( + PortableV2ErrorCode::InvalidStructure, + "portable edge table name is not UTF-8", + ) + }), + _ => Err(PortableV2Error::new( + PortableV2ErrorCode::InvalidStructure, + "portable edge table path is not normalized", + )), + }) + .collect::, _>>()? + .join("/"); + let semantic = match &routes { + Some(table) => table.semantic_relative_path(&relative), + None => crate::graph_files::legacy_inventory_logical_text(&relative), + } + .map_err(|error| storage(&error))?; + let Some(route) = crate::route_component::route_position(&semantic) + .map_err(|error| storage(&error))? + else { + continue; + }; + files.push((route.to_owned(), path)); + } + } + files.sort(); + Ok(files) +} + +#[cfg(test)] +mod tests { + use std::fs; + use std::path::Path; + use std::sync::atomic::AtomicBool; + + use arrow::array::{ArrayRef, FixedSizeBinaryArray, StringArray}; + use arrow::record_batch::RecordBatch; + use sha2::{Digest, Sha256}; + use std::sync::Arc; + use uuid::Uuid; + + use super::super::tests::supported; + use super::*; + use crate::{PortableV2Limits, import_complete_portable_v2}; + + fn fixed(ids: &[Uuid]) -> FixedSizeBinaryArray { + FixedSizeBinaryArray::try_from_iter(ids.iter().map(|id| id.as_bytes().as_slice())).unwrap() + } + + /// An initial construction of 8 nodes and 24 `KNOWS` edges, published as + /// the project's current generation. + fn construction_published_project() -> tempfile::TempDir { + let node_ids = (0..8_u128) + .map(|index| Uuid::from_u128(0x1000 + index)) + .collect::>(); + let edge_ids = (0..24_u128) + .map(|index| Uuid::from_u128(0x2000 + index)) + .collect::>(); + let nodes = RecordBatch::try_new( + crate::CONSTRUCTION_NODE_SCHEMA.clone(), + vec![ + Arc::new(fixed(&node_ids)) as ArrayRef, + Arc::new(StringArray::from(vec!["Person"; node_ids.len()])) as ArrayRef, + ], + ) + .unwrap(); + let sources = (0..edge_ids.len()) + .map(|index| node_ids[index % node_ids.len()]) + .collect::>(); + let targets = (0..edge_ids.len()) + .map(|index| node_ids[(index + 3) % node_ids.len()]) + .collect::>(); + let edges = RecordBatch::try_new( + crate::CONSTRUCTION_EDGE_SCHEMA.clone(), + vec![ + Arc::new(fixed(&edge_ids)) as ArrayRef, + Arc::new(StringArray::from(vec!["KNOWS"; edge_ids.len()])) as ArrayRef, + Arc::new(fixed(&sources)) as ArrayRef, + Arc::new(fixed(&targets)) as ArrayRef, + ], + ) + .unwrap(); + let source = tempfile::tempdir().unwrap(); + crate::open_or_initialize_project(source.path()).unwrap(); + let mut session = crate::GraphConstructionSession::open_with_mode( + source.path(), + Uuid::from_u128(0x31), + 0, + graphforge_core::OntologyMode::Exploratory, + crate::GraphConstructionBudgets::default(), + ) + .unwrap(); + session + .append(crate::ConstructionChunkKind::Node, "nodes", &nodes) + .unwrap(); + session + .append(crate::ConstructionChunkKind::Edge, "edges", &edges) + .unwrap(); + session.seal().unwrap(); + let shape = session.shape_canonical_with_cancellation(|| false).unwrap(); + let encoding = session.encode_canonical(&shape, 1).unwrap(); + session + .publish_canonical(&encoding, Uuid::from_u128(0x32), Uuid::from_u128(0x33)) + .unwrap(); + source + } + + fn index_entries(inventory: &crate::GraphFilesInventory) -> Vec { + inventory + .files + .iter() + .filter(|entry| entry.role == crate::GraphFileRole::Index) + .map(|entry| entry.relative_path.clone()) + .collect() + } + + /// The materialized tree at `root` carries a manifest stamped with its own + /// topology generation and a CSR that validates against its topology. + fn assert_index_current(root: &Path) { + let topology_generation = crate::read_topology_generation(root).unwrap(); + let rows = crate::adjacency::read_manifest(root).unwrap(); + assert!(!rows.is_empty()); + assert!( + rows.iter() + .all(|row| row.topology_generation == topology_generation) + ); + assert!( + crate::adjacency::validate_adjacency_index(root) + .unwrap() + .is_empty() + ); + } + + fn export_complete(source: &Path) -> (tempfile::TempDir, std::path::PathBuf) { + let generation = crate::resolve_project_generation(source).unwrap(); + let package_parent = tempfile::tempdir().unwrap(); + let package = package_parent.path().join("adjacency.gfproject"); + let limits = crate::PortableV2ExportLimits::default(); + let plan = crate::plan_complete_portable_v2(&generation, limits).unwrap(); + crate::export_complete_portable_v2( + &plan, + &package, + crate::PortableV2Output::Expanded, + limits, + &AtomicBool::new(false), + |_| {}, + ) + .unwrap(); + (package_parent, package) + } + + /// A complete package of a construction-published generation carries its + /// adjacency CSR, and the import must publish it unchanged rather than + /// building a second one (#1388). + #[test] + fn complete_import_publishes_the_derived_adjacency_index() { + let source = construction_published_project(); + let generation = crate::resolve_project_generation(source.path()).unwrap(); + let source_index = index_entries(&generation.graph_files_inventory().unwrap().unwrap()); + assert!(source_index.contains(&"indexes/adjacency/index_manifest.parquet".to_owned())); + + let (_package_parent, package) = export_complete(source.path()); + assert!( + package + .join("data/components/graph-data/graph-tree/indexes/adjacency") + .is_dir(), + "a complete package carries the generation's derived index" + ); + + let target = tempfile::tempdir().unwrap(); + import_complete_portable_v2( + &package, + target.path(), + Uuid::from_u128(0x34), + Uuid::from_u128(0x35), + &supported(), + PortableV2Limits::default(), + None, + ) + .unwrap(); + let imported = crate::resolve_project_generation(target.path()).unwrap(); + let inventory = imported.graph_files_inventory().unwrap().unwrap(); + assert_eq!(index_entries(&inventory), source_index); + let workspace = tempfile::tempdir().unwrap(); + crate::materialize_graph_objects(imported.container_root(), &inventory, workspace.path()) + .unwrap(); + assert_index_current(workspace.path()); + } + + /// A package without the index (a subset export, or one made before + /// #1388) gets it built into the verified stage tree before the compact + /// CAS append, so the imported generation still ships it. + #[test] + fn import_builds_the_index_a_package_lacks_and_never_twice() { + let source = construction_published_project(); + let generation = crate::resolve_project_generation(source.path()).unwrap(); + let inventory = generation.graph_files_inventory().unwrap().unwrap(); + let stage = tempfile::tempdir().unwrap(); + let tree = stage.path().join("graph-tree"); + fs::create_dir(&tree).unwrap(); + crate::materialize_graph_objects(generation.container_root(), &inventory, &tree).unwrap(); + fs::remove_dir_all(tree.join("indexes")).unwrap(); + let placeholder = crate::graph_files_root_participant(&crate::GraphFilesRootV2 { + format: "graphforge-graph-files-root".into(), + format_version: crate::graph_files::GRAPH_FILES_MAPPED_ROOT_RECORD_VERSION, + root_node_sha256: "0".repeat(64), + logical_file_count: 0, + logical_byte_length: 0, + }) + .unwrap(); + let participants = vec![ProjectFileParticipant { + participant: placeholder.clone(), + source: stage.path().join("graph-files.json"), + byte_length: placeholder.bytes.len() as u64, + content_sha256: Sha256::digest(&placeholder.bytes).into(), + }]; + let added = persist_import_adjacency(stage.path(), &tree, &participants, None).unwrap(); + assert_eq!(added, index_entries(&inventory).len()); + assert!(!stage.path().join(IMPORT_ADJACENCY_SPILL_ROOT).exists()); + assert_index_current(&tree); + assert_eq!( + persist_import_adjacency(stage.path(), &tree, &participants, None).unwrap(), + 0 + ); + } +} diff --git a/docs-site/astro.config.mjs b/docs-site/astro.config.mjs index 7ff1c1f71..c701e269a 100644 --- a/docs-site/astro.config.mjs +++ b/docs-site/astro.config.mjs @@ -359,6 +359,10 @@ export default defineConfig({ label: '0036 — The GraphForge release version contract', slug: 'adr/0036-release-version-contract', }, + { + label: '0037 — Derived adjacency is published with the generation', + slug: 'adr/0037-adjacency-published-with-generation', + }, { label: '0038 — Determinism belongs at the publication boundary', slug: 'adr/0038-determinism-at-the-publication-boundary', diff --git a/docs-site/scripts/sync-content.mjs b/docs-site/scripts/sync-content.mjs index ea730c65a..2c7402aea 100755 --- a/docs-site/scripts/sync-content.mjs +++ b/docs-site/scripts/sync-content.mjs @@ -157,6 +157,7 @@ const PAGES = [ 'adr/0032-research-project-authority.md', 'adr/0035-structured-stage-errors.md', 'adr/0036-release-version-contract.md', + 'adr/0037-adjacency-published-with-generation.md', 'adr/0038-determinism-at-the-publication-boundary.md', // END generated ADR records 'releases/roadmap.md', diff --git a/docs/adr/0037-adjacency-published-with-generation.md b/docs/adr/0037-adjacency-published-with-generation.md new file mode 100644 index 000000000..762c44be6 --- /dev/null +++ b/docs/adr/0037-adjacency-published-with-generation.md @@ -0,0 +1,81 @@ +--- +title: "ADR 0037: Derived adjacency is published with the generation" +adr: "0037" +status: "Accepted" +date: "2026-09-17" +superseded_by: null +--- + +# ADR 0037: Derived adjacency is published with the generation + +**Status:** Accepted + +**Date:** 2026-09-17 + +**Build target:** v0.6.0 + +**Implementation:** Accepted for the #1388 implementation. Initial +construction and complete portable import publish the index; appends keep the +lazy rebuild and are the recorded follow-up. + +**Related:** #1388, [ADR 0004](0004-adjacency-index.md), [ADR 0013](0013-project-generation-protocol.md) + +## Context + +ADR 0004 made the adjacency CSR a derived, rebuildable artifact under +`indexes/adjacency/`. The persistent provider serves it when its manifest +matches the project's `topology_generation` and rebuilds it lazily otherwise. +Host-scale ingest never created it: every query process rebuilt the entire CSR +into a per-process temporary directory and deleted it at exit. At S20 that is +about 16 s of CPU per process, paid four times per ladder rung; at S26 it was +0.58 h, the largest non-ingest term. The rebuild bypasses the read-path +counters, so the `adjacency` storage category reported zero bytes while a +one-hop query read gigabytes. + +## Decision + +The generation that publishes canonical topology also publishes its derived +adjacency CSR, as ordinary graph-files inventory entries. + +- **Construction.** The canonical encoder builds `indexes/adjacency/` from the + exact edge tables it has just encoded, stamps the generation it is about to + bind, and records every file as a SHA-256-declared artifact of the encoded + inventory. The existing publisher installs those artifacts into the project + object store like every other file. Only an initial construction + (`parent_topology_generation == 0`) is covered: an append carries the parent's + files forward without re-reading them, so its index would need the parent's + edge tables too. Until that lands, an append's carried-forward manifest reads + as stale and is never served; the provider rebuilds lazily as before. +- **Portable import.** A complete package of a generation carries its index + and imports it unchanged. Subset exports exclude `Index`-role files and older + packages predate the index; for those a complete import builds the CSR into + the verified package tree before the compact object-store append, so the + imported generation ships it. The v1 inventory contract, which is verified + file-for-file against the package tree, keeps the lazy rebuild. +- **Reading.** No read-path change. Hydration materializes the entries into + the workspace, the provider finds a fresh manifest and opens the CSR + presence-only (#1094), and a project without one rebuilds exactly as before. + +## Consequences + +- **Not a durable-format change.** The on-disk CSR format, the manifest, and + `Index`-role inventory entries all pre-date this record; explicit + `index("adjacency")` already published them. What changes is when the + artifact is produced. No version bump and no compatibility read path. +- **Determinism is preserved.** CSR bytes derive from `topology/` alone, the + shard directory takes its content digest as its name, and the manifest's + build time is the session's recorded clock. The encoded inventory authority + therefore stays reproducible across sessions and resume; the determinism + suite covers the new artifacts because it compares every encoded artifact. +- **Authentication.** Every published CSR object is digest-verified by the + open-time sweep like Topology, and shard payloads are additionally + authenticated on first row touch. The shard manifest (`*.csr.json`) and + `index_manifest.parquet` are covered by the sweep only; the provider treats + any disagreement as a stale index and rebuilds privately rather than serving. +- **Cost moves to publish.** One build per generation instead of one per + process, plus the artifact bytes in the object store and in every open-time + sweep. The union `_all` pair duplicates the per-relation pair for a + single-relation graph; sharing them is a possible later saving. +- **Follow-ups.** Append constructions; tombstoning a parent's stale index on + append; attributing the rebuild's I/O to the read-path counters, which #1422 + currently misses. diff --git a/docs/adr/README.md b/docs/adr/README.md index 59022506c..1a8a99dd0 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -42,6 +42,7 @@ Roadmap-only ADRs are not retained in this tree. | 0032 | [Research Branches share Project publication authority](0032-research-project-authority.md) | `0032-research-project-authority.md` | | 0035 | [Preserve stage diagnostics at public error boundaries](0035-structured-stage-errors.md) | `0035-structured-stage-errors.md` | | 0036 | [The GraphForge release version contract](0036-release-version-contract.md) | `0036-release-version-contract.md` | +| 0037 | [Derived adjacency is published with the generation](0037-adjacency-published-with-generation.md) | `0037-adjacency-published-with-generation.md` | | 0038 | [Determinism belongs at the publication boundary](0038-determinism-at-the-publication-boundary.md) | `0038-determinism-at-the-publication-boundary.md` | ## Superseded records diff --git a/docs/engineering/adrs/README.md b/docs/engineering/adrs/README.md index dd47fbde3..311f641e3 100644 --- a/docs/engineering/adrs/README.md +++ b/docs/engineering/adrs/README.md @@ -83,6 +83,7 @@ the Repository Policy job (#1390). | 0032 | Research Branches share Project publication authority | Accepted | [`../../adr/0032-research-project-authority.md`](../../adr/0032-research-project-authority.md) | | 0035 | Preserve stage diagnostics at public error boundaries | Accepted | [`../../adr/0035-structured-stage-errors.md`](../../adr/0035-structured-stage-errors.md) | | 0036 | The GraphForge release version contract | Accepted | [`../../adr/0036-release-version-contract.md`](../../adr/0036-release-version-contract.md) | +| 0037 | Derived adjacency is published with the generation | Accepted | [`../../adr/0037-adjacency-published-with-generation.md`](../../adr/0037-adjacency-published-with-generation.md) | | 0038 | Determinism belongs at the publication boundary | Accepted | [`../../adr/0038-determinism-at-the-publication-boundary.md`](../../adr/0038-determinism-at-the-publication-boundary.md) | ### Superseded