1use std::fs::File;
36use std::path::{Path, PathBuf};
37use std::sync::Arc;
38
39use arrow::array::{
40 Array, LargeListArray, RecordBatch, StringArray, StructArray, TimestampMicrosecondArray,
41 UInt64Array,
42};
43use arrow::buffer::{OffsetBuffer, ScalarBuffer};
44use arrow::datatypes::{DataType, Field};
45use arrow::ipc::reader::FileReader;
46use arrow::ipc::writer::FileWriter;
47use sha2::{Digest, Sha256};
48use tempfile::NamedTempFile;
49
50use graphforge_core::GfError;
51
52use crate::schemas::{ADJACENCY_CSR_SCHEMA, ADJACENCY_MANIFEST_SCHEMA, adjacency_entry_fields};
53use crate::staging::RewriteBatch;
54
55pub const ALL_RELATIONS_STEM: &str = "_all";
59
60pub const MANIFEST_FILE: &str = "index_manifest.parquet";
62
63fn storage_err(e: impl std::fmt::Display) -> GfError {
64 GfError::Storage(e.to_string())
65}
66
67#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
74pub enum Direction {
75 Out,
77 In,
79}
80
81impl Direction {
82 #[must_use]
85 pub fn as_str(self) -> &'static str {
86 match self {
87 Self::Out => "out",
88 Self::In => "in",
89 }
90 }
91
92 pub fn parse(s: &str) -> Result<Self, GfError> {
97 match s {
98 "out" => Ok(Self::Out),
99 "in" => Ok(Self::In),
100 other => Err(GfError::Storage(format!(
101 "invalid adjacency direction {other:?} (expected \"out\" or \"in\")"
102 ))),
103 }
104 }
105}
106
107#[derive(Clone, Debug, Default, PartialEq, Eq)]
120pub struct CsrIndex {
121 pub offsets: Vec<u64>,
123 pub edge_ids: Vec<u64>,
125 pub neighbor_ids: Vec<u64>,
127}
128
129impl CsrIndex {
130 #[must_use]
132 pub fn node_count(&self) -> u64 {
133 (self.offsets.len().max(1) - 1) as u64
134 }
135
136 #[must_use]
138 pub fn edge_count(&self) -> u64 {
139 self.edge_ids.len() as u64
140 }
141
142 fn validate(&self) -> Result<(), GfError> {
144 if self.offsets.first() != Some(&0) {
145 return Err(GfError::Storage(format!(
146 "invalid CSR: offsets must start with 0 (got {:?})",
147 self.offsets.first()
148 )));
149 }
150 if self.offsets.windows(2).any(|w| w[0] > w[1]) {
151 return Err(GfError::Storage(
152 "invalid CSR: offsets must be monotonically non-decreasing".to_owned(),
153 ));
154 }
155 let last = *self.offsets.last().unwrap_or(&0);
156 if last != self.edge_count() || self.edge_ids.len() != self.neighbor_ids.len() {
157 return Err(GfError::Storage(format!(
158 "invalid CSR: final offset {last} must equal target lengths \
159 (edge_ids: {}, neighbor_ids: {})",
160 self.edge_ids.len(),
161 self.neighbor_ids.len()
162 )));
163 }
164 Ok(())
165 }
166}
167
168#[derive(Clone, Debug, PartialEq, Eq)]
171pub struct AdjacencyManifestRow {
172 pub relation_type: String,
174 pub direction: Direction,
176 pub topology_generation: u64,
178 pub built_at_micros: i64,
181 pub node_count: u64,
183 pub edge_count: u64,
185}
186
187#[derive(Clone, Copy, Debug, PartialEq, Eq)]
189pub enum AdjacencyFreshnessState {
190 Current,
192 Missing,
194 Stale,
196 Incompatible,
198}
199
200impl AdjacencyFreshnessState {
201 #[must_use]
203 pub const fn as_str(self) -> &'static str {
204 match self {
205 Self::Current => "current",
206 Self::Missing => "missing",
207 Self::Stale => "stale",
208 Self::Incompatible => "incompatible",
209 }
210 }
211}
212
213#[derive(Clone, Copy, Debug, PartialEq, Eq)]
215pub enum AdjacencyFreshnessReason {
216 NotBuilt,
218 MixedArtifactGeneration,
220 IncompleteDeltaChain,
222 MissingCsr,
224 UnreadableArtifact,
226 ContentMismatch,
228 FutureArtifactGeneration,
230}
231
232impl AdjacencyFreshnessReason {
233 #[must_use]
235 pub const fn as_str(self) -> &'static str {
236 match self {
237 Self::NotBuilt => "not_built",
238 Self::MixedArtifactGeneration => "mixed_artifact_generation",
239 Self::IncompleteDeltaChain => "incomplete_delta_chain",
240 Self::MissingCsr => "missing_csr",
241 Self::UnreadableArtifact => "unreadable_artifact",
242 Self::ContentMismatch => "content_mismatch",
243 Self::FutureArtifactGeneration => "future_artifact_generation",
244 }
245 }
246}
247
248#[derive(Clone, Debug, PartialEq, Eq)]
250pub struct AdjacencyInspection {
251 pub source_generation: u64,
253 pub source_fingerprint: String,
255 pub artifact_generation: Option<u64>,
257 pub artifact_effective_generation: Option<u64>,
259 pub artifact_fingerprint: Option<String>,
261 pub state: AdjacencyFreshnessState,
263 pub reason: Option<AdjacencyFreshnessReason>,
265}
266
267#[must_use]
269pub fn adjacency_dir(project_dir: &Path) -> PathBuf {
270 project_dir.join("indexes").join("adjacency")
271}
272
273#[must_use]
276pub fn csr_path(project_dir: &Path, relation_type: &str, direction: Direction) -> PathBuf {
277 adjacency_dir(project_dir).join(format!("{relation_type}.{}.csr", direction.as_str()))
278}
279
280#[must_use]
282pub fn manifest_path(project_dir: &Path) -> PathBuf {
283 adjacency_dir(project_dir).join(MANIFEST_FILE)
284}
285
286pub fn write_csr(path: &Path, csr: &CsrIndex) -> Result<(), GfError> {
295 csr.validate()?;
296
297 let offsets: Vec<i64> = csr
298 .offsets
299 .iter()
300 .map(|&o| i64::try_from(o).map_err(storage_err))
301 .collect::<Result<_, _>>()?;
302 let entries = StructArray::new(
303 adjacency_entry_fields(),
304 vec![
305 Arc::new(UInt64Array::from(csr.edge_ids.clone())),
306 Arc::new(UInt64Array::from(csr.neighbor_ids.clone())),
307 ],
308 None,
309 );
310 let item_field = Arc::new(Field::new(
311 "item",
312 DataType::Struct(adjacency_entry_fields()),
313 false,
314 ));
315 let adjacency = LargeListArray::new(
316 item_field,
317 OffsetBuffer::new(ScalarBuffer::from(offsets)),
318 Arc::new(entries),
319 None,
320 );
321 let batch = RecordBatch::try_new(Arc::clone(&ADJACENCY_CSR_SCHEMA), vec![Arc::new(adjacency)])
322 .map_err(storage_err)?;
323
324 let parent = path.parent().ok_or_else(|| {
325 GfError::Storage(format!(
326 "CSR path {} has no parent directory",
327 path.display()
328 ))
329 })?;
330 std::fs::create_dir_all(parent).map_err(storage_err)?;
331 let file_name = path
332 .file_name()
333 .map_or_else(|| "csr".to_owned(), |n| n.to_string_lossy().into_owned());
334 let tmp = tempfile::Builder::new()
335 .prefix(&format!("{file_name}."))
336 .suffix(".tmp")
337 .tempfile_in(parent)
338 .map_err(storage_err)?;
339
340 let mut writer =
341 FileWriter::try_new(tmp.as_file(), &ADJACENCY_CSR_SCHEMA).map_err(storage_err)?;
342 writer.write(&batch).map_err(storage_err)?;
343 writer.finish().map_err(storage_err)?;
344 persist_temp(tmp, path)
345}
346
347pub fn read_csr(path: &Path) -> Result<CsrIndex, GfError> {
353 let file = File::open(path)
354 .map_err(|e| GfError::Storage(format!("cannot open CSR file {}: {e}", path.display())))?;
355 let reader = FileReader::try_new(file, None)
356 .map_err(|e| GfError::Storage(format!("invalid CSR file {}: {e}", path.display())))?;
357 if reader.schema().fields() != ADJACENCY_CSR_SCHEMA.fields() {
358 return Err(GfError::Storage(format!(
359 "CSR file {} has unexpected schema {:?}",
360 path.display(),
361 reader.schema()
362 )));
363 }
364
365 let mut csr = CsrIndex {
366 offsets: vec![0],
367 ..CsrIndex::default()
368 };
369 for batch in reader {
370 let batch = batch.map_err(storage_err)?;
371 let adjacency = batch
372 .column(0)
373 .as_any()
374 .downcast_ref::<LargeListArray>()
375 .ok_or_else(|| {
376 GfError::Storage(format!(
377 "CSR file {}: adjacency column is not a LargeList",
378 path.display()
379 ))
380 })?;
381 let entries = adjacency
382 .values()
383 .as_any()
384 .downcast_ref::<StructArray>()
385 .ok_or_else(|| {
386 GfError::Storage(format!(
387 "CSR file {}: adjacency entries are not a Struct",
388 path.display()
389 ))
390 })?;
391 let (edge_ids, neighbor_ids) = (
392 uint64_column(entries.column(0), "edge_id")?,
393 uint64_column(entries.column(1), "neighbor_id")?,
394 );
395 let value_offsets = adjacency.value_offsets();
399 for row in 0..adjacency.len() {
400 let start = usize::try_from(value_offsets[row]).map_err(storage_err)?;
401 let end = usize::try_from(value_offsets[row + 1]).map_err(storage_err)?;
402 for entry in start..end {
403 csr.edge_ids.push(edge_ids.value(entry));
404 csr.neighbor_ids.push(neighbor_ids.value(entry));
405 }
406 csr.offsets.push(csr.edge_count());
407 }
408 }
409 csr.validate()?;
410 Ok(csr)
411}
412
413pub fn write_manifest(project_dir: &Path, rows: &[AdjacencyManifestRow]) -> Result<(), GfError> {
422 let relation_types: StringArray = rows
423 .iter()
424 .map(|r| Some(r.relation_type.as_str()))
425 .collect();
426 let directions: StringArray = rows.iter().map(|r| Some(r.direction.as_str())).collect();
427 let generations: Vec<u64> = rows.iter().map(|r| r.topology_generation).collect();
428 let built_ats: Vec<i64> = rows.iter().map(|r| r.built_at_micros).collect();
429 let node_counts: Vec<u64> = rows.iter().map(|r| r.node_count).collect();
430 let edge_counts: Vec<u64> = rows.iter().map(|r| r.edge_count).collect();
431 let batch = RecordBatch::try_new(
432 Arc::clone(&ADJACENCY_MANIFEST_SCHEMA),
433 vec![
434 Arc::new(relation_types),
435 Arc::new(directions),
436 Arc::new(UInt64Array::from(generations)),
437 Arc::new(TimestampMicrosecondArray::from(built_ats).with_timezone("UTC")),
438 Arc::new(UInt64Array::from(node_counts)),
439 Arc::new(UInt64Array::from(edge_counts)),
440 ],
441 )
442 .map_err(storage_err)?;
443
444 let mut staged = RewriteBatch::new();
445 staged.stage(
446 &manifest_path(project_dir),
447 Arc::clone(&ADJACENCY_MANIFEST_SCHEMA),
448 &batch,
449 )?;
450 staged.commit()
451}
452
453pub fn read_manifest(project_dir: &Path) -> Result<Vec<AdjacencyManifestRow>, GfError> {
462 let path = manifest_path(project_dir);
463 if !path.exists() {
464 return Ok(Vec::new());
465 }
466 let file = File::open(&path).map_err(storage_err)?;
467 let reader = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(file)
468 .map_err(storage_err)?
469 .build()
470 .map_err(storage_err)?;
471
472 let mut rows = Vec::new();
473 for batch in reader {
474 let batch = batch.map_err(storage_err)?;
475 if batch.schema().fields() != ADJACENCY_MANIFEST_SCHEMA.fields() {
476 return Err(GfError::Storage(format!(
477 "adjacency manifest {} has unexpected schema {:?}",
478 path.display(),
479 batch.schema()
480 )));
481 }
482 let relation_types = string_column(batch.column(0), "relation_type")?;
483 let directions = string_column(batch.column(1), "direction")?;
484 let generations = uint64_column(batch.column(2), "topology_generation")?;
485 let built_ats = batch
486 .column(3)
487 .as_any()
488 .downcast_ref::<TimestampMicrosecondArray>()
489 .ok_or_else(|| {
490 GfError::Storage("adjacency manifest: built_at is not a timestamp".to_owned())
491 })?;
492 let node_counts = uint64_column(batch.column(4), "node_count")?;
493 let edge_counts = uint64_column(batch.column(5), "edge_count")?;
494 for i in 0..batch.num_rows() {
495 rows.push(AdjacencyManifestRow {
496 relation_type: relation_types.value(i).to_owned(),
497 direction: Direction::parse(directions.value(i))?,
498 topology_generation: generations.value(i),
499 built_at_micros: built_ats.value(i),
500 node_count: node_counts.value(i),
501 edge_count: edge_counts.value(i),
502 });
503 }
504 }
505 Ok(rows)
506}
507
508pub(crate) type BuildEntry = (u64, u64, u64);
515
516pub fn build_adjacency_index(
547 project_dir: &Path,
548 built_at_micros: i64,
549) -> Result<Vec<AdjacencyManifestRow>, GfError> {
550 build_adjacency_index_with_checkpoint(project_dir, built_at_micros, || Ok(()))
551}
552
553pub fn build_adjacency_index_with_checkpoint(
555 project_dir: &Path,
556 built_at_micros: i64,
557 mut checkpoint: impl FnMut() -> Result<(), GfError>,
558) -> Result<Vec<AdjacencyManifestRow>, GfError> {
559 build_adjacency_index_into(project_dir, project_dir, built_at_micros, &mut checkpoint)
560}
561
562pub fn build_adjacency_index_into(
565 source_project_dir: &Path,
566 artifact_project_dir: &Path,
567 built_at_micros: i64,
568 mut checkpoint: impl FnMut() -> Result<(), GfError>,
569) -> Result<Vec<AdjacencyManifestRow>, GfError> {
570 checkpoint()?;
571 let generation = crate::generation::read_topology_generation(source_project_dir)?;
573 let (groups, union_out) = collect_adjacency_groups(source_project_dir)?;
574
575 let adjacency = adjacency_dir(artifact_project_dir);
576 std::fs::create_dir_all(&adjacency).map_err(storage_err)?;
577
578 let mut manifest = Vec::new();
579 {
580 let mut write_pair = |stem: &str, entries: &[BuildEntry]| -> Result<(), GfError> {
581 for direction in [Direction::Out, Direction::In] {
582 checkpoint()?;
583 let csr = csr_from_entries(entries, direction);
584 write_csr(&csr_path(artifact_project_dir, stem, direction), &csr)?;
585 manifest.push(AdjacencyManifestRow {
586 relation_type: stem.to_owned(),
587 direction,
588 topology_generation: generation,
589 built_at_micros,
590 node_count: csr.node_count(),
591 edge_count: csr.edge_count(),
592 });
593 }
594 Ok(())
595 };
596 for (stem, entries) in &groups {
597 write_pair(stem, entries)?;
598 }
599 write_pair(ALL_RELATIONS_STEM, &union_out)?;
600 }
601 checkpoint()?;
602
603 write_manifest(artifact_project_dir, &manifest)?;
606 checkpoint()?;
607
608 crate::adjacency_delta::prune_delta_segments(artifact_project_dir, generation);
614 Ok(manifest)
615}
616
617#[allow(clippy::type_complexity)]
623fn collect_adjacency_groups(
624 project_dir: &Path,
625) -> Result<
626 (
627 std::collections::BTreeMap<String, Vec<BuildEntry>>,
628 Vec<BuildEntry>,
629 ),
630 GfError,
631> {
632 use std::collections::BTreeMap;
633
634 let mut groups: BTreeMap<String, Vec<BuildEntry>> = BTreeMap::new();
635 let mut union_out: Vec<BuildEntry> = Vec::new();
636 for path in crate::mutator::parquet_files_in(project_dir, "topology/edges")? {
637 let Some(stem) = path.file_stem().and_then(|s| s.to_str()).map(str::to_owned) else {
638 continue;
639 };
640 let schema = crate::catalog::discover_parquet_schema(&path).ok_or_else(|| {
644 GfError::Storage(format!(
645 "adjacency build: cannot read parquet schema for {}",
646 path.display()
647 ))
648 })?;
649 let batches = crate::catalog::read_parquet_or_empty(&path, schema).map_err(storage_err)?;
650 let exploratory = stem == "_exploratory";
651 for batch in &batches {
652 let edge_ids = uint64_column(named_column(batch, "edge_id")?, "edge_id")?;
653 let src_ids = uint64_column(named_column(batch, "src_id")?, "src_id")?;
654 let dst_ids = uint64_column(named_column(batch, "dst_id")?, "dst_id")?;
655 let rel_names = if exploratory {
656 Some(string_column(
657 named_column(batch, "rel_type_name")?,
658 "rel_type_name",
659 )?)
660 } else {
661 None
662 };
663 for i in 0..batch.num_rows() {
664 let entry = (src_ids.value(i), edge_ids.value(i), dst_ids.value(i));
665 union_out.push(entry);
666 let rel = rel_names.map_or(stem.as_str(), |names| names.value(i));
667 if usable_stem(rel) {
668 groups.entry(rel.to_owned()).or_default().push(entry);
669 }
670 }
671 }
672 }
673 Ok((groups, union_out))
674}
675
676#[derive(Clone, Debug, PartialEq, Eq)]
682pub enum AdjacencyValidationIssue {
683 StaleGeneration {
686 manifest: u64,
688 current: u64,
690 },
691 MissingCsr {
693 rel: String,
695 direction: Direction,
697 },
698 UnreadableCsr {
701 rel: String,
703 direction: Direction,
705 error: String,
707 },
708 Mismatch {
711 rel: String,
713 direction: Direction,
715 },
716}
717
718pub fn validate_adjacency_index(
734 project_dir: &Path,
735) -> Result<Vec<AdjacencyValidationIssue>, GfError> {
736 validate_adjacency_index_against(project_dir, project_dir)
737}
738
739pub fn validate_adjacency_index_against(
741 source_project_dir: &Path,
742 artifact_project_dir: &Path,
743) -> Result<Vec<AdjacencyValidationIssue>, GfError> {
744 let manifest = read_manifest(artifact_project_dir)?;
745 if manifest.is_empty() {
746 return Ok(Vec::new()); }
748
749 let mut issues = Vec::new();
750 let current = crate::generation::read_topology_generation(source_project_dir)?;
751
752 let base = manifest.first().map(|r| r.topology_generation);
759 let uniform = base.is_some_and(|b| manifest.iter().all(|r| r.topology_generation == b));
760 let chain = match base {
761 Some(b) if uniform && b < current => {
762 crate::adjacency_delta::read_delta_chain(artifact_project_dir, b, current)
763 }
764 _ => None,
765 };
766 let delta_covered = chain.is_some();
767 let chain = chain.unwrap_or_default();
768
769 if !delta_covered
770 && let Some(stale) = manifest
771 .iter()
772 .find(|r| r.topology_generation != current)
773 .map(|r| r.topology_generation)
774 {
775 issues.push(AdjacencyValidationIssue::StaleGeneration {
776 manifest: stale,
777 current,
778 });
779 }
780
781 let (groups, union_out) = collect_adjacency_groups(source_project_dir)?;
782 for row in &manifest {
783 let expected_entries: &[BuildEntry] = if row.relation_type == ALL_RELATIONS_STEM {
784 &union_out
785 } else {
786 groups
787 .get(&row.relation_type)
788 .map_or(&[][..], Vec::as_slice)
789 };
790 let expected = csr_from_entries(expected_entries, row.direction);
791 let path = csr_path(artifact_project_dir, &row.relation_type, row.direction);
792 if !path.exists() {
793 issues.push(AdjacencyValidationIssue::MissingCsr {
794 rel: row.relation_type.clone(),
795 direction: row.direction,
796 });
797 continue;
798 }
799 match read_csr(&path) {
800 Ok(base_csr) => {
803 let actual = if delta_covered {
804 crate::adjacency_delta::apply_delta_segments(
805 &base_csr,
806 &row.relation_type,
807 row.direction,
808 &chain,
809 )
810 } else {
811 base_csr
812 };
813 if actual != expected {
814 issues.push(AdjacencyValidationIssue::Mismatch {
815 rel: row.relation_type.clone(),
816 direction: row.direction,
817 });
818 }
819 }
820 Err(e) => issues.push(AdjacencyValidationIssue::UnreadableCsr {
821 rel: row.relation_type.clone(),
822 direction: row.direction,
823 error: e.to_string(),
824 }),
825 }
826 }
827 Ok(issues)
828}
829
830pub fn inspect_adjacency_index(project_dir: &Path) -> Result<AdjacencyInspection, GfError> {
836 let source_generation = crate::generation::read_topology_generation(project_dir)?;
837 let (_, mut source_entries) = collect_adjacency_groups(project_dir)?;
838 let source_fingerprint = entries_fingerprint(&mut source_entries);
839 let Ok(manifest) = read_manifest(project_dir) else {
840 return Ok(inspection_without_artifact(
841 source_generation,
842 source_fingerprint,
843 AdjacencyFreshnessState::Incompatible,
844 AdjacencyFreshnessReason::UnreadableArtifact,
845 ));
846 };
847 if manifest.is_empty() {
848 let built = manifest_path(project_dir).exists();
849 return Ok(inspection_without_artifact(
850 source_generation,
851 source_fingerprint,
852 if built {
853 AdjacencyFreshnessState::Incompatible
854 } else {
855 AdjacencyFreshnessState::Missing
856 },
857 if built {
858 AdjacencyFreshnessReason::UnreadableArtifact
859 } else {
860 AdjacencyFreshnessReason::NotBuilt
861 },
862 ));
863 }
864 let base = manifest[0].topology_generation;
865 if manifest.iter().any(|row| row.topology_generation != base) {
866 return Ok(AdjacencyInspection {
867 source_generation,
868 source_fingerprint,
869 artifact_generation: None,
870 artifact_effective_generation: None,
871 artifact_fingerprint: None,
872 state: AdjacencyFreshnessState::Incompatible,
873 reason: Some(AdjacencyFreshnessReason::MixedArtifactGeneration),
874 });
875 }
876 if base > source_generation {
877 return Ok(AdjacencyInspection {
878 source_generation,
879 source_fingerprint,
880 artifact_generation: Some(base),
881 artifact_effective_generation: None,
882 artifact_fingerprint: None,
883 state: AdjacencyFreshnessState::Incompatible,
884 reason: Some(AdjacencyFreshnessReason::FutureArtifactGeneration),
885 });
886 }
887 let chain = if base < source_generation {
888 match crate::adjacency_delta::read_delta_chain(project_dir, base, source_generation) {
889 Some(chain) => chain,
890 None => {
891 return Ok(AdjacencyInspection {
892 source_generation,
893 source_fingerprint,
894 artifact_generation: Some(base),
895 artifact_effective_generation: None,
896 artifact_fingerprint: None,
897 state: AdjacencyFreshnessState::Stale,
898 reason: Some(AdjacencyFreshnessReason::IncompleteDeltaChain),
899 });
900 }
901 }
902 } else {
903 Vec::new()
904 };
905 let union_path = csr_path(project_dir, ALL_RELATIONS_STEM, Direction::Out);
906 if !union_path.exists() {
907 return Ok(AdjacencyInspection {
908 source_generation,
909 source_fingerprint,
910 artifact_generation: Some(base),
911 artifact_effective_generation: Some(source_generation),
912 artifact_fingerprint: None,
913 state: AdjacencyFreshnessState::Incompatible,
914 reason: Some(AdjacencyFreshnessReason::MissingCsr),
915 });
916 }
917 let Ok(base_csr) = read_csr(&union_path) else {
918 return Ok(AdjacencyInspection {
919 source_generation,
920 source_fingerprint,
921 artifact_generation: Some(base),
922 artifact_effective_generation: Some(source_generation),
923 artifact_fingerprint: None,
924 state: AdjacencyFreshnessState::Incompatible,
925 reason: Some(AdjacencyFreshnessReason::UnreadableArtifact),
926 });
927 };
928 inspect_effective_artifact(
929 project_dir,
930 source_generation,
931 source_fingerprint,
932 base,
933 &base_csr,
934 &chain,
935 )
936}
937
938fn inspect_effective_artifact(
939 project_dir: &Path,
940 source_generation: u64,
941 source_fingerprint: String,
942 base: u64,
943 base_csr: &CsrIndex,
944 chain: &[crate::adjacency_delta::DeltaSegment],
945) -> Result<AdjacencyInspection, GfError> {
946 let effective = crate::adjacency_delta::apply_delta_segments(
947 base_csr,
948 ALL_RELATIONS_STEM,
949 Direction::Out,
950 chain,
951 );
952 let Some(mut artifact_entries) = entries_from_out_csr(&effective) else {
953 return Ok(AdjacencyInspection {
954 source_generation,
955 source_fingerprint,
956 artifact_generation: Some(base),
957 artifact_effective_generation: Some(source_generation),
958 artifact_fingerprint: None,
959 state: AdjacencyFreshnessState::Incompatible,
960 reason: Some(AdjacencyFreshnessReason::UnreadableArtifact),
961 });
962 };
963 let artifact_fingerprint = entries_fingerprint(&mut artifact_entries);
964 let validation_issues = validate_adjacency_index(project_dir)?;
965 let validation_reason = validation_issues.first().map(|issue| match issue {
966 AdjacencyValidationIssue::StaleGeneration { .. } => {
967 AdjacencyFreshnessReason::IncompleteDeltaChain
968 }
969 AdjacencyValidationIssue::MissingCsr { .. } => AdjacencyFreshnessReason::MissingCsr,
970 AdjacencyValidationIssue::UnreadableCsr { .. } => {
971 AdjacencyFreshnessReason::UnreadableArtifact
972 }
973 AdjacencyValidationIssue::Mismatch { .. } => AdjacencyFreshnessReason::ContentMismatch,
974 });
975 let (state, reason) =
976 if artifact_fingerprint == source_fingerprint && validation_reason.is_none() {
977 (AdjacencyFreshnessState::Current, None)
978 } else {
979 (
980 AdjacencyFreshnessState::Incompatible,
981 Some(validation_reason.unwrap_or(AdjacencyFreshnessReason::ContentMismatch)),
982 )
983 };
984 Ok(AdjacencyInspection {
985 source_generation,
986 source_fingerprint,
987 artifact_generation: Some(base),
988 artifact_effective_generation: Some(source_generation),
989 artifact_fingerprint: Some(artifact_fingerprint),
990 state,
991 reason,
992 })
993}
994
995fn inspection_without_artifact(
996 source_generation: u64,
997 source_fingerprint: String,
998 state: AdjacencyFreshnessState,
999 reason: AdjacencyFreshnessReason,
1000) -> AdjacencyInspection {
1001 AdjacencyInspection {
1002 source_generation,
1003 source_fingerprint,
1004 artifact_generation: None,
1005 artifact_effective_generation: None,
1006 artifact_fingerprint: None,
1007 state,
1008 reason: Some(reason),
1009 }
1010}
1011
1012fn entries_from_out_csr(csr: &CsrIndex) -> Option<Vec<BuildEntry>> {
1013 let mut entries = Vec::with_capacity(csr.edge_ids.len());
1014 for (src_index, offsets) in csr.offsets.windows(2).enumerate() {
1015 let src = u64::try_from(src_index).ok()?;
1016 let start = usize::try_from(offsets[0]).ok()?;
1017 let end = usize::try_from(offsets[1]).ok()?;
1018 if start > end || end > csr.edge_ids.len() || end > csr.neighbor_ids.len() {
1019 return None;
1020 }
1021 for index in start..end {
1022 entries.push((src, csr.edge_ids[index], csr.neighbor_ids[index]));
1023 }
1024 }
1025 Some(entries)
1026}
1027
1028fn entries_fingerprint(entries: &mut [BuildEntry]) -> String {
1029 entries.sort_unstable();
1030 let mut digest = Sha256::new();
1031 digest.update(b"graphforge/adjacency-topology/v1\0");
1032 for &(src, edge, dst) in entries.iter() {
1033 digest.update(src.to_le_bytes());
1034 digest.update(edge.to_le_bytes());
1035 digest.update(dst.to_le_bytes());
1036 }
1037 format!("sha256:{:x}", digest.finalize())
1038}
1039
1040pub(crate) fn csr_from_entries(entries: &[BuildEntry], direction: Direction) -> CsrIndex {
1045 let mut keyed: Vec<(u64, u64, u64)> = entries
1046 .iter()
1047 .map(|&(src, edge, dst)| match direction {
1048 Direction::Out => (src, edge, dst),
1049 Direction::In => (dst, edge, src),
1050 })
1051 .collect();
1052 keyed.sort_unstable_by_key(|&(key, edge, _)| (key, edge));
1053
1054 let node_count = keyed.last().map_or(0, |&(key, _, _)| key + 1);
1055 let mut csr = CsrIndex {
1056 offsets: Vec::with_capacity(usize::try_from(node_count).unwrap_or(0) + 1),
1057 edge_ids: Vec::with_capacity(keyed.len()),
1058 neighbor_ids: Vec::with_capacity(keyed.len()),
1059 };
1060 csr.offsets.push(0);
1061 let mut next = 0usize; for node in 0..node_count {
1063 while next < keyed.len() && keyed[next].0 == node {
1064 csr.edge_ids.push(keyed[next].1);
1065 csr.neighbor_ids.push(keyed[next].2);
1066 next += 1;
1067 }
1068 csr.offsets.push(csr.edge_count());
1069 }
1070 csr
1071}
1072
1073pub(crate) fn usable_stem(rel: &str) -> bool {
1077 if rel == ALL_RELATIONS_STEM {
1078 return false;
1079 }
1080 let mut comps = Path::new(rel).components();
1081 matches!(comps.next(), Some(std::path::Component::Normal(_))) && comps.next().is_none()
1082}
1083
1084fn named_column<'a>(
1086 batch: &'a arrow::record_batch::RecordBatch,
1087 name: &str,
1088) -> Result<&'a arrow::array::ArrayRef, GfError> {
1089 batch
1090 .column_by_name(name)
1091 .ok_or_else(|| GfError::Storage(format!("adjacency build: missing column {name}")))
1092}
1093
1094fn persist_temp(tmp: NamedTempFile, path: &Path) -> Result<(), GfError> {
1096 tmp.persist(path)
1097 .map(|_| ())
1098 .map_err(|e| storage_err(e.error))
1099}
1100
1101fn uint64_column<'a>(
1102 column: &'a arrow::array::ArrayRef,
1103 name: &str,
1104) -> Result<&'a UInt64Array, GfError> {
1105 column
1106 .as_any()
1107 .downcast_ref::<UInt64Array>()
1108 .ok_or_else(|| GfError::Storage(format!("adjacency: {name} column is not UInt64")))
1109}
1110
1111fn string_column<'a>(
1112 column: &'a arrow::array::ArrayRef,
1113 name: &str,
1114) -> Result<&'a StringArray, GfError> {
1115 column
1116 .as_any()
1117 .downcast_ref::<StringArray>()
1118 .ok_or_else(|| GfError::Storage(format!("adjacency: {name} column is not Utf8")))
1119}
1120
1121#[cfg(test)]
1126mod tests {
1127 use tempfile::TempDir;
1128
1129 use super::*;
1130
1131 fn sample_csr() -> CsrIndex {
1133 CsrIndex {
1134 offsets: vec![0, 2, 2, 4],
1135 edge_ids: vec![10, 11, 12, 13],
1136 neighbor_ids: vec![1, 2, 0, 1],
1137 }
1138 }
1139
1140 #[test]
1141 fn entries_from_out_csr_rejects_malformed_offset_bounds() {
1142 let past_targets = CsrIndex {
1143 offsets: vec![0, 2],
1144 edge_ids: vec![7],
1145 neighbor_ids: vec![9],
1146 };
1147 assert!(entries_from_out_csr(&past_targets).is_none());
1148
1149 let descending = CsrIndex {
1150 offsets: vec![1, 0],
1151 edge_ids: vec![7],
1152 neighbor_ids: vec![9],
1153 };
1154 assert!(entries_from_out_csr(&descending).is_none());
1155
1156 let mismatched_columns = CsrIndex {
1157 offsets: vec![0, 1],
1158 edge_ids: vec![7],
1159 neighbor_ids: vec![],
1160 };
1161 assert!(entries_from_out_csr(&mismatched_columns).is_none());
1162
1163 let missing_offsets = CsrIndex {
1164 offsets: vec![],
1165 edge_ids: vec![],
1166 neighbor_ids: vec![],
1167 };
1168 assert_eq!(entries_from_out_csr(&missing_offsets), Some(Vec::new()));
1169 }
1170
1171 #[test]
1172 fn wave10_effective_entry_decode_rejects_malformed_csr_bounds() {
1173 for malformed in [
1174 CsrIndex {
1175 offsets: vec![0, 2],
1176 edge_ids: vec![1],
1177 neighbor_ids: vec![2],
1178 },
1179 CsrIndex {
1180 offsets: vec![1, 0],
1181 edge_ids: vec![1],
1182 neighbor_ids: vec![2],
1183 },
1184 CsrIndex {
1185 offsets: vec![0, 1],
1186 edge_ids: vec![1],
1187 neighbor_ids: vec![],
1188 },
1189 ] {
1190 assert!(entries_from_out_csr(&malformed).is_none());
1191 }
1192 }
1193
1194 #[test]
1195 fn public_csr_writer_rejects_every_inconsistent_topology_shape_without_file() {
1196 let dir = TempDir::new().unwrap();
1197 let path = csr_path(dir.path(), "BROKEN", Direction::Out);
1198 for malformed in [
1199 CsrIndex {
1200 offsets: vec![],
1201 edge_ids: vec![],
1202 neighbor_ids: vec![],
1203 },
1204 CsrIndex {
1205 offsets: vec![0, 2],
1206 edge_ids: vec![1],
1207 neighbor_ids: vec![2],
1208 },
1209 CsrIndex {
1210 offsets: vec![0, 1],
1211 edge_ids: vec![1],
1212 neighbor_ids: vec![],
1213 },
1214 CsrIndex {
1215 offsets: vec![1],
1216 edge_ids: vec![],
1217 neighbor_ids: vec![],
1218 },
1219 ] {
1220 assert_eq!(write_csr(&path, &malformed).unwrap_err().code(), "GF_IO");
1221 assert!(!path.exists());
1222 }
1223 }
1224
1225 #[test]
1226 fn csr_round_trip_preserves_offsets_and_targets() {
1227 let dir = TempDir::new().unwrap();
1228 let path = csr_path(dir.path(), "KNOWS", Direction::Out);
1229 let csr = sample_csr();
1230 write_csr(&path, &csr).unwrap();
1231 assert_eq!(read_csr(&path).unwrap(), csr);
1232 }
1233
1234 #[test]
1235 fn empty_graph_round_trips_as_offsets_zero() {
1236 let dir = TempDir::new().unwrap();
1237 let path = csr_path(dir.path(), "KNOWS", Direction::In);
1238 let csr = CsrIndex {
1239 offsets: vec![0],
1240 ..CsrIndex::default()
1241 };
1242 write_csr(&path, &csr).unwrap();
1243 let back = read_csr(&path).unwrap();
1244 assert_eq!(back, csr);
1245 assert_eq!(back.node_count(), 0);
1246 assert_eq!(back.edge_count(), 0);
1247 }
1248
1249 #[test]
1250 fn node_with_no_neighbors_round_trips() {
1251 let dir = TempDir::new().unwrap();
1252 let path = csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out);
1253 let csr = sample_csr(); write_csr(&path, &csr).unwrap();
1255 let back = read_csr(&path).unwrap();
1256 assert_eq!(back.offsets[1], back.offsets[2], "node 1 has no neighbors");
1257 assert_eq!(back.node_count(), 3);
1258 assert_eq!(back.edge_count(), 4);
1259 }
1260
1261 #[test]
1262 fn write_csr_replaces_existing_file_atomically() {
1263 let dir = TempDir::new().unwrap();
1264 let path = csr_path(dir.path(), "KNOWS", Direction::Out);
1265 write_csr(&path, &sample_csr()).unwrap();
1266
1267 let newer = CsrIndex {
1268 offsets: vec![0, 1],
1269 edge_ids: vec![99],
1270 neighbor_ids: vec![0],
1271 };
1272 write_csr(&path, &newer).unwrap();
1273 assert_eq!(read_csr(&path).unwrap(), newer, "second write wins");
1274
1275 let temps = std::fs::read_dir(adjacency_dir(dir.path()))
1276 .unwrap()
1277 .filter_map(Result::ok)
1278 .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
1279 .count();
1280 assert_eq!(temps, 0, "no temp residue");
1281 }
1282
1283 #[test]
1284 fn read_csr_rejects_wrong_schema() {
1285 let dir = TempDir::new().unwrap();
1286 let path = dir.path().join("bogus.csr");
1287 std::fs::write(&path, b"not an arrow ipc file").unwrap();
1288 assert!(matches!(read_csr(&path), Err(GfError::Storage(_))));
1289 }
1290
1291 #[test]
1292 fn read_csr_missing_file_is_an_error() {
1293 let dir = TempDir::new().unwrap();
1294 let path = csr_path(dir.path(), "ABSENT", Direction::Out);
1295 assert!(matches!(read_csr(&path), Err(GfError::Storage(_))));
1296 }
1297
1298 #[test]
1299 fn write_csr_rejects_invalid_offsets() {
1300 let dir = TempDir::new().unwrap();
1301 let path = dir.path().join("bad.csr");
1302
1303 let non_monotone = CsrIndex {
1305 offsets: vec![0, 3, 1],
1306 edge_ids: vec![1],
1307 neighbor_ids: vec![1],
1308 };
1309 assert!(matches!(
1310 write_csr(&path, &non_monotone),
1311 Err(GfError::Storage(_))
1312 ));
1313
1314 let length_mismatch = CsrIndex {
1316 offsets: vec![0, 2],
1317 edge_ids: vec![1],
1318 neighbor_ids: vec![1],
1319 };
1320 assert!(matches!(
1321 write_csr(&path, &length_mismatch),
1322 Err(GfError::Storage(_))
1323 ));
1324
1325 let empty_offsets = CsrIndex::default();
1327 assert!(matches!(
1328 write_csr(&path, &empty_offsets),
1329 Err(GfError::Storage(_))
1330 ));
1331
1332 let ragged = CsrIndex {
1334 offsets: vec![0, 2],
1335 edge_ids: vec![1, 2],
1336 neighbor_ids: vec![1],
1337 };
1338 assert!(matches!(
1339 write_csr(&path, &ragged),
1340 Err(GfError::Storage(_))
1341 ));
1342
1343 assert!(!path.exists(), "no file written for invalid CSR");
1344 }
1345
1346 #[test]
1347 fn manifest_round_trip_multi_relation() {
1348 const TS: i64 = 1_700_000_000_000_000;
1349 let dir = TempDir::new().unwrap();
1350 let rows = vec![
1351 AdjacencyManifestRow {
1352 relation_type: "WORKS_AT".to_owned(),
1353 direction: Direction::Out,
1354 topology_generation: 7,
1355 built_at_micros: TS,
1356 node_count: 100,
1357 edge_count: 250,
1358 },
1359 AdjacencyManifestRow {
1360 relation_type: "WORKS_AT".to_owned(),
1361 direction: Direction::In,
1362 topology_generation: 7,
1363 built_at_micros: TS,
1364 node_count: 100,
1365 edge_count: 250,
1366 },
1367 AdjacencyManifestRow {
1368 relation_type: "OWNS".to_owned(),
1369 direction: Direction::Out,
1370 topology_generation: 7,
1371 built_at_micros: TS + 1,
1372 node_count: 40,
1373 edge_count: 41,
1374 },
1375 AdjacencyManifestRow {
1376 relation_type: ALL_RELATIONS_STEM.to_owned(),
1377 direction: Direction::Out,
1378 topology_generation: 7,
1379 built_at_micros: TS + 2,
1380 node_count: 100,
1381 edge_count: 291,
1382 },
1383 ];
1384 write_manifest(dir.path(), &rows).unwrap();
1385 assert_eq!(read_manifest(dir.path()).unwrap(), rows);
1386 }
1387
1388 #[test]
1389 fn read_manifest_absent_returns_empty() {
1390 let dir = TempDir::new().unwrap();
1391 assert_eq!(read_manifest(dir.path()).unwrap(), Vec::new());
1392 }
1393
1394 #[test]
1395 fn write_manifest_replaces_existing() {
1396 let dir = TempDir::new().unwrap();
1397 let first = vec![AdjacencyManifestRow {
1398 relation_type: "KNOWS".to_owned(),
1399 direction: Direction::Out,
1400 topology_generation: 1,
1401 built_at_micros: 0,
1402 node_count: 1,
1403 edge_count: 1,
1404 }];
1405 write_manifest(dir.path(), &first).unwrap();
1406
1407 let second = vec![AdjacencyManifestRow {
1408 relation_type: "KNOWS".to_owned(),
1409 direction: Direction::Out,
1410 topology_generation: 2,
1411 built_at_micros: 1,
1412 node_count: 2,
1413 edge_count: 3,
1414 }];
1415 write_manifest(dir.path(), &second).unwrap();
1416 assert_eq!(read_manifest(dir.path()).unwrap(), second);
1417 }
1418
1419 #[test]
1420 fn read_manifest_rejects_wrong_schema() {
1421 let dir = TempDir::new().unwrap();
1422 let schema = Arc::new(arrow::datatypes::Schema::new(vec![Field::new(
1424 "v",
1425 DataType::Int64,
1426 false,
1427 )]));
1428 let batch = RecordBatch::try_new(
1429 Arc::clone(&schema),
1430 vec![Arc::new(arrow::array::Int64Array::from(vec![1]))],
1431 )
1432 .unwrap();
1433 let mut staged = RewriteBatch::new();
1434 staged
1435 .stage(&manifest_path(dir.path()), schema, &batch)
1436 .unwrap();
1437 staged.commit().unwrap();
1438
1439 assert!(matches!(
1440 read_manifest(dir.path()),
1441 Err(GfError::Storage(_))
1442 ));
1443 }
1444
1445 #[test]
1446 fn direction_round_trips_through_str() {
1447 for d in [Direction::Out, Direction::In] {
1448 assert_eq!(Direction::parse(d.as_str()).unwrap(), d);
1449 }
1450 assert!(matches!(
1451 Direction::parse("sideways"),
1452 Err(GfError::Storage(_))
1453 ));
1454 }
1455
1456 #[test]
1457 fn csr_path_layout() {
1458 let p = csr_path(Path::new("/proj"), "WORKS_AT", Direction::In);
1459 assert_eq!(p, Path::new("/proj/indexes/adjacency/WORKS_AT.in.csr"));
1460 assert_eq!(
1461 manifest_path(Path::new("/proj")),
1462 Path::new("/proj/indexes/adjacency/index_manifest.parquet")
1463 );
1464 }
1465
1466 use crate::GraphWriter;
1471 use graphforge_core::OntologyMode;
1472 use graphforge_core::TypeId;
1473 use graphforge_core::uuid::{Uuid, new_v7, to_bytes};
1474
1475 const BUILD_TS: i64 = 1_700_000_000_000_000;
1477
1478 fn write_diamond(dir: &Path) -> [u64; 4] {
1481 let mut w = GraphWriter::open_at(dir, OntologyMode::Strict, BUILD_TS).unwrap();
1482 let uuids: Vec<Uuid> = (0..4).map(|_| new_v7()).collect();
1483 let ids: Vec<u64> = uuids
1484 .iter()
1485 .map(|u| w.create_node(*u, TypeId(0)).unwrap())
1486 .collect();
1487 let (a, b, c, d) = (&uuids[0], &uuids[1], &uuids[2], &uuids[3]);
1488 for (src, dst) in [(a, b), (a, c), (b, d), (c, d), (a, b), (d, d)] {
1489 w.create_edge(new_v7(), "KNOWS", src, dst).unwrap();
1490 }
1491 w.flush().unwrap();
1492 [ids[0], ids[1], ids[2], ids[3]]
1493 }
1494
1495 #[test]
1496 fn build_is_deterministic_and_stamps_pre_scan_generation() {
1497 let dir = TempDir::new().unwrap();
1498 write_diamond(dir.path()); let rows = build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1501 assert!(rows.iter().all(|r| r.topology_generation == 1));
1502 assert_eq!(rows.len(), 4);
1504
1505 let knows_out = std::fs::read(csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap();
1506 let all_in =
1507 std::fs::read(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::In)).unwrap();
1508
1509 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1511 assert_eq!(
1512 std::fs::read(csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap(),
1513 knows_out
1514 );
1515 assert_eq!(
1516 std::fs::read(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::In)).unwrap(),
1517 all_in
1518 );
1519
1520 let csr = read_csr(&csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap();
1522 let knows_manifest = read_manifest(dir.path())
1523 .unwrap()
1524 .into_iter()
1525 .find(|r| r.relation_type == "KNOWS" && r.direction == Direction::Out)
1526 .unwrap();
1527 assert_eq!(knows_manifest.node_count, csr.node_count());
1528 assert_eq!(knows_manifest.edge_count, 6);
1529 let windows: Vec<&[u64]> = csr
1530 .offsets
1531 .windows(2)
1532 .map(|w| &csr.edge_ids[w[0] as usize..w[1] as usize])
1533 .collect();
1534 for per_node in windows {
1535 assert!(
1536 per_node.windows(2).all(|w| w[0] <= w[1]),
1537 "edge ids ascending per node"
1538 );
1539 }
1540 }
1541
1542 #[test]
1543 fn private_staging_never_changes_the_reader_visible_artifact() {
1544 let dir = TempDir::new().unwrap();
1545 write_diamond(dir.path());
1546 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1547 let prior_manifest = std::fs::read(manifest_path(dir.path())).unwrap();
1548 let stage = TempDir::new_in(dir.path().parent().unwrap()).unwrap();
1549 let mut checkpoints = 0;
1550
1551 build_adjacency_index_into(dir.path(), stage.path(), BUILD_TS + 1, || {
1552 checkpoints += 1;
1553 assert_eq!(
1554 std::fs::read(manifest_path(dir.path())).unwrap(),
1555 prior_manifest
1556 );
1557 assert_eq!(
1558 inspect_adjacency_index(dir.path()).unwrap().state,
1559 AdjacencyFreshnessState::Current
1560 );
1561 Ok(())
1562 })
1563 .unwrap();
1564 assert!(checkpoints >= 4);
1565 assert!(
1566 validate_adjacency_index_against(dir.path(), stage.path())
1567 .unwrap()
1568 .is_empty()
1569 );
1570 }
1571
1572 #[test]
1573 fn post_delete_sparse_ids_build_round_trips() {
1574 let dir = TempDir::new().unwrap();
1575 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, BUILD_TS).unwrap();
1576 let uuids: Vec<Uuid> = (0..4).map(|_| new_v7()).collect();
1577 let ids: Vec<u64> = uuids
1578 .iter()
1579 .map(|u| w.create_node(*u, TypeId(0)).unwrap())
1580 .collect();
1581 for pair in uuids.windows(2) {
1582 w.create_edge(new_v7(), "KNOWS", &pair[0], &pair[1])
1583 .unwrap();
1584 }
1585 w.flush().unwrap();
1586
1587 let node_set: std::collections::HashSet<[u8; 16]> =
1589 std::iter::once(to_bytes(&uuids[1])).collect();
1590 let incident = crate::incident_edge_uuids(dir.path(), &node_set).unwrap();
1591 let edge_set: std::collections::HashSet<[u8; 16]> = incident.into_iter().collect();
1592 crate::delete_nodes_and_edges(dir.path(), &node_set, &edge_set).unwrap();
1593
1594 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1595 let csr = read_csr(&csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap();
1596 let gap = usize::try_from(ids[1]).unwrap();
1597 assert_eq!(
1598 csr.offsets[gap],
1599 csr.offsets[gap + 1],
1600 "deleted id is an empty range"
1601 );
1602 let n3 = usize::try_from(ids[2]).unwrap();
1604 let (s, e) = (csr.offsets[n3] as usize, csr.offsets[n3 + 1] as usize);
1605 assert_eq!(&csr.neighbor_ids[s..e], &[ids[3]]);
1606 }
1607
1608 #[test]
1609 fn exploratory_rows_group_by_rel_type_and_union_covers_all() {
1610 let dir = TempDir::new().unwrap();
1611 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, BUILD_TS).unwrap();
1612 let (a, b, c) = (new_v7(), new_v7(), new_v7());
1613 let ids: Vec<u64> = [a, b, c]
1614 .iter()
1615 .map(|u| w.create_node(*u, TypeId(0)).unwrap())
1616 .collect();
1617 w.create_edge(new_v7(), "KNOWS", &a, &b).unwrap();
1618 w.create_edge(new_v7(), "OWNS", &a, &c).unwrap();
1619 w.flush().unwrap();
1620
1621 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1622
1623 let knows = read_csr(&csr_path(dir.path(), "KNOWS", Direction::Out)).unwrap();
1624 assert_eq!(knows.edge_count(), 1, "decoy OWNS row excluded");
1625 let owns = read_csr(&csr_path(dir.path(), "OWNS", Direction::Out)).unwrap();
1626 assert_eq!(owns.edge_count(), 1);
1627
1628 let all = read_csr(&csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out)).unwrap();
1629 assert_eq!(all.edge_count(), 2, "union covers both rel types");
1630 let row = usize::try_from(ids[0]).unwrap();
1631 let (s, e) = (all.offsets[row] as usize, all.offsets[row + 1] as usize);
1632 assert_eq!(&all.edge_ids[s..e], &[1, 2], "union in edge_id order");
1633 assert_eq!(&all.neighbor_ids[s..e], &[ids[1], ids[2]]);
1634 }
1635
1636 #[test]
1637 fn hostile_and_reserved_stems_are_skipped_but_counted_in_union() {
1638 let dir = TempDir::new().unwrap();
1639 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Exploratory, BUILD_TS).unwrap();
1640 let (a, b, c) = (new_v7(), new_v7(), new_v7());
1641 for u in [a, b, c] {
1642 w.create_node(u, TypeId(0)).unwrap();
1643 }
1644 w.create_edge(new_v7(), "a/b", &a, &b).unwrap();
1645 w.create_edge(new_v7(), ALL_RELATIONS_STEM, &a, &c).unwrap();
1646 w.flush().unwrap();
1647
1648 let rows = build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1649 assert!(rows.iter().all(|r| r.relation_type == ALL_RELATIONS_STEM));
1651 assert_eq!(rows.len(), 2);
1652 let all = read_csr(&csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out)).unwrap();
1653 assert_eq!(
1654 all.edge_count(),
1655 2,
1656 "skipped rels still flow into the union"
1657 );
1658 assert!(
1659 !csr_path(dir.path(), "a/b", Direction::Out).exists(),
1660 "no nested path written for the separator-bearing rel name"
1661 );
1662 assert!(!csr_path(dir.path(), "a", Direction::Out).exists());
1663 }
1664
1665 #[test]
1666 fn build_failure_leaves_no_manifest() {
1667 let dir = TempDir::new().unwrap();
1668 write_diamond(dir.path());
1669 let adj = adjacency_dir(dir.path());
1671 std::fs::create_dir_all(&adj).unwrap();
1672 let mut perms = std::fs::metadata(&adj).unwrap().permissions();
1673 perms.set_readonly(true);
1674 std::fs::set_permissions(&adj, perms.clone()).unwrap();
1675
1676 let result = build_adjacency_index(dir.path(), BUILD_TS);
1677 perms.set_readonly(false);
1679 std::fs::set_permissions(&adj, perms).unwrap();
1680
1681 assert!(result.is_err());
1682 assert!(!manifest_path(dir.path()).exists(), "manifest written last");
1683 assert_eq!(read_manifest(dir.path()).unwrap(), Vec::new());
1684 }
1685
1686 #[test]
1687 fn empty_project_builds_union_pair_and_manifest() {
1688 let dir = TempDir::new().unwrap();
1689 let rows = build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1690 assert_eq!(rows.len(), 2);
1691 assert!(rows.iter().all(|r| {
1692 r.relation_type == ALL_RELATIONS_STEM
1693 && r.topology_generation == 0
1694 && r.node_count == 0
1695 && r.edge_count == 0
1696 }));
1697 let all = read_csr(&csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out)).unwrap();
1698 assert_eq!(all.offsets, vec![0]);
1699 assert_eq!(read_manifest(dir.path()).unwrap(), rows);
1700 }
1701
1702 #[test]
1707 fn validate_reports_clean_on_valid_index() {
1708 let dir = TempDir::new().unwrap();
1709 write_diamond(dir.path());
1710 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1711 assert_eq!(validate_adjacency_index(dir.path()).unwrap(), Vec::new());
1712 }
1713
1714 #[test]
1715 fn validate_reports_clean_on_absent_index() {
1716 let dir = TempDir::new().unwrap();
1717 write_diamond(dir.path());
1718 assert_eq!(validate_adjacency_index(dir.path()).unwrap(), Vec::new());
1719 }
1720
1721 #[test]
1722 fn inspection_moves_from_missing_to_current_with_stable_identity() {
1723 let dir = TempDir::new().unwrap();
1724 write_diamond(dir.path());
1725 let missing = inspect_adjacency_index(dir.path()).unwrap();
1726 assert_eq!(missing.state, AdjacencyFreshnessState::Missing);
1727 assert_eq!(missing.reason, Some(AdjacencyFreshnessReason::NotBuilt));
1728 assert!(missing.source_fingerprint.starts_with("sha256:"));
1729
1730 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1731 let current = inspect_adjacency_index(dir.path()).unwrap();
1732 assert_eq!(current.state, AdjacencyFreshnessState::Current);
1733 assert_eq!(current.reason, None);
1734 assert_eq!(current.artifact_generation, Some(current.source_generation));
1735 assert_eq!(
1736 current.artifact_fingerprint.as_deref(),
1737 Some(current.source_fingerprint.as_str())
1738 );
1739 }
1740
1741 #[test]
1742 fn inspection_checks_every_manifest_referenced_csr() {
1743 let dir = TempDir::new().unwrap();
1744 write_diamond(dir.path());
1745 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1746 std::fs::write(csr_path(dir.path(), "KNOWS", Direction::In), b"garbage").unwrap();
1747
1748 let inspection = inspect_adjacency_index(dir.path()).unwrap();
1749 assert_eq!(inspection.state, AdjacencyFreshnessState::Incompatible);
1750 assert_eq!(
1751 inspection.reason,
1752 Some(AdjacencyFreshnessReason::UnreadableArtifact)
1753 );
1754 }
1755
1756 #[test]
1757 fn inspection_distinguishes_manifest_union_absence_and_union_corruption_after_reopen() {
1758 for case in ["manifest", "missing-union", "corrupt-union"] {
1759 let dir = TempDir::new().unwrap();
1760 write_diamond(dir.path());
1761 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1762 match case {
1763 "manifest" => std::fs::write(manifest_path(dir.path()), b"corrupt").unwrap(),
1764 "missing-union" => {
1765 std::fs::remove_file(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out))
1766 .unwrap();
1767 }
1768 "corrupt-union" => {
1769 std::fs::write(
1770 csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out),
1771 b"corrupt",
1772 )
1773 .unwrap();
1774 }
1775 _ => unreachable!(),
1776 }
1777
1778 let inspection = inspect_adjacency_index(dir.path()).unwrap();
1779 assert_eq!(inspection.state, AdjacencyFreshnessState::Incompatible);
1780 assert_eq!(
1781 inspection.reason,
1782 Some(if case == "missing-union" {
1783 AdjacencyFreshnessReason::MissingCsr
1784 } else {
1785 AdjacencyFreshnessReason::UnreadableArtifact
1786 })
1787 );
1788 assert_eq!(inspection.artifact_fingerprint, None);
1789 }
1790 }
1791
1792 #[test]
1793 fn inspection_accepts_only_a_complete_delta_chain_as_current() {
1794 let dir = TempDir::new().unwrap();
1795 write_diamond(dir.path());
1796 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1797 crate::generation::bump_topology_generation(dir.path()).unwrap();
1798
1799 let inspection = inspect_adjacency_index(dir.path()).unwrap();
1800 assert_eq!(inspection.state, AdjacencyFreshnessState::Stale);
1801 assert_eq!(
1802 inspection.reason,
1803 Some(AdjacencyFreshnessReason::IncompleteDeltaChain)
1804 );
1805 }
1806
1807 #[test]
1808 fn freshness_vocabulary_and_manifest_generation_failures_are_exact() {
1809 assert_eq!(AdjacencyFreshnessState::Current.as_str(), "current");
1810 assert_eq!(AdjacencyFreshnessState::Missing.as_str(), "missing");
1811 assert_eq!(AdjacencyFreshnessState::Stale.as_str(), "stale");
1812 assert_eq!(
1813 AdjacencyFreshnessState::Incompatible.as_str(),
1814 "incompatible"
1815 );
1816 for (reason, token) in [
1817 (AdjacencyFreshnessReason::NotBuilt, "not_built"),
1818 (
1819 AdjacencyFreshnessReason::MixedArtifactGeneration,
1820 "mixed_artifact_generation",
1821 ),
1822 (
1823 AdjacencyFreshnessReason::IncompleteDeltaChain,
1824 "incomplete_delta_chain",
1825 ),
1826 (AdjacencyFreshnessReason::MissingCsr, "missing_csr"),
1827 (
1828 AdjacencyFreshnessReason::UnreadableArtifact,
1829 "unreadable_artifact",
1830 ),
1831 (
1832 AdjacencyFreshnessReason::ContentMismatch,
1833 "content_mismatch",
1834 ),
1835 (
1836 AdjacencyFreshnessReason::FutureArtifactGeneration,
1837 "future_artifact_generation",
1838 ),
1839 ] {
1840 assert_eq!(reason.as_str(), token);
1841 }
1842
1843 let mixed = TempDir::new().unwrap();
1844 write_diamond(mixed.path());
1845 build_adjacency_index(mixed.path(), BUILD_TS).unwrap();
1846 let mut manifest = read_manifest(mixed.path()).unwrap();
1847 manifest[0].topology_generation += 1;
1848 write_manifest(mixed.path(), &manifest).unwrap();
1849 let inspection = inspect_adjacency_index(mixed.path()).unwrap();
1850 assert_eq!(inspection.state, AdjacencyFreshnessState::Incompatible);
1851 assert_eq!(
1852 inspection.reason,
1853 Some(AdjacencyFreshnessReason::MixedArtifactGeneration)
1854 );
1855 assert_eq!(inspection.artifact_generation, None);
1856
1857 let future = TempDir::new().unwrap();
1858 write_diamond(future.path());
1859 build_adjacency_index(future.path(), BUILD_TS).unwrap();
1860 let mut manifest = read_manifest(future.path()).unwrap();
1861 for row in &mut manifest {
1862 row.topology_generation += 1;
1863 }
1864 write_manifest(future.path(), &manifest).unwrap();
1865 let inspection = inspect_adjacency_index(future.path()).unwrap();
1866 assert_eq!(inspection.state, AdjacencyFreshnessState::Incompatible);
1867 assert_eq!(
1868 inspection.reason,
1869 Some(AdjacencyFreshnessReason::FutureArtifactGeneration)
1870 );
1871 assert_eq!(inspection.artifact_generation, Some(2));
1872 assert_eq!(inspection.artifact_fingerprint, None);
1873 }
1874
1875 #[test]
1876 fn validate_detects_corrupted_csr() {
1877 let dir = TempDir::new().unwrap();
1878 write_diamond(dir.path());
1879 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1880
1881 let bogus = CsrIndex {
1883 offsets: vec![0, 1],
1884 edge_ids: vec![99],
1885 neighbor_ids: vec![1],
1886 };
1887 write_csr(&csr_path(dir.path(), "KNOWS", Direction::Out), &bogus).unwrap();
1888
1889 let issues = validate_adjacency_index(dir.path()).unwrap();
1890 assert_eq!(
1891 issues,
1892 vec![AdjacencyValidationIssue::Mismatch {
1893 rel: "KNOWS".to_owned(),
1894 direction: Direction::Out,
1895 }]
1896 );
1897 }
1898
1899 #[test]
1900 fn validate_detects_unreadable_csr() {
1901 let dir = TempDir::new().unwrap();
1902 write_diamond(dir.path());
1903 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1904 std::fs::write(csr_path(dir.path(), "KNOWS", Direction::In), b"garbage").unwrap();
1905
1906 let issues = validate_adjacency_index(dir.path()).unwrap();
1907 assert_eq!(issues.len(), 1);
1908 assert!(matches!(
1909 &issues[0],
1910 AdjacencyValidationIssue::UnreadableCsr { rel, direction: Direction::In, .. }
1911 if rel == "KNOWS"
1912 ));
1913 }
1914
1915 #[test]
1916 fn validate_detects_missing_csr() {
1917 let dir = TempDir::new().unwrap();
1918 write_diamond(dir.path());
1919 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1920 std::fs::remove_file(csr_path(dir.path(), ALL_RELATIONS_STEM, Direction::Out)).unwrap();
1921
1922 let issues = validate_adjacency_index(dir.path()).unwrap();
1923 assert_eq!(
1924 issues,
1925 vec![AdjacencyValidationIssue::MissingCsr {
1926 rel: ALL_RELATIONS_STEM.to_owned(),
1927 direction: Direction::Out,
1928 }]
1929 );
1930 }
1931
1932 #[test]
1933 fn validate_detects_stale_generation_only_once() {
1934 let dir = TempDir::new().unwrap();
1935 write_diamond(dir.path());
1936 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1937 crate::generation::bump_topology_generation(dir.path()).unwrap();
1940
1941 let issues = validate_adjacency_index(dir.path()).unwrap();
1942 assert_eq!(
1943 issues,
1944 vec![AdjacencyValidationIssue::StaleGeneration {
1945 manifest: 1,
1946 current: 2,
1947 }]
1948 );
1949 }
1950
1951 #[test]
1956 fn validate_clean_on_delta_covered_index() {
1957 let dir = TempDir::new().unwrap();
1958 write_diamond(dir.path());
1959 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1960
1961 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, BUILD_TS).unwrap();
1963 let (a, b) = (new_v7(), new_v7());
1964 w.create_node(a, TypeId(0)).unwrap();
1965 w.create_node(b, TypeId(0)).unwrap();
1966 w.create_edge(new_v7(), "KNOWS", &a, &b).unwrap();
1967 w.flush().unwrap();
1968
1969 assert!(
1970 validate_adjacency_index(dir.path()).unwrap().is_empty(),
1971 "base + intact chain == current rebuild ⇒ no issues"
1972 );
1973 let inspection = inspect_adjacency_index(dir.path()).unwrap();
1974 assert_eq!(inspection.state, AdjacencyFreshnessState::Current);
1975 assert_eq!(
1976 inspection.artifact_effective_generation,
1977 Some(inspection.source_generation)
1978 );
1979 assert_eq!(
1980 inspection.artifact_fingerprint.as_deref(),
1981 Some(inspection.source_fingerprint.as_str())
1982 );
1983 }
1984
1985 #[test]
1989 fn validate_detects_mismatch_under_delta_chain() {
1990 let dir = TempDir::new().unwrap();
1991 write_diamond(dir.path());
1992 build_adjacency_index(dir.path(), BUILD_TS).unwrap();
1993 let mut w = GraphWriter::open_at(dir.path(), OntologyMode::Strict, BUILD_TS).unwrap();
1994 let (a, b) = (new_v7(), new_v7());
1995 w.create_node(a, TypeId(0)).unwrap();
1996 w.create_node(b, TypeId(0)).unwrap();
1997 w.create_edge(new_v7(), "KNOWS", &a, &b).unwrap();
1998 w.flush().unwrap();
1999
2000 let knows_out = csr_path(dir.path(), "KNOWS", Direction::Out);
2004 write_csr(&knows_out, &csr_from_entries(&[], Direction::Out)).unwrap();
2005
2006 let issues = validate_adjacency_index(dir.path()).unwrap();
2007 assert!(
2008 issues.contains(&AdjacencyValidationIssue::Mismatch {
2009 rel: "KNOWS".to_owned(),
2010 direction: Direction::Out,
2011 }),
2012 "corruption under a chain is still detected: {issues:?}"
2013 );
2014 let inspection = inspect_adjacency_index(dir.path()).unwrap();
2015 assert_eq!(inspection.state, AdjacencyFreshnessState::Incompatible);
2016 assert_eq!(
2017 inspection.artifact_effective_generation,
2018 Some(inspection.source_generation)
2019 );
2020 assert!(inspection.artifact_fingerprint.is_some());
2021 }
2022
2023 #[test]
2024 fn wave13_csr_io_rejects_parentless_destination_and_wrong_arrow_schema() {
2025 let empty = CsrIndex {
2026 offsets: vec![0],
2027 edge_ids: vec![],
2028 neighbor_ids: vec![],
2029 };
2030 assert!(write_csr(Path::new("/"), &empty).is_err());
2031
2032 let root = TempDir::new().unwrap();
2033 let path = root.path().join("wrong-schema.arrow");
2034 let schema = Arc::new(arrow::datatypes::Schema::new(vec![Field::new(
2035 "wrong",
2036 DataType::UInt64,
2037 false,
2038 )]));
2039 let batch = RecordBatch::try_new(
2040 Arc::clone(&schema),
2041 vec![Arc::new(UInt64Array::from(vec![1]))],
2042 )
2043 .unwrap();
2044 let file = std::fs::File::create(&path).unwrap();
2045 let mut writer = FileWriter::try_new(file, &schema).unwrap();
2046 writer.write(&batch).unwrap();
2047 writer.finish().unwrap();
2048 assert!(read_csr(&path).is_err());
2049 }
2050}