From df5721be9a930b545f2cf923a5a22ca698410f70 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 07:44:06 -0600 Subject: [PATCH 1/7] feat(storage): add versioned sharded CSR format --- crates/graphforge-storage/src/adjacency.rs | 261 +++++++++++++++++++++ 1 file changed, 261 insertions(+) diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index a7ad654c8..2f46729c6 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -44,6 +44,7 @@ use arrow::buffer::{OffsetBuffer, ScalarBuffer}; use arrow::datatypes::{DataType, Field}; use arrow::ipc::reader::FileReader; use arrow::ipc::writer::FileWriter; +use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use tempfile::NamedTempFile; @@ -60,6 +61,212 @@ pub const ALL_RELATIONS_STEM: &str = "_all"; /// File name of the adjacency index manifest within `indexes/adjacency/`. pub const MANIFEST_FILE: &str = "index_manifest.parquet"; +const SHARDED_CSR_VERSION: u32 = 1; + +#[derive(Clone, Debug, Serialize, Deserialize)] +struct CsrShardRecord { + first_node: u64, + node_count: u64, + edge_count: u64, + file: String, + sha256: String, +} + +#[derive(Clone, Debug, Serialize, Deserialize)] +struct CsrShardManifest { + format: String, + version: u32, + node_count: u64, + edge_count: u64, + shard_dir: String, + shards: Vec, +} + +/// Bounded reader for a versioned sharded CSR. Opening validates the small +/// manifest; row access reads and authenticates only the containing shard. +#[derive(Clone, Debug)] +pub struct ShardedCsrIndex { + root: PathBuf, + manifest: CsrShardManifest, +} + +impl ShardedCsrIndex { + /// Open the shard manifest beside the legacy logical `.csr` path. + pub fn open(path: &Path) -> Result { + let manifest_path = path.with_extension("csr.json"); + let bytes = std::fs::read(&manifest_path).map_err(storage_err)?; + let manifest: CsrShardManifest = serde_json::from_slice(&bytes).map_err(storage_err)?; + if manifest.format != "graphforge.csr-shards" || manifest.version != SHARDED_CSR_VERSION { + return Err(GfError::Storage(format!( + "unsupported sharded CSR manifest {}", + manifest_path.display() + ))); + } + let mut expected = 0_u64; + let mut edges = 0_u64; + for shard in &manifest.shards { + if shard.first_node != expected || shard.node_count == 0 { + return Err(GfError::Storage("invalid CSR shard boundary chain".into())); + } + expected = expected.saturating_add(shard.node_count); + edges = edges.saturating_add(shard.edge_count); + if Path::new(&shard.file).components().count() != 1 { + return Err(GfError::Storage("invalid CSR shard file name".into())); + } + } + if expected != manifest.node_count || edges != manifest.edge_count { + return Err(GfError::Storage( + "CSR shard manifest counts disagree".into(), + )); + } + Ok(Self { + root: path + .parent() + .unwrap_or_else(|| Path::new(".")) + .join(&manifest.shard_dir), + manifest, + }) + } + + /// Total logical source rows across all shards. + #[must_use] + pub const fn node_count(&self) -> u64 { + self.manifest.node_count + } + + /// Total adjacency entries across all shards. + #[must_use] + pub const fn edge_count(&self) -> u64 { + self.manifest.edge_count + } + + /// Read one logical row without loading unrelated shards. + pub fn row(&self, node_id: u64) -> Result, GfError> { + if node_id >= self.manifest.node_count { + return Ok(Vec::new()); + } + let shard = self + .manifest + .shards + .binary_search_by(|candidate| { + if node_id < candidate.first_node { + std::cmp::Ordering::Greater + } else if node_id >= candidate.first_node.saturating_add(candidate.node_count) { + std::cmp::Ordering::Less + } else { + std::cmp::Ordering::Equal + } + }) + .map_err(|_| GfError::Storage("CSR shard boundary lookup failed".into()))?; + let record = &self.manifest.shards[shard]; + let path = self.root.join(&record.file); + let bytes = std::fs::read(&path).map_err(|error| { + GfError::Storage(format!("missing CSR shard {}: {error}", path.display())) + })?; + if sha256_hex(&bytes) != record.sha256 { + return Err(GfError::Storage(format!( + "CSR shard checksum mismatch: {}", + record.file + ))); + } + let csr = read_csr_bytes(&bytes, &path)?; + if csr.node_count() != record.node_count || csr.edge_count() != record.edge_count { + return Err(GfError::Storage(format!( + "CSR shard count mismatch: {}", + record.file + ))); + } + Ok(csr.row(node_id - record.first_node).iter().collect()) + } +} + +/// Write a versioned checksummed shard set and publish its manifest last. +/// This compatibility helper accepts an in-memory CSR; streaming builders use +/// the same shard writer while producing one bounded shard at a time. +pub fn write_sharded_csr(path: &Path, csr: &CsrIndex, max_edges: usize) -> Result<(), GfError> { + csr.validate()?; + let parent = path + .parent() + .ok_or_else(|| GfError::Storage("CSR path has no parent".into()))?; + std::fs::create_dir_all(parent).map_err(storage_err)?; + let stem = path + .file_name() + .and_then(|name| name.to_str()) + .unwrap_or("csr"); + let shard_dir = format!("{stem}.{}.d", uuid::Uuid::new_v4().as_simple()); + let root = parent.join(&shard_dir); + std::fs::create_dir(&root).map_err(storage_err)?; + let result = (|| { + let mut records = Vec::new(); + let mut first_node = 0_u64; + let mut shard = CsrIndex { + offsets: vec![0], + ..CsrIndex::default() + }; + for node in 0..csr.node_count() { + let row = csr.row(node); + if shard.node_count() > 0 + && shard.edge_ids.len().saturating_add(row.len()) > max_edges.max(1) + { + records.push(write_csr_shard(&root, first_node, &shard, records.len())?); + first_node = node; + shard = CsrIndex { + offsets: vec![0], + ..CsrIndex::default() + }; + } + shard.edge_ids.extend_from_slice(row.edge_ids); + shard.neighbor_ids.extend_from_slice(row.neighbor_ids); + shard.offsets.push(shard.edge_count()); + } + if shard.node_count() > 0 || records.is_empty() { + records.push(write_csr_shard(&root, first_node, &shard, records.len())?); + } + let manifest = CsrShardManifest { + format: "graphforge.csr-shards".into(), + version: SHARDED_CSR_VERSION, + node_count: csr.node_count(), + edge_count: csr.edge_count(), + shard_dir, + shards: records, + }; + let bytes = serde_json::to_vec_pretty(&manifest).map_err(storage_err)?; + let manifest_path = path.with_extension("csr.json"); + let mut temp = tempfile::Builder::new() + .prefix(stem) + .suffix(".json.tmp") + .tempfile_in(parent) + .map_err(storage_err)?; + use std::io::Write as _; + temp.write_all(&bytes).map_err(storage_err)?; + temp.as_file().sync_all().map_err(storage_err)?; + persist_temp(temp, &manifest_path) + })(); + if result.is_err() { + let _ = std::fs::remove_dir_all(&root); + } + result +} + +fn write_csr_shard( + root: &Path, + first_node: u64, + shard: &CsrIndex, + ordinal: usize, +) -> Result { + let file = format!("{ordinal:020}.csr"); + let path = root.join(&file); + write_csr(&path, shard)?; + let bytes = std::fs::read(&path).map_err(storage_err)?; + Ok(CsrShardRecord { + first_node, + node_count: shard.node_count(), + edge_count: shard.edge_count(), + file, + sha256: sha256_hex(&bytes), + }) +} + fn storage_err(e: impl std::fmt::Display) -> GfError { GfError::Storage(e.to_string()) } @@ -425,6 +632,20 @@ pub fn read_csr(path: &Path) -> Result { .map_err(|e| GfError::Storage(format!("cannot open CSR file {}: {e}", path.display())))?; let reader = FileReader::try_new(file, None) .map_err(|e| GfError::Storage(format!("invalid CSR file {}: {e}", path.display())))?; + decode_csr(reader, path) +} + +fn read_csr_bytes(bytes: &[u8], path: &Path) -> Result { + let reader = FileReader::try_new(std::io::Cursor::new(bytes), None).map_err(|error| { + GfError::Storage(format!("invalid CSR shard {}: {error}", path.display())) + })?; + decode_csr(reader, path) +} + +fn decode_csr( + reader: FileReader, + path: &Path, +) -> Result { if reader.schema().fields() != ADJACENCY_CSR_SCHEMA.fields() { return Err(GfError::Storage(format!( "CSR file {} has unexpected schema {:?}", @@ -481,6 +702,16 @@ pub fn read_csr(path: &Path) -> Result { Ok(csr) } +fn sha256_hex(bytes: &[u8]) -> String { + use std::fmt::Write as _; + Sha256::digest(bytes) + .iter() + .fold(String::with_capacity(64), |mut output, byte| { + let _ = write!(output, "{byte:02x}"); + output + }) +} + /// Replace `index_manifest.parquet` with `rows`, atomically. /// /// Per the build ordering convention (module docs), call this **after** all @@ -1823,6 +2054,36 @@ mod tests { } } + #[test] + fn sharded_csr_crosses_boundaries_and_rejects_missing_or_corrupt_shards() { + let directory = TempDir::new().unwrap(); + let path = directory.path().join("KNOWS.out.csr"); + let expected = sample_csr(); + write_sharded_csr(&path, &expected, 2).unwrap(); + let reader = ShardedCsrIndex::open(&path).unwrap(); + assert_eq!((reader.node_count(), reader.edge_count()), (3, 4)); + for node in 0..expected.node_count() { + assert_eq!( + reader.row(node).unwrap(), + expected.row(node).iter().collect::>() + ); + } + + let first = reader.root.join(&reader.manifest.shards[0].file); + let original = std::fs::read(&first).unwrap(); + std::fs::write(&first, b"corrupt").unwrap(); + assert!(reader.row(0).unwrap_err().to_string().contains("checksum")); + std::fs::write(&first, original).unwrap(); + std::fs::remove_file(&first).unwrap(); + assert!( + reader + .row(0) + .unwrap_err() + .to_string() + .contains("missing CSR shard") + ); + } + #[test] fn entries_from_out_csr_rejects_malformed_offset_bounds() { let past_targets = CsrIndex { From d645ba2a66476769c384071bc023330f40a6274c Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 07:45:15 -0600 Subject: [PATCH 2/7] refactor(storage): expose streamed CSR merge sink --- crates/graphforge-storage/src/adjacency.rs | 72 ++++++++++------------ 1 file changed, 31 insertions(+), 41 deletions(-) diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index 2f46729c6..77f8b47ca 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -1310,15 +1310,24 @@ fn merge_keyed_runs_to_csr( runs: &[PathBuf], checkpoint: &mut dyn FnMut() -> Result<(), GfError>, ) -> Result { + let mut keyed = Vec::new(); + merge_keyed_runs(runs, checkpoint, &mut |entry| { + keyed.push(entry); + Ok(()) + })?; + Ok(csr_from_sorted_keyed(&keyed)) +} + +fn merge_keyed_runs( + runs: &[PathBuf], + checkpoint: &mut dyn FnMut() -> Result<(), GfError>, + emit: &mut dyn FnMut((u64, u64, u64)) -> Result<(), GfError>, +) -> Result<(), GfError> { use std::cmp::Reverse; use std::collections::BinaryHeap; if runs.is_empty() { - return Ok(CsrIndex { - offsets: vec![0], - edge_ids: Vec::new(), - neighbor_ids: Vec::new(), - }); + return Ok(()); } let mut cursors: Vec = runs @@ -1333,12 +1342,6 @@ fn merge_keyed_runs_to_csr( } } - let mut csr = CsrIndex { - offsets: vec![0], - edge_ids: Vec::new(), - neighbor_ids: Vec::new(), - }; - let mut current_node: Option = None; let mut seen = 0u64; while let Some(Reverse((key, edge, neighbor, idx))) = heap.pop() { if seen.is_multiple_of(65_536) { @@ -1346,43 +1349,30 @@ fn merge_keyed_runs_to_csr( } seen += 1; - match current_node { - None => { - for _ in 0..key { - csr.offsets.push(csr.edge_count()); - } - current_node = Some(key); - } - Some(node) if key > node => { - // Close `node`, then emit empty rows for (node+1)..key. - csr.offsets.push(csr.edge_count()); - let mut fill = node + 1; - while fill < key { - csr.offsets.push(csr.edge_count()); - fill += 1; - } - current_node = Some(key); - } - Some(node) if key == node => {} - Some(node) => { - return Err(GfError::Storage(format!( - "adjacency spill merge produced non-sorted keys ({node} then {key})" - ))); - } - } - - csr.edge_ids.push(edge); - csr.neighbor_ids.push(neighbor); + emit((key, edge, neighbor))?; cursors[idx].pull()?; if let Some((k, e, n)) = cursors[idx].current { heap.push(Reverse((k, e, n, idx))); } } - if current_node.is_some() { - csr.offsets.push(csr.edge_count()); + Ok(()) +} + +fn csr_from_sorted_keyed(keyed: &[(u64, u64, u64)]) -> CsrIndex { + let mut csr = CsrIndex { + offsets: vec![0], + ..CsrIndex::default() + }; + for &(key, edge, neighbor) in keyed { + while csr.node_count() <= key { + csr.offsets.push(csr.edge_count()); + } + csr.edge_ids.push(edge); + csr.neighbor_ids.push(neighbor); + *csr.offsets.last_mut().expect("offset") = csr.edge_count(); } - Ok(csr) + csr } fn stream_build_groups( From eed6430f047e12903a1ef8bc2d3800cf4d264d9a Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 08:03:17 -0600 Subject: [PATCH 3/7] feat(storage): stream adjacency into bounded CSR shards --- crates/graphforge-api/src/search_index.rs | 3 + crates/graphforge-exec/src/adjacency.rs | 123 +++- crates/graphforge-storage/src/adjacency.rs | 742 +++++++++++++++++---- docs/book/architecture/storage.md | 47 +- 4 files changed, 747 insertions(+), 168 deletions(-) diff --git a/crates/graphforge-api/src/search_index.rs b/crates/graphforge-api/src/search_index.rs index 790ad5c2d..c2d64ab9a 100644 --- a/crates/graphforge-api/src/search_index.rs +++ b/crates/graphforge-api/src/search_index.rs @@ -377,6 +377,9 @@ impl GraphForge { spill_dir, spill_max_bytes: policy.spill_max_bytes, memory_budget_bytes: Some(policy.memory_budget_bytes), + shard_max_edges: graphforge_storage::adjacency::DEFAULT_CSR_SHARD_EDGES, + shard_max_nodes: graphforge_storage::adjacency::DEFAULT_CSR_SHARD_NODES, + merge_fan_in: graphforge_storage::adjacency::DEFAULT_ADJACENCY_MERGE_FAN_IN, } .effective() } diff --git a/crates/graphforge-exec/src/adjacency.rs b/crates/graphforge-exec/src/adjacency.rs index 8ce34d095..e53abe865 100644 --- a/crates/graphforge-exec/src/adjacency.rs +++ b/crates/graphforge-exec/src/adjacency.rs @@ -28,7 +28,8 @@ use arrow::record_batch::RecordBatch; use graphforge_core::{GfError, OntologyMode}; use graphforge_ir::Direction; use graphforge_storage::adjacency::{ - self as csr, ALL_RELATIONS_STEM, AdjacencyManifestRow, CsrIndex, CsrRow, build_adjacency_index, + self as csr, ALL_RELATIONS_STEM, AdjacencyManifestRow, CsrIndex, CsrRow, ShardedCsrIndex, + build_adjacency_index, }; use graphforge_storage::adjacency_delta::{ CsrDeltaOverlay, DeltaSegment, overlay_delta_segments, read_delta_chain, @@ -94,6 +95,12 @@ enum AdjacencyInner { Empty, Map(HashMap>), Csr(Arc), + Sharded(Arc), + ShardedOverlay { + base: Arc, + replaced: HashMap>, + node_extent: u64, + }, Overlay(CsrDeltaOverlay), Undirected { out: Arc, @@ -101,6 +108,8 @@ enum AdjacencyInner { }, } +type ShardedReplacementRows = HashMap>; + /// Borrowed or owned neighbor row for one node. #[derive(Clone, Debug)] pub struct NeighborRow<'a> { @@ -281,6 +290,10 @@ impl AdjacencyInner { Self::Empty => true, Self::Map(map) => map.is_empty(), Self::Csr(csr) => csr.edge_count() == 0, + Self::Sharded(csr) => csr.edge_count() == 0, + Self::ShardedOverlay { base, replaced, .. } => { + base.edge_count() == 0 && replaced.values().all(Vec::is_empty) + } Self::Overlay(overlay) => { overlay.base.edge_count() == 0 && overlay.replaced.values().all(Vec::is_empty) } @@ -293,6 +306,16 @@ impl AdjacencyInner { Self::Empty => NeighborRow::pairs(&[]), Self::Map(map) => NeighborRow::pairs(map.get(&node_id).map_or(&[], Vec::as_slice)), Self::Csr(csr) => NeighborRow::csr(csr.row(node_id)), + Self::Sharded(csr) => NeighborRow::owned( + csr.row(node_id) + .expect("authenticated immutable CSR shard changed after open"), + ), + Self::ShardedOverlay { base, replaced, .. } => { + NeighborRow::owned(replaced.get(&node_id).cloned().unwrap_or_else(|| { + base.row(node_id) + .expect("authenticated immutable CSR shard changed after open") + })) + } Self::Overlay(overlay) => match overlay.row(node_id) { graphforge_storage::adjacency_delta::OverlayRow::Base(row) => NeighborRow::csr(row), graphforge_storage::adjacency_delta::OverlayRow::Replaced(entries) => { @@ -308,15 +331,19 @@ impl AdjacencyInner { fn backing(&self) -> AdjacencyBacking { match self { Self::Empty | Self::Map(_) => AdjacencyBacking::ScanHashMap, - Self::Csr(_) => AdjacencyBacking::CsrNative, - Self::Overlay(_) => AdjacencyBacking::CsrOverlay, + Self::Csr(_) | Self::Sharded(_) => AdjacencyBacking::CsrNative, + Self::Overlay(_) | Self::ShardedOverlay { .. } => AdjacencyBacking::CsrOverlay, Self::Undirected { .. } => AdjacencyBacking::CsrUndirected, } } fn base_csr_entries_expanded(&self) -> u64 { match self { - Self::Empty | Self::Csr(_) | Self::Overlay(_) => 0, + Self::Empty + | Self::Csr(_) + | Self::Sharded(_) + | Self::Overlay(_) + | Self::ShardedOverlay { .. } => 0, Self::Map(map) => u64::try_from(map.values().map(Vec::len).sum::()).unwrap_or(0), Self::Undirected { out, inbound } => out .base_csr_entries_expanded() @@ -327,6 +354,9 @@ impl AdjacencyInner { fn overlay_row_count(&self) -> u64 { match self { Self::Overlay(overlay) => overlay.overlay_row_count(), + Self::ShardedOverlay { replaced, .. } => { + u64::try_from(replaced.len()).unwrap_or(u64::MAX) + } Self::Undirected { out, inbound } => out .overlay_row_count() .saturating_add(inbound.overlay_row_count()), @@ -341,7 +371,9 @@ impl AdjacencyInner { map.keys().copied().max().map_or(0, |m| m.saturating_add(1)) }), Self::Csr(csr) => csr.node_count(), + Self::Sharded(csr) => csr.node_count(), Self::Overlay(overlay) => overlay.node_extent, + Self::ShardedOverlay { node_extent, .. } => *node_extent, Self::Undirected { out, inbound } => out.node_extent().max(inbound.node_extent()), } } @@ -359,11 +391,21 @@ impl AdjacencyInner { visit(node_id, NeighborRow::csr(csr.row(node_id))); } } + Self::Sharded(csr) => { + for node_id in 0..csr.node_count() { + visit(node_id, self.neighbors(node_id)); + } + } Self::Overlay(overlay) => { for node_id in 0..overlay.node_extent { visit(node_id, self.neighbors(node_id)); } } + Self::ShardedOverlay { node_extent, .. } => { + for node_id in 0..*node_extent { + visit(node_id, self.neighbors(node_id)); + } + } Self::Undirected { out, inbound } => { let extent = out.node_extent().max(inbound.node_extent()); for node_id in 0..extent { @@ -690,7 +732,35 @@ impl PersistentAdjacencyProvider { deltas: &[DeltaSegment], ) -> Result { let directed = |d: csr::Direction| -> Result { - let base = Arc::new(csr::read_csr(&csr::csr_path(&self.dir, stem, d))?); + let path = csr::csr_path(&self.dir, stem, d); + if csr::sharded_csr_exists(&path) { + let base = Arc::new(ShardedCsrIndex::open(&path)?); + if let Some(row) = rows + .iter() + .find(|r| r.relation_type == stem && r.direction == d) + && (base.node_count() != row.node_count || base.edge_count() != row.edge_count) + { + return Err(GfError::Storage( + "adjacency sharded CSR disagrees with manifest counts (torn read)".into(), + )); + } + if deltas.is_empty() { + return Ok(AdjacencyInner::Sharded(base)); + } + let (replaced, node_extent) = sharded_overlay_rows(&base, stem, d, deltas)?; + if replaced.is_empty() { + return Ok(AdjacencyInner::Sharded(base)); + } + return Ok(AdjacencyInner::ShardedOverlay { + base, + replaced, + node_extent, + }); + } + // Legacy single-batch CSR migration path. A successful rebuild + // publishes sharded v1 files; old projects remain readable until + // that rebuild occurs. + let base = Arc::new(csr::read_csr(&path)?); if deltas.is_empty() { return Ok(AdjacencyInner::Csr(base)); } @@ -842,6 +912,44 @@ impl PersistentAdjacencyProvider { } } +fn sharded_overlay_rows( + base: &ShardedCsrIndex, + stem: &str, + direction: csr::Direction, + chain: &[DeltaSegment], +) -> Result<(ShardedReplacementRows, u64), GfError> { + let take_all = stem == ALL_RELATIONS_STEM; + let mut by_key: HashMap> = HashMap::new(); + let mut max_key = base.node_count().saturating_sub(1); + for segment in chain { + for edge in &segment.edges { + if take_all || edge.rel_type_name == stem { + let (key, neighbor) = match direction { + csr::Direction::Out => (edge.src_id, edge.dst_id), + csr::Direction::In => (edge.dst_id, edge.src_id), + }; + by_key + .entry(key) + .or_default() + .push((edge.edge_id, neighbor)); + max_key = max_key.max(key); + } + } + } + for (key, delta) in &mut by_key { + let mut combined = base.row(*key)?; + combined.append(delta); + combined.sort_unstable_by_key(|&(edge_id, _)| edge_id); + *delta = combined; + } + let extent = if by_key.is_empty() { + base.node_count() + } else { + base.node_count().max(max_key.saturating_add(1)) + }; + Ok((by_key, extent)) +} + impl AdjacencyProvider for PersistentAdjacencyProvider { fn adjacency( &self, @@ -903,7 +1011,10 @@ impl AdjacencyProvider for PersistentAdjacencyProvider { IndexState::Ready { fresh: true, rows, .. } => { - let files_exist = |d: csr::Direction| csr::csr_path(&self.dir, &stem, d).exists(); + let files_exist = |d: csr::Direction| { + let path = csr::csr_path(&self.dir, &stem, d); + path.exists() || csr::sharded_csr_exists(&path) + }; let present = Self::rows_cover(&rows, &stem, direction) && match direction { Direction::Out => files_exist(csr::Direction::Out), diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index 77f8b47ca..398363aa8 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -9,28 +9,31 @@ //! ```text //! indexes/adjacency/ //! ├── index_manifest.parquet ADJACENCY_MANIFEST_SCHEMA (Parquet) -//! ├── WORKS_AT.out.csr ADJACENCY_CSR_SCHEMA (Arrow IPC) -//! ├── WORKS_AT.in.csr -//! └── _all.out.csr union across relation types +//! ├── WORKS_AT.out.csr.json versioned/checksummed shard manifest +//! ├── WORKS_AT.out.csr.shards-.d/ +//! ├── WORKS_AT.in.csr.json +//! └── _all.out.csr.json union across relation types //! ``` //! //! # Build ordering convention //! -//! Builders MUST write all CSR files first and `index_manifest.parquet` -//! **last**. A crash mid-build then leaves the manifest absent or carrying the +//! Builders MUST write immutable shard files, publish each shard manifest +//! atomically, and write `index_manifest.parquet` **last**. A crash mid-build +//! then leaves the manifest absent or carrying the //! old `topology_generation`, so the index reads as stale and the provider //! falls back to scan-and-build — a torn build can cost a rebuild, never //! correctness. //! //! # CSR encoding //! -//! A `.csr` file is a single-batch Arrow IPC file with one column, +//! Each bounded shard `.csr` file is Arrow IPC with one column, //! `adjacency: LargeList` and one row per -//! surrogate `node_id` in `0..node_count`. The list offsets buffer is the CSR +//! local surrogate range. A high-degree logical row may continue in the next +//! shard. The list offsets buffer is the CSR //! offsets array; the struct child is the targets array. See //! [`ADJACENCY_CSR_SCHEMA`] and `docs/book/architecture/storage.md` §Derived -//! Indexes. The in-memory [`CsrIndex`] exposes the logical `offsets`/`targets` -//! model directly, so consumers never deal with the list encoding. +//! Indexes. [`ShardedCsrIndex`] resolves only the shard(s) containing a requested +//! row. Legacy single-batch [`CsrIndex`] files remain readable until rebuild. use std::fs::File; use std::path::{Path, PathBuf}; @@ -62,6 +65,10 @@ pub const ALL_RELATIONS_STEM: &str = "_all"; pub const MANIFEST_FILE: &str = "index_manifest.parquet"; const SHARDED_CSR_VERSION: u32 = 1; +/// Default maximum adjacency entries materialized in one persisted CSR shard. +pub const DEFAULT_CSR_SHARD_EDGES: usize = 1_048_576; +/// Default maximum local CSR rows (offset entries minus one) per shard. +pub const DEFAULT_CSR_SHARD_NODES: usize = 1_048_576; #[derive(Clone, Debug, Serialize, Deserialize)] struct CsrShardRecord { @@ -82,6 +89,23 @@ struct CsrShardManifest { shards: Vec, } +/// Aggregate bounded-resource evidence for one adjacency build. +#[derive(Clone, Debug, Default, Eq, PartialEq)] +pub struct AdjacencyBuildMetrics { + /// Projected source rows consumed from Parquet. + pub source_rows: u64, + /// Sorted spill runs written across relation and union accumulators. + pub spill_runs: u64, + /// Peak bytes charged to the spill session. + pub spill_bytes: u64, + /// Persisted CSR shards written across every relation/direction pair. + pub csr_shards: u64, + /// Largest number of entries retained by a CSR shard sink. + pub peak_shard_edges: u64, + /// Largest number of local CSR rows retained by a shard sink. + pub peak_shard_nodes: u64, +} + /// Bounded reader for a versioned sharded CSR. Opening validates the small /// manifest; row access reads and authenticates only the containing shard. #[derive(Clone, Debug)] @@ -102,30 +126,54 @@ impl ShardedCsrIndex { manifest_path.display() ))); } - let mut expected = 0_u64; + let mut prior_first = None; let mut edges = 0_u64; + let root = path + .parent() + .unwrap_or_else(|| Path::new(".")) + .join(&manifest.shard_dir); for shard in &manifest.shards { - if shard.first_node != expected || shard.node_count == 0 { - return Err(GfError::Storage("invalid CSR shard boundary chain".into())); + if shard.node_count == 0 + || shard.first_node.saturating_add(shard.node_count) > manifest.node_count + || prior_first.is_some_and(|prior| shard.first_node < prior) + { + return Err(GfError::Storage( + "invalid CSR shard boundary ordering".into(), + )); } - expected = expected.saturating_add(shard.node_count); + prior_first = Some(shard.first_node); edges = edges.saturating_add(shard.edge_count); if Path::new(&shard.file).components().count() != 1 { return Err(GfError::Storage("invalid CSR shard file name".into())); } + // Authenticate and structurally validate every bounded shard at + // open. Runtime row reads then cannot discover latent corruption + // through the infallible traversal interface. + let shard_path = root.join(&shard.file); + let shard_bytes = std::fs::read(&shard_path).map_err(|error| { + GfError::Storage(format!("missing CSR shard {}: {error}", shard.file)) + })?; + if sha256_hex(&shard_bytes) != shard.sha256 { + return Err(GfError::Storage(format!( + "CSR shard checksum mismatch: {}", + shard.file + ))); + } + let decoded = read_csr_bytes(&shard_bytes, &shard_path)?; + if decoded.node_count() != shard.node_count || decoded.edge_count() != shard.edge_count + { + return Err(GfError::Storage(format!( + "CSR shard count mismatch: {}", + shard.file + ))); + } } - if expected != manifest.node_count || edges != manifest.edge_count { + if edges != manifest.edge_count { return Err(GfError::Storage( "CSR shard manifest counts disagree".into(), )); } - Ok(Self { - root: path - .parent() - .unwrap_or_else(|| Path::new(".")) - .join(&manifest.shard_dir), - manifest, - }) + Ok(Self { root, manifest }) } /// Total logical source rows across all shards. @@ -145,20 +193,32 @@ impl ShardedCsrIndex { if node_id >= self.manifest.node_count { return Ok(Vec::new()); } - let shard = self + // A single high-degree row may span adjacent hard-capped shards. The + // manifest is ordered by `first_node`, so stop once starts pass the key. + let end = self .manifest .shards - .binary_search_by(|candidate| { - if node_id < candidate.first_node { - std::cmp::Ordering::Greater - } else if node_id >= candidate.first_node.saturating_add(candidate.node_count) { - std::cmp::Ordering::Less - } else { - std::cmp::Ordering::Equal - } - }) - .map_err(|_| GfError::Storage("CSR shard boundary lookup failed".into()))?; - let record = &self.manifest.shards[shard]; + .partition_point(|record| record.first_node <= node_id); + let mut start = end; + while start > 0 { + let prior = &self.manifest.shards[start - 1]; + if node_id >= prior.first_node.saturating_add(prior.node_count) { + break; + } + start -= 1; + } + let mut output = Vec::new(); + for record in &self.manifest.shards[start..end] { + output.extend(self.read_record_row(record, node_id)?); + } + Ok(output) + } + + fn read_record_row( + &self, + record: &CsrShardRecord, + node_id: u64, + ) -> Result, GfError> { let path = self.root.join(&record.file); let bytes = std::fs::read(&path).map_err(|error| { GfError::Storage(format!("missing CSR shard {}: {error}", path.display())) @@ -180,72 +240,190 @@ impl ShardedCsrIndex { } } +/// Whether the versioned sharded representation is published for `path`. +#[must_use] +pub fn sharded_csr_exists(path: &Path) -> bool { + path.with_extension("csr.json").is_file() +} + +fn csr_artifact_exists(path: &Path) -> bool { + path.is_file() || sharded_csr_exists(path) +} + /// Write a versioned checksummed shard set and publish its manifest last. /// This compatibility helper accepts an in-memory CSR; streaming builders use /// the same shard writer while producing one bounded shard at a time. pub fn write_sharded_csr(path: &Path, csr: &CsrIndex, max_edges: usize) -> Result<(), GfError> { csr.validate()?; - let parent = path - .parent() - .ok_or_else(|| GfError::Storage("CSR path has no parent".into()))?; - std::fs::create_dir_all(parent).map_err(storage_err)?; - let stem = path - .file_name() - .and_then(|name| name.to_str()) - .unwrap_or("csr"); - let shard_dir = format!("{stem}.{}.d", uuid::Uuid::new_v4().as_simple()); - let root = parent.join(&shard_dir); - std::fs::create_dir(&root).map_err(storage_err)?; - let result = (|| { - let mut records = Vec::new(); - let mut first_node = 0_u64; - let mut shard = CsrIndex { + let mut writer = ShardedCsrWriter::create(path, max_edges, DEFAULT_CSR_SHARD_NODES)?; + for node in 0..csr.node_count() { + for (edge, neighbor) in csr.row(node).iter() { + writer.emit((node, edge, neighbor))?; + } + } + writer.finish(csr.node_count()).map(|_| ()) +} + +/// Row-emitting bounded sink used directly by the external-run merge. +struct ShardedCsrWriter { + path: PathBuf, + root: PathBuf, + shard_dir: String, + max_edges: usize, + max_nodes: usize, + records: Vec, + shard: CsrIndex, + first_node: Option, + last_key: Option, + edge_count: u64, + peak_shard_edges: u64, + peak_shard_nodes: u64, + finished: bool, +} + +impl ShardedCsrWriter { + fn create(path: &Path, max_edges: usize, max_nodes: usize) -> Result { + let parent = path + .parent() + .ok_or_else(|| GfError::Storage("CSR path has no parent".into()))?; + std::fs::create_dir_all(parent).map_err(storage_err)?; + let stem = path.file_name().and_then(|n| n.to_str()).unwrap_or("csr"); + let shard_dir = format!("{stem}.{}.d", uuid::Uuid::new_v4().as_simple()); + let root = parent.join(&shard_dir); + std::fs::create_dir(&root).map_err(storage_err)?; + Ok(Self { + path: path.to_path_buf(), + root, + shard_dir, + max_edges: max_edges.max(1), + max_nodes: max_nodes.max(1), + records: Vec::new(), + shard: CsrIndex { + offsets: vec![0], + ..CsrIndex::default() + }, + first_node: None, + last_key: None, + edge_count: 0, + peak_shard_edges: 0, + peak_shard_nodes: 0, + finished: false, + }) + } + + fn emit(&mut self, (key, edge, neighbor): (u64, u64, u64)) -> Result<(), GfError> { + if self.shard.edge_ids.len() >= self.max_edges + || self.first_node.is_some_and(|first| { + key.saturating_sub(first) >= u64::try_from(self.max_nodes).unwrap_or(u64::MAX) + }) + { + self.flush()?; + } + let first = *self.first_node.get_or_insert(key); + if key < self.last_key.unwrap_or(key) { + return Err(GfError::Storage( + "CSR shard sink received unsorted row keys".into(), + )); + } + let local = key - first; + while self.shard.node_count() <= local { + self.shard.offsets.push(self.shard.edge_count()); + } + self.shard.edge_ids.push(edge); + self.shard.neighbor_ids.push(neighbor); + *self.shard.offsets.last_mut().expect("offset exists") = self.shard.edge_count(); + self.last_key = Some(key); + self.edge_count = self.edge_count.saturating_add(1); + self.peak_shard_edges = self.peak_shard_edges.max(self.shard.edge_count()); + self.peak_shard_nodes = self.peak_shard_nodes.max(self.shard.node_count()); + Ok(()) + } + + fn flush(&mut self) -> Result<(), GfError> { + let Some(first) = self.first_node else { + return Ok(()); + }; + self.records.push(write_csr_shard( + &self.root, + first, + &self.shard, + self.records.len(), + )?); + self.shard = CsrIndex { offsets: vec![0], ..CsrIndex::default() }; - for node in 0..csr.node_count() { - let row = csr.row(node); - if shard.node_count() > 0 - && shard.edge_ids.len().saturating_add(row.len()) > max_edges.max(1) - { - records.push(write_csr_shard(&root, first_node, &shard, records.len())?); - first_node = node; - shard = CsrIndex { - offsets: vec![0], - ..CsrIndex::default() - }; - } - shard.edge_ids.extend_from_slice(row.edge_ids); - shard.neighbor_ids.extend_from_slice(row.neighbor_ids); - shard.offsets.push(shard.edge_count()); - } - if shard.node_count() > 0 || records.is_empty() { - records.push(write_csr_shard(&root, first_node, &shard, records.len())?); + self.first_node = None; + self.last_key = None; + Ok(()) + } + + fn finish(mut self, node_count: u64) -> Result<(u64, u64, u64), GfError> { + use std::io::Write as _; + + self.flush()?; + let mut identity = Sha256::new(); + identity.update(b"graphforge/csr-shards/v1\0"); + identity.update(node_count.to_le_bytes()); + identity.update(self.edge_count.to_le_bytes()); + for record in &self.records { + identity.update(record.first_node.to_le_bytes()); + identity.update(record.node_count.to_le_bytes()); + identity.update(record.edge_count.to_le_bytes()); + identity.update(record.sha256.as_bytes()); + } + let digest = sha256_hex(&identity.finalize()); + let stem = self + .path + .file_name() + .and_then(|n| n.to_str()) + .unwrap_or("csr"); + let stable_dir = format!("{stem}.shards-{}.d", &digest[..24]); + let stable_root = self + .path + .parent() + .expect("validated parent") + .join(&stable_dir); + if stable_root.exists() { + std::fs::remove_dir_all(&self.root).map_err(storage_err)?; + } else { + std::fs::rename(&self.root, &stable_root).map_err(storage_err)?; } + self.root = stable_root; + self.shard_dir = stable_dir; let manifest = CsrShardManifest { format: "graphforge.csr-shards".into(), version: SHARDED_CSR_VERSION, - node_count: csr.node_count(), - edge_count: csr.edge_count(), - shard_dir, - shards: records, + node_count, + edge_count: self.edge_count, + shard_dir: self.shard_dir.clone(), + shards: std::mem::take(&mut self.records), }; let bytes = serde_json::to_vec_pretty(&manifest).map_err(storage_err)?; - let manifest_path = path.with_extension("csr.json"); + let parent = self.path.parent().expect("validated parent"); let mut temp = tempfile::Builder::new() .prefix(stem) .suffix(".json.tmp") .tempfile_in(parent) .map_err(storage_err)?; - use std::io::Write as _; temp.write_all(&bytes).map_err(storage_err)?; temp.as_file().sync_all().map_err(storage_err)?; - persist_temp(temp, &manifest_path) - })(); - if result.is_err() { - let _ = std::fs::remove_dir_all(&root); + persist_temp(temp, &self.path.with_extension("csr.json"))?; + self.finished = true; + Ok(( + manifest.shards.len() as u64, + self.peak_shard_edges, + self.peak_shard_nodes, + )) + } +} + +impl Drop for ShardedCsrWriter { + fn drop(&mut self) { + if !self.finished { + let _ = std::fs::remove_dir_all(&self.root); + } } - result } fn write_csr_shard( @@ -619,7 +797,15 @@ pub fn write_csr(path: &Path, csr: &CsrIndex) -> Result<(), GfError> { FileWriter::try_new(tmp.as_file(), &ADJACENCY_CSR_SCHEMA).map_err(storage_err)?; writer.write(&batch).map_err(storage_err)?; writer.finish().map_err(storage_err)?; - persist_temp(tmp, path) + persist_temp(tmp, path)?; + // Explicit legacy writes are a supported migration/testing seam. Make the + // representation choice unambiguous so a stale sharded manifest cannot + // shadow the newly published single-batch file. + let sharded_manifest = path.with_extension("csr.json"); + if sharded_manifest.exists() { + std::fs::remove_file(sharded_manifest).map_err(storage_err)?; + } + Ok(()) } /// Read a CSR file written by [`write_csr`] back into a [`CsrIndex`]. @@ -628,6 +814,22 @@ pub fn write_csr(path: &Path, csr: &CsrIndex) -> Result<(), GfError> { /// Returns [`GfError::Storage`] if the file is missing, is not an Arrow IPC /// file with [`ADJACENCY_CSR_SCHEMA`], or decodes to an invalid CSR. pub fn read_csr(path: &Path) -> Result { + if sharded_csr_exists(path) { + let sharded = ShardedCsrIndex::open(path)?; + let mut csr = CsrIndex { + offsets: vec![0], + ..CsrIndex::default() + }; + for node in 0..sharded.node_count() { + for (edge, neighbor) in sharded.row(node)? { + csr.edge_ids.push(edge); + csr.neighbor_ids.push(neighbor); + } + csr.offsets.push(csr.edge_count()); + } + csr.validate()?; + return Ok(csr); + } let file = File::open(path) .map_err(|e| GfError::Storage(format!("cannot open CSR file {}: {e}", path.display())))?; let reader = FileReader::try_new(file, None) @@ -824,6 +1026,8 @@ pub const DEFAULT_ADJACENCY_BATCH_SIZE: usize = 8_192; /// copies. Peak working set remains a function of this budget (and the /// configured memory/spill caps), not total edge count. pub const DEFAULT_ADJACENCY_CHUNK_ROWS: usize = 1_048_576; +/// Maximum sorted runs opened concurrently by one merge pass. +pub const DEFAULT_ADJACENCY_MERGE_FAN_IN: usize = 64; /// Spill subdirectory name under the artifact adjacency directory when no /// explicit spill root is configured. @@ -857,6 +1061,13 @@ pub struct AdjacencyBuildOptions { /// Soft memory budget used to shrink [`chunk_rows`](Self::chunk_rows) when /// set. Does not replace the hard spill-byte cap. pub memory_budget_bytes: Option, + /// Hard upper bound on adjacency entries retained by one CSR shard sink. + /// A single high-degree row is split across consecutive shards when needed. + pub shard_max_edges: usize, + /// Hard upper bound on local CSR rows (offset entries minus one) per shard. + pub shard_max_nodes: usize, + /// Maximum spill runs opened in one k-way merge pass (minimum 2). + pub merge_fan_in: usize, } impl Default for AdjacencyBuildOptions { @@ -867,6 +1078,9 @@ impl Default for AdjacencyBuildOptions { spill_dir: None, spill_max_bytes: None, memory_budget_bytes: None, + shard_max_edges: DEFAULT_CSR_SHARD_EDGES, + shard_max_nodes: DEFAULT_CSR_SHARD_NODES, + merge_fan_in: DEFAULT_ADJACENCY_MERGE_FAN_IN, } } } @@ -878,6 +1092,9 @@ impl AdjacencyBuildOptions { pub fn effective(&self) -> Self { let mut out = self.clone(); out.batch_size = out.batch_size.max(1); + out.shard_max_edges = out.shard_max_edges.max(1); + out.shard_max_nodes = out.shard_max_nodes.max(1); + out.merge_fan_in = out.merge_fan_in.max(2); let mut chunk = out.chunk_rows.max(1); if let Some(budget) = out.memory_budget_bytes.filter(|b| *b > 0) { // Leave headroom for out+in keyed copies (~2×) plus CSR/merge state. @@ -970,6 +1187,24 @@ pub fn build_adjacency_index_into_with_options( options: &AdjacencyBuildOptions, mut checkpoint: impl FnMut() -> Result<(), GfError>, ) -> Result, GfError> { + build_adjacency_index_into_with_metrics( + source_project_dir, + artifact_project_dir, + built_at_micros, + options, + &mut checkpoint, + ) + .map(|(manifest, _)| manifest) +} + +/// Bounded build returning explicit source/spill/shard resource counters. +pub fn build_adjacency_index_into_with_metrics( + source_project_dir: &Path, + artifact_project_dir: &Path, + built_at_micros: i64, + options: &AdjacencyBuildOptions, + mut checkpoint: impl FnMut() -> Result<(), GfError>, +) -> Result<(Vec, AdjacencyBuildMetrics), GfError> { checkpoint()?; // Generation BEFORE the scan — see the race note in the doc comment. let generation = crate::generation::read_topology_generation(source_project_dir)?; @@ -988,25 +1223,39 @@ pub fn build_adjacency_index_into_with_options( base.join(format!("build-{}", uuid::Uuid::new_v4().as_simple())) }; let mut spill = SpillSession::create(&spill_root)?.with_max_bytes(options.spill_max_bytes); + let mut metrics = AdjacencyBuildMetrics::default(); let build_result = (|| { - let mut groups = - stream_build_groups(source_project_dir, &options, &mut spill, &mut checkpoint)?; + let mut groups = stream_build_groups( + source_project_dir, + &options, + &mut spill, + &mut metrics, + &mut checkpoint, + )?; checkpoint()?; let mut manifest = Vec::new(); let mut write_pair = |stem: &str, group: &mut EntryGroup| -> Result<(), GfError> { for direction in [Direction::Out, Direction::In] { checkpoint()?; - let csr = group.finish_csr(direction, &mut spill, &mut checkpoint)?; - write_csr(&csr_path(artifact_project_dir, stem, direction), &csr)?; + let outcome = group.finish_sharded_csr( + direction, + &csr_path(artifact_project_dir, stem, direction), + &options, + &mut spill, + &mut checkpoint, + )?; + metrics.csr_shards = metrics.csr_shards.saturating_add(outcome.shards); + metrics.peak_shard_edges = metrics.peak_shard_edges.max(outcome.peak_shard_edges); + metrics.peak_shard_nodes = metrics.peak_shard_nodes.max(outcome.peak_shard_nodes); manifest.push(AdjacencyManifestRow { relation_type: stem.to_owned(), direction, topology_generation: generation, built_at_micros, - node_count: csr.node_count(), - edge_count: csr.edge_count(), + node_count: outcome.node_count, + edge_count: outcome.edge_count, }); } Ok(()) @@ -1035,13 +1284,15 @@ pub fn build_adjacency_index_into_with_options( // survive, so the new base + those is immediately fresh. Manifest first, // prune after: a crash between leaves dead segments a later prune removes. crate::adjacency_delta::prune_delta_segments(artifact_project_dir, generation); - Ok(manifest) + metrics.spill_runs = spill.run_counter; + metrics.spill_bytes = spill.peak_bytes; + Ok((manifest, metrics.clone())) })(); match build_result { - Ok(manifest) => { + Ok(result) => { spill.cleanup(); - Ok(manifest) + Ok(result) } Err(error) => { spill.cleanup(); @@ -1054,7 +1305,8 @@ pub fn build_adjacency_index_into_with_options( /// and failure cannot leave temporary runs behind as a published artifact. struct SpillSession { root: PathBuf, - bytes_written: u64, + bytes_current: u64, + peak_bytes: u64, max_bytes: Option, run_counter: u64, cleaned: bool, @@ -1067,7 +1319,8 @@ impl SpillSession { std::fs::create_dir_all(root).map_err(storage_err)?; Ok(Self { root: root.to_path_buf(), - bytes_written: 0, + bytes_current: 0, + peak_bytes: 0, max_bytes: None, run_counter: 0, cleaned: false, @@ -1087,9 +1340,10 @@ impl SpillSession { } fn account_write(&mut self, bytes: u64) -> Result<(), GfError> { - self.bytes_written = self.bytes_written.saturating_add(bytes); + self.bytes_current = self.bytes_current.saturating_add(bytes); + self.peak_bytes = self.peak_bytes.max(self.bytes_current); if let Some(max) = self.max_bytes - && self.bytes_written > max + && self.bytes_current > max { return Err(resource_limit(format!( "adjacency build spill exceeded max_bytes ({max})" @@ -1098,6 +1352,13 @@ impl SpillSession { Ok(()) } + fn remove_run(&mut self, path: &Path) -> Result<(), GfError> { + let bytes = std::fs::metadata(path).map_err(storage_err)?.len(); + std::fs::remove_file(path).map_err(storage_err)?; + self.bytes_current = self.bytes_current.saturating_sub(bytes); + Ok(()) + } + fn cleanup(&mut self) { if self.cleaned { return; @@ -1195,33 +1456,84 @@ impl EntryGroup { Ok(()) } - fn finish_csr( + fn finish_sharded_csr( &mut self, direction: Direction, + path: &Path, + options: &AdjacencyBuildOptions, spill: &mut SpillSession, checkpoint: &mut dyn FnMut() -> Result<(), GfError>, - ) -> Result { - let runs = match direction { - Direction::Out => &self.out_runs, - Direction::In => &self.in_runs, + ) -> Result { + let had_runs = match direction { + Direction::Out => !self.out_runs.is_empty(), + Direction::In => !self.in_runs.is_empty(), }; - if runs.is_empty() { - // Fast path: never spilled — identical to the pre-#336 builder. - return Ok(csr_from_entries(&self.buffer, direction)); - } - // Spill residual as a final run so merge sees only sorted runs. - if !self.buffer.is_empty() { + if had_runs && !self.buffer.is_empty() { self.flush(spill, checkpoint)?; } - let runs = match direction { - Direction::Out => &self.out_runs, - Direction::In => &self.in_runs, + if had_runs { + let runs = match direction { + Direction::Out => &mut self.out_runs, + Direction::In => &mut self.in_runs, + }; + compact_keyed_runs( + runs, + options.merge_fan_in, + &self.label, + direction, + spill, + checkpoint, + )?; + } + let mut writer = + ShardedCsrWriter::create(path, options.shard_max_edges, options.shard_max_nodes)?; + let mut max_key = None::; + let mut emit = |entry: (u64, u64, u64)| { + max_key = Some(max_key.map_or(entry.0, |prior| prior.max(entry.0))); + writer.emit(entry) }; - checkpoint()?; - merge_keyed_runs_to_csr(runs, checkpoint) + if had_runs { + let runs = match direction { + Direction::Out => &self.out_runs, + Direction::In => &self.in_runs, + }; + merge_keyed_runs(runs, checkpoint, &mut emit)?; + } else { + // The no-spill fast path remains bounded by `chunk_rows`. + let mut keyed: Vec<_> = match direction { + Direction::Out => self.buffer.clone(), + Direction::In => self + .buffer + .iter() + .map(|&(src, edge, dst)| (dst, edge, src)) + .collect(), + }; + keyed.sort_unstable_by_key(|&(key, edge, _)| (key, edge)); + for entry in keyed { + emit(entry)?; + } + } + let node_count = max_key.map_or(0, |key| key.saturating_add(1)); + let edge_count = writer.edge_count; + let (shards, peak_shard_edges, peak_shard_nodes) = writer.finish(node_count)?; + Ok(ShardedWriteOutcome { + node_count, + edge_count, + shards, + peak_shard_edges, + peak_shard_nodes, + }) } } +struct ShardedWriteOutcome { + node_count: u64, + edge_count: u64, + shards: u64, + peak_shard_edges: u64, + peak_shard_nodes: u64, +} + fn write_keyed_run( path: &Path, entries: &[(u64, u64, u64)], @@ -1306,16 +1618,82 @@ impl RunCursor { } } -fn merge_keyed_runs_to_csr( - runs: &[PathBuf], +fn compact_keyed_runs( + runs: &mut Vec, + fan_in: usize, + label: &str, + direction: Direction, + spill: &mut SpillSession, checkpoint: &mut dyn FnMut() -> Result<(), GfError>, -) -> Result { - let mut keyed = Vec::new(); - merge_keyed_runs(runs, checkpoint, &mut |entry| { - keyed.push(entry); - Ok(()) +) -> Result<(), GfError> { + let fan_in = fan_in.max(2); + while runs.len() > fan_in { + let mut next = Vec::with_capacity(runs.len().div_ceil(fan_in)); + for chunk in runs.chunks(fan_in) { + checkpoint()?; + if chunk.len() == 1 { + next.push(chunk[0].clone()); + continue; + } + let output = spill.next_run_path(label, direction); + merge_keyed_runs_to_run(chunk, &output, spill, checkpoint)?; + for input in chunk { + spill.remove_run(input)?; + } + next.push(output); + } + *runs = next; + } + Ok(()) +} + +fn keyed_run_count(path: &Path) -> Result { + use std::io::Read; + let mut file = std::io::BufReader::new(std::fs::File::open(path).map_err(storage_err)?); + let mut header = [0_u8; 20]; + file.read_exact(&mut header).map_err(storage_err)?; + if &header[..8] != SPILL_RUN_MAGIC + || u32::from_le_bytes(header[8..12].try_into().expect("four bytes")) != SPILL_RUN_VERSION + { + return Err(GfError::Storage(format!( + "adjacency spill run {} has invalid header", + path.display() + ))); + } + Ok(u64::from_le_bytes( + header[12..20].try_into().expect("eight bytes"), + )) +} + +fn merge_keyed_runs_to_run( + inputs: &[PathBuf], + output: &Path, + spill: &mut SpillSession, + checkpoint: &mut dyn FnMut() -> Result<(), GfError>, +) -> Result<(), GfError> { + use std::io::{BufWriter, Write}; + let count = inputs.iter().try_fold(0_u64, |total, path| { + keyed_run_count(path).map(|count| total.saturating_add(count)) })?; - Ok(csr_from_sorted_keyed(&keyed)) + let bytes = 20_u64.saturating_add(count.saturating_mul(BYTES_PER_KEYED_ENTRY)); + spill.account_write(bytes)?; + let mut writer = + BufWriter::with_capacity(1 << 20, std::fs::File::create(output).map_err(storage_err)?); + writer.write_all(SPILL_RUN_MAGIC).map_err(storage_err)?; + writer + .write_all(&SPILL_RUN_VERSION.to_le_bytes()) + .map_err(storage_err)?; + writer + .write_all(&count.to_le_bytes()) + .map_err(storage_err)?; + merge_keyed_runs(inputs, checkpoint, &mut |(key, edge, neighbor)| { + writer.write_all(&key.to_le_bytes()).map_err(storage_err)?; + writer.write_all(&edge.to_le_bytes()).map_err(storage_err)?; + writer + .write_all(&neighbor.to_le_bytes()) + .map_err(storage_err) + })?; + writer.flush().map_err(storage_err) } fn merge_keyed_runs( @@ -1359,26 +1737,11 @@ fn merge_keyed_runs( Ok(()) } -fn csr_from_sorted_keyed(keyed: &[(u64, u64, u64)]) -> CsrIndex { - let mut csr = CsrIndex { - offsets: vec![0], - ..CsrIndex::default() - }; - for &(key, edge, neighbor) in keyed { - while csr.node_count() <= key { - csr.offsets.push(csr.edge_count()); - } - csr.edge_ids.push(edge); - csr.neighbor_ids.push(neighbor); - *csr.offsets.last_mut().expect("offset") = csr.edge_count(); - } - csr -} - fn stream_build_groups( project_dir: &Path, options: &AdjacencyBuildOptions, spill: &mut SpillSession, + metrics: &mut AdjacencyBuildMetrics, checkpoint: &mut dyn FnMut() -> Result<(), GfError>, ) -> Result, GfError> { spill.max_bytes = options.spill_max_bytes; @@ -1406,6 +1769,7 @@ fn stream_build_groups( None }; for i in 0..batch.num_rows() { + metrics.source_rows = metrics.source_rows.saturating_add(1); let entry = (src_ids.value(i), edge_ids.value(i), dst_ids.value(i)); groups .get_mut(ALL_RELATIONS_STEM) @@ -1686,7 +2050,7 @@ pub fn validate_adjacency_index_against( }; let expected = csr_from_entries(expected_entries, row.direction); let path = csr_path(artifact_project_dir, &row.relation_type, row.direction); - if !path.exists() { + if !csr_artifact_exists(&path) { issues.push(AdjacencyValidationIssue::MissingCsr { rel: row.relation_type.clone(), direction: row.direction, @@ -1800,7 +2164,7 @@ pub fn inspect_adjacency_index(project_dir: &Path) -> Result>() + ); + } + + #[test] + fn sparse_surrogate_gap_cannot_expand_one_shard_offsets() { + let directory = TempDir::new().unwrap(); + let path = directory.path().join("KNOWS.out.csr"); + let mut writer = ShardedCsrWriter::create(&path, 8, 2).unwrap(); + writer.emit((0, 1, 100)).unwrap(); + writer.emit((1_000_000, 2, 0)).unwrap(); + let (shards, _, peak_nodes) = writer.finish(1_000_001).unwrap(); + assert_eq!(shards, 2); + assert!(peak_nodes <= 2); + let reader = ShardedCsrIndex::open(&path).unwrap(); + assert_eq!(reader.row(0).unwrap(), vec![(1, 100)]); + assert_eq!(reader.row(999_999).unwrap(), Vec::<(u64, u64)>::new()); + assert_eq!(reader.row(1_000_000).unwrap(), vec![(2, 0)]); + } + #[test] fn entries_from_out_csr_rejects_malformed_offset_bounds() { let past_targets = CsrIndex { @@ -2457,18 +2869,19 @@ mod tests { // KNOWS out/in + _all out/in. assert_eq!(rows.len(), 4); - let knows_out = std::fs::read(csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap(); - let all_in = - std::fs::read(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::In)).unwrap(); + let knows_path = csr_path(dir.path(), "KNOWS", Direction::Out); + let all_in_path = csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::In); + let knows_out = std::fs::read(knows_path.with_extension("csr.json")).unwrap(); + let all_in = std::fs::read(all_in_path.with_extension("csr.json")).unwrap(); // Rebuild: byte-identical CSR files (R-ADJ-2). build_adjacency_index(dir.path(), BUILD_TS).unwrap(); assert_eq!( - std::fs::read(csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap(), + std::fs::read(knows_path.with_extension("csr.json")).unwrap(), knows_out ); assert_eq!( - std::fs::read(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::In)).unwrap(), + std::fs::read(all_in_path.with_extension("csr.json")).unwrap(), all_in ); @@ -2698,7 +3111,13 @@ mod tests { let dir = TempDir::new().unwrap(); write_diamond(dir.path()); build_adjacency_index(dir.path(), BUILD_TS).unwrap(); - std::fs::write(csr_path(dir.path(), "KNOWS", Direction::In), b"garbage").unwrap(); + let path = csr_path(dir.path(), "KNOWS", Direction::In); + let reader = ShardedCsrIndex::open(&path).unwrap(); + std::fs::write( + reader.root.join(&reader.manifest.shards[0].file), + b"garbage", + ) + .unwrap(); let inspection = inspect_adjacency_index(dir.path()).unwrap(); assert_eq!(inspection.state, AdjacencyFreshnessState::Incompatible); @@ -2717,12 +3136,17 @@ mod tests { match case { "manifest" => std::fs::write(manifest_path(dir.path()), b"corrupt").unwrap(), "missing-union" => { - std::fs::remove_file(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out)) - .unwrap(); + std::fs::remove_file( + csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out) + .with_extension("csr.json"), + ) + .unwrap(); } "corrupt-union" => { + let path = csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out); + let reader = ShardedCsrIndex::open(&path).unwrap(); std::fs::write( - csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out), + reader.root.join(&reader.manifest.shards[0].file), b"corrupt", ) .unwrap(); @@ -2856,7 +3280,13 @@ mod tests { let dir = TempDir::new().unwrap(); write_diamond(dir.path()); build_adjacency_index(dir.path(), BUILD_TS).unwrap(); - std::fs::write(csr_path(dir.path(), "KNOWS", Direction::In), b"garbage").unwrap(); + let path = csr_path(dir.path(), "KNOWS", Direction::In); + let reader = ShardedCsrIndex::open(&path).unwrap(); + std::fs::write( + reader.root.join(&reader.manifest.shards[0].file), + b"garbage", + ) + .unwrap(); let issues = validate_adjacency_index(dir.path()).unwrap(); assert_eq!(issues.len(), 1); @@ -2872,7 +3302,10 @@ mod tests { let dir = TempDir::new().unwrap(); write_diamond(dir.path()); build_adjacency_index(dir.path(), BUILD_TS).unwrap(); - std::fs::remove_file(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out)).unwrap(); + std::fs::remove_file( + csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out).with_extension("csr.json"), + ) + .unwrap(); let issues = validate_adjacency_index(dir.path()).unwrap(); assert_eq!( @@ -3126,8 +3559,11 @@ mod tests { spill_dir: None, spill_max_bytes: None, memory_budget_bytes: None, + shard_max_edges: 2, + shard_max_nodes: 2, + merge_fan_in: 2, }; - build_adjacency_index_into_with_options( + let (_, metrics) = build_adjacency_index_into_with_metrics( dir.path(), dir.path(), BUILD_TS, @@ -3135,6 +3571,16 @@ mod tests { &mut || Ok(()), ) .unwrap(); + assert_eq!(metrics.source_rows, 6); + assert!(metrics.spill_runs >= 8); + assert!(metrics.csr_shards >= 4); + assert!(metrics.peak_shard_edges <= 2); + assert!(metrics.peak_shard_nodes <= 2); + assert!(sharded_csr_exists(&csr_path( + dir.path(), + "KNOWS", + Direction::Out + ))); assert_eq!( read_csr(&csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap(), @@ -3175,6 +3621,9 @@ mod tests { spill_dir: Some(spill.clone()), spill_max_bytes: None, memory_budget_bytes: None, + shard_max_edges: 2, + shard_max_nodes: 2, + merge_fan_in: 2, }; let mut checkpoints = 0usize; let err = build_adjacency_index_into_with_options( @@ -3233,6 +3682,9 @@ mod tests { spill_dir: None, spill_max_bytes: Some(1), // impossible for any real run file memory_budget_bytes: None, + shard_max_edges: 2, + shard_max_nodes: 2, + merge_fan_in: 2, }; let err = build_adjacency_index_into_with_options( dir.path(), diff --git a/docs/book/architecture/storage.md b/docs/book/architecture/storage.md index e8311781d..4190a03a7 100644 --- a/docs/book/architecture/storage.md +++ b/docs/book/architecture/storage.md @@ -221,12 +221,13 @@ changes results — only speed. indexes/ └── adjacency/ ├── index_manifest.parquet - ├── WORKS_AT.out.csr - ├── WORKS_AT.in.csr - ├── OWNS.out.csr - ├── OWNS.in.csr - ├── _all.out.csr # union across relation types (for via=None) - └── _all.in.csr + ├── WORKS_AT.out.csr.json # versioned shard-set manifest + ├── WORKS_AT.out.csr.shards-.d/ + │ ├── 00000000000000000000.csr + │ └── ... + ├── WORKS_AT.in.csr.json + ├── OWNS.{out,in}.csr.json + └── _all.{out,in}.csr.json # union across relation types ``` The builder (`graphforge_storage::adjacency::build_adjacency_index`) writes one `{out, in}` pair per @@ -249,14 +250,15 @@ make the result stale, never falsely fresh. | `node_count` | `UInt64` | Number of source nodes covered (CSR row count) | | `edge_count` | `UInt64` | Number of `(edge, neighbor)` entries | -**CSR file (`..csr`)** — single-batch Arrow IPC, one column, one row per -surrogate `node_id ∈ 0..node_count`: +**Sharded CSR (`..csr.json`)** — a versioned JSON manifest names an +immutable, content-addressed shard directory. Each bounded shard is Arrow IPC with one +column and covers a contiguous local surrogate range: | Column | Arrow type | Notes | |---|---|---| | `adjacency` | `LargeList` | Row `i` holds the adjacency entries of surrogate `node_id = i`, in CSR order | -This is the CSR structure in its idiomatic Arrow encoding — the two logical arrays cannot be +Within each shard this is the CSR structure in its idiomatic Arrow encoding — the two logical arrays cannot be two top-level columns because a RecordBatch requires equal column lengths. The list's offsets buffer **is** the CSR offsets array (length `node_count + 1`, `Int64`, starting at 0, monotone), and the flattened struct child **is** the targets array (length `edge_count`): @@ -267,13 +269,16 @@ Conventions: - **Empty graph**: a zero-row batch — logical `offsets == [0]`, empty targets. The offsets array is never empty. - **Node with no neighbors**: an empty list (`offsets[i] == offsets[i+1]`). -- CSR rows cover exactly `node_id ∈ 0..node_count`; surrogates beyond `node_count` simply - have no entries. +- The shard manifest records format/version, total node/edge counts, ordered boundaries, + per-shard counts, and SHA-256 checksums. A row may span consecutive shards when a + high-degree vertex exceeds the configured hard edge cap; readers concatenate those + fragments in deterministic `(key, edge_id)` order. +- Logical CSR rows cover exactly `node_id ∈ 0..node_count`; surrogates beyond `node_count` + have no entries. Empty interior rows need no physical shard bytes. - In-memory consumers (`graphforge_exec::AdjacencyProvider`) keep the logical - `offsets` / parallel `edge_ids`+`neighbor_ids` model on a persisted hit - (#340 CSR-native views); the list encoding remains a file-format detail - (`graphforge_storage::adjacency::CsrIndex`). Scan-build fallback still - materializes a hash map for oracle parity. + a `ShardedCsrIndex` on a persisted hit and materialize only the requested bounded + shard row. Legacy single-batch `.csr` files remain readable and migrate on rebuild. + Scan-build fallback still materializes a hash map for oracle parity. ### Rebuild and versioning semantics @@ -302,8 +307,16 @@ Conventions: tie-break makes the CSR bytes reproducible from `topology/` alone. `_all.{out,in}.csr` are the same sorts over the union of all typed files plus `_exploratory.parquet`. The manifest's `built_at` is excluded from the determinism guarantee. -- **Build ordering.** Builders write all CSR files first and `index_manifest.parquet` - **last**, so a torn build reads as stale (absent/old manifest), never as fresh. +- **Bounded build.** Projected Parquet batches feed sorted spill runs. Bounded-fan-in merge + passes (64 runs by default) emit rows directly into hard-capped shard sinks; they never reconstruct complete edge/neighbor + arrays. Both edge entries and local offset rows have hard shard caps. + `AdjacencyBuildMetrics` exposes source rows, spill runs/bytes, shard count, and peak + shard entries/rows for scale evidence. +- **Build ordering.** Builders write immutable shard directories, atomically publish each + shard-set manifest, and write `index_manifest.parquet` **last**. The public facade builds + in a same-filesystem private directory, validates it, then swaps the complete adjacency + directory under its visibility lock. Cancellation or failure leaves the prior directory + active and removes unpublished spill/build state. - **Loader semantics** (`graphforge_exec::PersistentAdjacencyProvider`). Freshness requires a non-empty manifest whose graph source fingerprint equals the pinned graph participant. Fresh + row present ⇒ load From 1f59a1626cbe246931d4cc36ac0dd2b3b3933020 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 08:07:18 -0600 Subject: [PATCH 4/7] docs(storage): document bounded CSR shard lifecycle --- crates/graphforge-storage/src/adjacency.rs | 15 ++++++++ docs/book/architecture/storage.md | 43 +++++++++++----------- 2 files changed, 37 insertions(+), 21 deletions(-) diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index 398363aa8..eee3074d1 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -2486,6 +2486,21 @@ mod tests { assert_eq!(reader.row(1_000_000).unwrap(), vec![(2, 0)]); } + #[test] + fn legacy_single_batch_csr_migrates_to_shards_on_rebuild() { + let directory = TempDir::new().unwrap(); + let path = directory.path().join("KNOWS.out.csr"); + let expected = sample_csr(); + + write_csr(&path, &expected).unwrap(); + assert!(!sharded_csr_exists(&path)); + assert_eq!(read_csr(&path).unwrap(), expected); + + write_sharded_csr(&path, &expected, 2).unwrap(); + assert!(sharded_csr_exists(&path)); + assert_eq!(read_csr(&path).unwrap(), expected); + } + #[test] fn entries_from_out_csr_rejects_malformed_offset_bounds() { let past_targets = CsrIndex { diff --git a/docs/book/architecture/storage.md b/docs/book/architecture/storage.md index 0e9429a79..51ac8b5b1 100644 --- a/docs/book/architecture/storage.md +++ b/docs/book/architecture/storage.md @@ -243,9 +243,9 @@ The builder (`graphforge_storage::adjacency::build_adjacency_index`) writes one relation type plus the `_all` union pair, then the manifest **last**. Relation names unusable as file stems (path separators, `..`, the reserved `_all`) are skipped — those relations are served by scan-build, but their rows still flow into the union index. The -manifest is stamped with the pinned project-generation UUID and graph source -fingerprint read **before** the edge scan. A newer publication can therefore -make the result stale, never falsely fresh. +manifest is stamped with the `topology_generation` counter read **before** the +edge scan. A concurrent topology mutation can therefore make the result stale, +never falsely fresh. **`index_manifest.parquet`** @@ -253,8 +253,7 @@ make the result stale, never falsely fresh. |---|---|---| | `relation_type` | `Utf8` | Relation type name, or `_all` for the union index | | `direction` | `Utf8` | `"out"` \| `"in"` | -| `project_generation_uuid` | `FixedSizeBinary(16)` | Committed generation pinned by the builder | -| `graph_source_fingerprint` | `FixedSizeBinary(32)` | Canonical graph participant fingerprint | +| `topology_generation` | `UInt64` | Counter pinned before the source scan | | `built_at` | `Timestamp(Microseconds, UTC)` | | | `node_count` | `UInt64` | Number of source nodes covered (CSR row count) | | `edge_count` | `UInt64` | Number of `(edge, neighbor)` entries | @@ -284,41 +283,43 @@ Conventions: fragments in deterministic `(key, edge_id)` order. - Logical CSR rows cover exactly `node_id ∈ 0..node_count`; surrogates beyond `node_count` have no entries. Empty interior rows need no physical shard bytes. -- In-memory consumers (`graphforge_exec::AdjacencyProvider`) keep the logical - a `ShardedCsrIndex` on a persisted hit and materialize only the requested bounded - shard row. Legacy single-batch `.csr` files remain readable and migrate on rebuild. +- In-memory consumers (`graphforge_exec::AdjacencyProvider`) keep a + `ShardedCsrIndex` on a persisted hit and materialize only the requested logical row + from its bounded shard fragments. Legacy single-batch `.csr` files remain readable + and migrate on rebuild. Scan-build fallback still materializes a hash map for oracle parity. ### Rebuild and versioning semantics - **Source of truth.** A CSR is always reconstructable from `topology/edges/.parquet` alone, deterministically. -- **Generation identity.** The adjacency manifest records the committed - project-generation UUID and graph source fingerprint from the reader's - pinned snapshot. There is no absent-counter or generation-zero meaning. +- **Generation identity.** The adjacency manifest records the topology counter + pinned before the source scan. A complete delta chain may advance an older + base to the current topology counter without copying the base CSR. - **Publication rule.** A graph mutation and its source fingerprint publish in the same immutable generation. `CURRENT` changes only after every participant is durable and validated. - **Crash-safety invariant.** A reader sees either the prior complete graph generation or the new complete graph generation. A failed or interrupted write never exposes a counter/data mismatch or committed prefix. -- **Staleness detection.** The provider compares the manifest's source - fingerprint with the pinned graph participant fingerprint. A corrupt - accelerator is always stale, never fresh. +- **Staleness detection.** The provider compares the manifest's topology counter + with the current counter and validates any required bounded delta chain. A + corrupt accelerator is never served as a hit. - **Fallback.** On mismatch (or absent index), the provider scans the typed edge tables and builds the adjacency in memory — yielding identical results, only slower. A stale or missing index can therefore never cause incorrect output. - **Rebuild triggers.** Lazy on first traversal when the `indexes/adjacency/` capability is - present, or explicit via `forge.index("adjacency", ...)`. Incremental rebuild - (append-delta + compaction) is deferred to v0.5.1. + present, or explicit via `forge.index("adjacency", ...)`. Append-only commits + publish bounded delta segments; a full rebuild compacts them into sharded bases. - **Determinism (R-ADJ-2).** Full rebuild streams each typed edge file once; `out` entries sort by `(src_id, edge_id)` and `in` entries by `(dst_id, edge_id)` — the `edge_id` - tie-break makes the CSR bytes reproducible from `topology/` alone. `_all.{out,in}.csr` are + tie-break makes shard bytes reproducible from `topology/` alone. `_all.{out,in}.csr.json` are the same sorts over the union of all typed files plus `_exploratory.parquet`. The manifest's `built_at` is excluded from the determinism guarantee. - **Bounded build.** Projected Parquet batches feed sorted spill runs. Bounded-fan-in merge - passes (64 runs by default) emit rows directly into hard-capped shard sinks; they never reconstruct complete edge/neighbor - arrays. Both edge entries and local offset rows have hard shard caps. + passes (64 runs by default) emit rows directly into hard-capped shard sinks; they never + reconstruct complete edge/neighbor arrays. Both edge entries and local offset rows have + hard shard caps. `AdjacencyBuildMetrics` exposes source rows, spill runs/bytes, shard count, and peak shard entries/rows for scale evidence. - **Build ordering.** Builders write immutable shard directories, atomically publish each @@ -327,8 +328,8 @@ Conventions: directory under its visibility lock. Cancellation or failure leaves the prior directory active and removes unpublished spill/build state. - **Loader semantics** (`graphforge_exec::PersistentAdjacencyProvider`). - Freshness requires a non-empty manifest whose graph source fingerprint - equals the pinned graph participant. Fresh + row present ⇒ load + Freshness requires a non-empty manifest whose topology generation is current + directly or through a complete bounded delta chain. Fresh + row present ⇒ load (`adjacency=hit`); stale or torn ⇒ lazy rebuild, then serve; fresh but **no row** for the requested relation ⇒ scan-build *without* rebuild (rebuilding cannot add an unknown relation — prevents a rebuild-per-query loop); a From bff461b95aeba03c6ff7b18d95eda591d8138e2d Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 08:19:52 -0600 Subject: [PATCH 5/7] fix(storage): harden CSR shard publication and reuse --- crates/graphforge-api/src/search_index.rs | 40 +++---- crates/graphforge-storage/src/adjacency.rs | 117 ++++++++++++++++++++- 2 files changed, 134 insertions(+), 23 deletions(-) diff --git a/crates/graphforge-api/src/search_index.rs b/crates/graphforge-api/src/search_index.rs index c2d64ab9a..b40d76fc0 100644 --- a/crates/graphforge-api/src/search_index.rs +++ b/crates/graphforge-api/src/search_index.rs @@ -360,28 +360,32 @@ impl GraphForge { .join(graphforge_storage::adjacency::ADJACENCY_SPILL_DIR_NAME), ) }; - graphforge_storage::adjacency::AdjacencyBuildOptions { + let mut options = graphforge_storage::adjacency::AdjacencyBuildOptions::default(); + options.chunk_rows = { // Scale chunk with the instance memory budget so large Explicit // policies spill fewer runs (agent-host >200M evidence). Cap keeps // a single flush allocation bounded; `effective()` still shrinks // when the budget is tighter than this starting point. - chunk_rows: { - const MAX_CHUNK_ROWS: u64 = 16_777_216; - let budget_entries = policy - .memory_budget_bytes - .saturating_div(24u64.saturating_mul(4)); - usize::try_from(budget_entries.clamp(1, MAX_CHUNK_ROWS)) - .unwrap_or(graphforge_storage::adjacency::DEFAULT_ADJACENCY_CHUNK_ROWS) - }, - batch_size: policy.batch_size.max(1), - spill_dir, - spill_max_bytes: policy.spill_max_bytes, - memory_budget_bytes: Some(policy.memory_budget_bytes), - shard_max_edges: graphforge_storage::adjacency::DEFAULT_CSR_SHARD_EDGES, - shard_max_nodes: graphforge_storage::adjacency::DEFAULT_CSR_SHARD_NODES, - merge_fan_in: graphforge_storage::adjacency::DEFAULT_ADJACENCY_MERGE_FAN_IN, - } - .effective() + const MAX_CHUNK_ROWS: u64 = 16_777_216; + let budget_entries = policy + .memory_budget_bytes + .saturating_div(24u64.saturating_mul(4)); + usize::try_from(budget_entries.clamp(1, MAX_CHUNK_ROWS)) + .unwrap_or(graphforge_storage::adjacency::DEFAULT_ADJACENCY_CHUNK_ROWS) + }; + options.batch_size = policy.batch_size.max(1); + options.spill_dir = spill_dir; + options.spill_max_bytes = policy.spill_max_bytes; + options.memory_budget_bytes = Some(policy.memory_budget_bytes); + options.shard_max_edges = { + let budget_entries = policy + .memory_budget_bytes + .saturating_div(24u64.saturating_mul(8)); + usize::try_from(budget_entries) + .unwrap_or(graphforge_storage::adjacency::DEFAULT_CSR_SHARD_EDGES) + .clamp(1, graphforge_storage::adjacency::DEFAULT_CSR_SHARD_EDGES) + }; + options.effective() } /// Inspect the derived adjacency artifact without mutating it. diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index eee3074d1..752809789 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -112,6 +112,9 @@ pub struct AdjacencyBuildMetrics { pub struct ShardedCsrIndex { root: PathBuf, manifest: CsrShardManifest, + // One decoded shard is enough to make sequential row traversal O(shards) + // while keeping reader memory bounded by the configured shard limit. + cache: std::sync::Arc>>, } impl ShardedCsrIndex { @@ -173,7 +176,11 @@ impl ShardedCsrIndex { "CSR shard manifest counts disagree".into(), )); } - Ok(Self { root, manifest }) + Ok(Self { + root, + manifest, + cache: std::sync::Arc::new(std::sync::Mutex::new(None)), + }) } /// Total logical source rows across all shards. @@ -219,6 +226,15 @@ impl ShardedCsrIndex { record: &CsrShardRecord, node_id: u64, ) -> Result, GfError> { + let mut cache = self + .cache + .lock() + .map_err(|_| GfError::Storage("CSR shard cache lock poisoned".into()))?; + if let Some((file, csr)) = cache.as_ref() + && file == &record.file + { + return Ok(csr.row(node_id - record.first_node).iter().collect()); + } let path = self.root.join(&record.file); let bytes = std::fs::read(&path).map_err(|error| { GfError::Storage(format!("missing CSR shard {}: {error}", path.display())) @@ -236,7 +252,9 @@ impl ShardedCsrIndex { record.file ))); } - Ok(csr.row(node_id - record.first_node).iter().collect()) + let output = csr.row(node_id - record.first_node).iter().collect(); + *cache = Some((record.file.clone(), csr)); + Ok(output) } } @@ -278,6 +296,7 @@ struct ShardedCsrWriter { edge_count: u64, peak_shard_edges: u64, peak_shard_nodes: u64, + owned_root: Option, finished: bool, } @@ -293,6 +312,7 @@ impl ShardedCsrWriter { std::fs::create_dir(&root).map_err(storage_err)?; Ok(Self { path: path.to_path_buf(), + owned_root: Some(root.clone()), root, shard_dir, max_edges: max_edges.max(1), @@ -384,10 +404,15 @@ impl ShardedCsrWriter { .parent() .expect("validated parent") .join(&stable_dir); - if stable_root.exists() { + if stable_root.exists() && shard_set_matches(&stable_root, &self.records) { std::fs::remove_dir_all(&self.root).map_err(storage_err)?; + self.owned_root = None; } else { + if stable_root.exists() { + std::fs::remove_dir_all(&stable_root).map_err(storage_err)?; + } std::fs::rename(&self.root, &stable_root).map_err(storage_err)?; + self.owned_root = Some(stable_root.clone()); } self.root = stable_root; self.shard_dir = stable_dir; @@ -420,12 +445,33 @@ impl ShardedCsrWriter { impl Drop for ShardedCsrWriter { fn drop(&mut self) { - if !self.finished { - let _ = std::fs::remove_dir_all(&self.root); + if !self.finished + && let Some(root) = self.owned_root.as_ref() + { + let _ = std::fs::remove_dir_all(root); } } } +fn shard_set_matches(root: &Path, records: &[CsrShardRecord]) -> bool { + for record in records { + let path = root.join(&record.file); + let Ok(bytes) = std::fs::read(&path) else { + return false; + }; + if sha256_hex(&bytes) != record.sha256 { + return false; + } + let Ok(csr) = read_csr_bytes(&bytes, &path) else { + return false; + }; + if csr.node_count() != record.node_count || csr.edge_count() != record.edge_count { + return false; + } + } + true +} + fn write_csr_shard( root: &Path, first_node: u64, @@ -803,7 +849,15 @@ pub fn write_csr(path: &Path, csr: &CsrIndex) -> Result<(), GfError> { // shadow the newly published single-batch file. let sharded_manifest = path.with_extension("csr.json"); if sharded_manifest.exists() { + let shard_root = std::fs::read(&sharded_manifest) + .ok() + .and_then(|bytes| serde_json::from_slice::(&bytes).ok()) + .filter(|manifest| Path::new(&manifest.shard_dir).components().count() == 1) + .and_then(|manifest| path.parent().map(|parent| parent.join(manifest.shard_dir))); std::fs::remove_file(sharded_manifest).map_err(storage_err)?; + if let Some(shard_root) = shard_root { + let _ = std::fs::remove_dir_all(shard_root); + } } Ok(()) } @@ -1046,6 +1100,7 @@ const BYTES_PER_KEYED_ENTRY: u64 = 24; /// project-local `.spill` root) and are removed on success, failure, or /// cancellation. #[derive(Clone, Debug, Eq, PartialEq)] +#[non_exhaustive] pub struct AdjacencyBuildOptions { /// Maximum projected edge rows retained in memory per relation/union /// accumulator before flushing a sorted spill run. @@ -2438,6 +2493,52 @@ mod tests { ); } + #[test] + fn sequential_rows_reuse_only_the_current_authenticated_shard() { + let directory = TempDir::new().unwrap(); + let path = directory.path().join("KNOWS.out.csr"); + let expected = sample_csr(); + write_sharded_csr(&path, &expected, 8).unwrap(); + let reader = ShardedCsrIndex::open(&path).unwrap(); + assert_eq!(reader.row(0).unwrap(), vec![(10, 1), (11, 2)]); + std::fs::remove_file(reader.root.join(&reader.manifest.shards[0].file)).unwrap(); + assert_eq!(reader.row(2).unwrap(), vec![(12, 0), (13, 1)]); + } + + #[test] + fn deterministic_rebuild_repairs_a_corrupt_stable_shard_set() { + let directory = TempDir::new().unwrap(); + let path = directory.path().join("KNOWS.out.csr"); + let expected = sample_csr(); + write_sharded_csr(&path, &expected, 2).unwrap(); + let reader = ShardedCsrIndex::open(&path).unwrap(); + std::fs::write( + reader.root.join(&reader.manifest.shards[0].file), + b"corrupt", + ) + .unwrap(); + + write_sharded_csr(&path, &expected, 2).unwrap(); + assert_eq!(read_csr(&path).unwrap(), expected); + } + + #[test] + fn failed_manifest_republish_does_not_delete_reused_live_shards() { + let directory = TempDir::new().unwrap(); + let path = directory.path().join("KNOWS.out.csr"); + let manifest_path = path.with_extension("csr.json"); + let saved_manifest = directory.path().join("saved-manifest.json"); + let expected = sample_csr(); + write_sharded_csr(&path, &expected, 2).unwrap(); + std::fs::rename(&manifest_path, &saved_manifest).unwrap(); + std::fs::create_dir(&manifest_path).unwrap(); + + assert!(write_sharded_csr(&path, &expected, 2).is_err()); + std::fs::remove_dir(&manifest_path).unwrap(); + std::fs::rename(&saved_manifest, &manifest_path).unwrap(); + assert_eq!(read_csr(&path).unwrap(), expected); + } + #[test] fn high_degree_row_spans_hard_capped_shards_in_order() { let directory = TempDir::new().unwrap(); @@ -2497,8 +2598,14 @@ mod tests { assert_eq!(read_csr(&path).unwrap(), expected); write_sharded_csr(&path, &expected, 2).unwrap(); + let shard_root = ShardedCsrIndex::open(&path).unwrap().root; assert!(sharded_csr_exists(&path)); assert_eq!(read_csr(&path).unwrap(), expected); + + write_csr(&path, &expected).unwrap(); + assert!(!shard_root.exists()); + assert!(!sharded_csr_exists(&path)); + assert_eq!(read_csr(&path).unwrap(), expected); } #[test] From 32902602bdebaa92137ec4192afa0e80a3fbced9 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 08:30:52 -0600 Subject: [PATCH 6/7] fix(storage): reject unsafe CSR manifest paths --- crates/graphforge-storage/src/adjacency.rs | 39 ++++++++++++++++++++-- 1 file changed, 37 insertions(+), 2 deletions(-) diff --git a/crates/graphforge-storage/src/adjacency.rs b/crates/graphforge-storage/src/adjacency.rs index 752809789..345499955 100644 --- a/crates/graphforge-storage/src/adjacency.rs +++ b/crates/graphforge-storage/src/adjacency.rs @@ -131,6 +131,9 @@ impl ShardedCsrIndex { } let mut prior_first = None; let mut edges = 0_u64; + if !is_normal_path_component(&manifest.shard_dir) { + return Err(GfError::Storage("invalid CSR shard directory name".into())); + } let root = path .parent() .unwrap_or_else(|| Path::new(".")) @@ -146,7 +149,7 @@ impl ShardedCsrIndex { } prior_first = Some(shard.first_node); edges = edges.saturating_add(shard.edge_count); - if Path::new(&shard.file).components().count() != 1 { + if !is_normal_path_component(&shard.file) { return Err(GfError::Storage("invalid CSR shard file name".into())); } // Authenticate and structurally validate every bounded shard at @@ -268,6 +271,12 @@ fn csr_artifact_exists(path: &Path) -> bool { path.is_file() || sharded_csr_exists(path) } +fn is_normal_path_component(value: &str) -> bool { + let mut components = Path::new(value).components(); + matches!(components.next(), Some(std::path::Component::Normal(_))) + && components.next().is_none() +} + /// Write a versioned checksummed shard set and publish its manifest last. /// This compatibility helper accepts an in-memory CSR; streaming builders use /// the same shard writer while producing one bounded shard at a time. @@ -852,7 +861,7 @@ pub fn write_csr(path: &Path, csr: &CsrIndex) -> Result<(), GfError> { let shard_root = std::fs::read(&sharded_manifest) .ok() .and_then(|bytes| serde_json::from_slice::(&bytes).ok()) - .filter(|manifest| Path::new(&manifest.shard_dir).components().count() == 1) + .filter(|manifest| is_normal_path_component(&manifest.shard_dir)) .and_then(|manifest| path.parent().map(|parent| parent.join(manifest.shard_dir))); std::fs::remove_file(sharded_manifest).map_err(storage_err)?; if let Some(shard_root) = shard_root { @@ -2608,6 +2617,32 @@ mod tests { assert_eq!(read_csr(&path).unwrap(), expected); } + #[test] + fn legacy_cleanup_rejects_parent_directory_manifest_paths() { + let directory = TempDir::new().unwrap(); + let adjacency = directory.path().join("adjacency"); + std::fs::create_dir(&adjacency).unwrap(); + let sentinel = directory.path().join("must-remain"); + std::fs::write(&sentinel, b"safe").unwrap(); + let path = adjacency.join("KNOWS.out.csr"); + let malicious = CsrShardManifest { + format: "graphforge.csr-shards".into(), + version: SHARDED_CSR_VERSION, + node_count: 0, + edge_count: 0, + shard_dir: "..".into(), + shards: Vec::new(), + }; + std::fs::write( + path.with_extension("csr.json"), + serde_json::to_vec(&malicious).unwrap(), + ) + .unwrap(); + + write_csr(&path, &sample_csr()).unwrap(); + assert_eq!(std::fs::read(&sentinel).unwrap(), b"safe"); + } + #[test] fn entries_from_out_csr_rejects_malformed_offset_bounds() { let past_targets = CsrIndex { From 701e0869243f5e3a9c411518edc5dd08b8980a9f Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Wed, 19 Aug 2026 08:45:04 -0600 Subject: [PATCH 7/7] test(exec): exercise missing sharded CSR recovery --- .../tests/persistent_adjacency.rs | 19 ++++++------------- 1 file changed, 6 insertions(+), 13 deletions(-) diff --git a/crates/graphforge-exec/tests/persistent_adjacency.rs b/crates/graphforge-exec/tests/persistent_adjacency.rs index 7669f2914..15d47c6aa 100644 --- a/crates/graphforge-exec/tests/persistent_adjacency.rs +++ b/crates/graphforge-exec/tests/persistent_adjacency.rs @@ -236,16 +236,16 @@ fn corrupt_generation_counter_is_miss_without_rebuild() { } #[test] -fn missing_csr_file_with_fresh_manifest_rebuilds_and_repairs() { +fn missing_csr_artifact_with_fresh_manifest_rebuilds_and_repairs() { let dir = TempDir::new().unwrap(); write_diamond(dir.path()); build_adjacency_index(dir.path(), TS).unwrap(); - std::fs::remove_file(graphforge_storage::adjacency::csr_path( + let csr_path = graphforge_storage::adjacency::csr_path( dir.path(), "KNOWS", graphforge_storage::adjacency::Direction::Out, - )) - .unwrap(); + ); + std::fs::remove_file(csr_path.with_extension("csr.json")).unwrap(); let provider = persistent(dir.path(), OntologyMode::Strict); assert_eq!( @@ -259,15 +259,8 @@ fn missing_csr_file_with_fresh_manifest_rebuilds_and_repairs() { .adjacency("KNOWS", Direction::Out) .unwrap() ); - // The lazy rebuild restored the file. - assert!( - graphforge_storage::adjacency::csr_path( - dir.path(), - "KNOWS", - graphforge_storage::adjacency::Direction::Out - ) - .exists() - ); + // The lazy rebuild restored the sharded artifact. + assert!(graphforge_storage::adjacency::sharded_csr_exists(&csr_path)); } // ---------------------------------------------------------------------------