1use crate::storage::adjacency_overlay::{FrozenCsrSegment, L0CsrSegment};
37use crate::storage::csr::MainCsr;
38use crate::storage::direction::Direction;
39use crate::storage::manager::StorageManager;
40use crate::storage::shadow_csr::{ShadowCsr, ShadowEdge};
41use dashmap::DashMap;
42use parking_lot::RwLock;
43use std::collections::{HashMap, HashSet};
44use std::sync::Arc;
45use std::sync::atomic::{AtomicUsize, Ordering};
46use uni_common::core::id::{Eid, Vid};
47
48#[derive(Debug, Default)]
67pub struct PinnedVersions {
68 counts: parking_lot::Mutex<std::collections::BTreeMap<u64, usize>>,
69}
70
71impl PinnedVersions {
72 pub fn pin(self: &Arc<Self>, version: u64) -> PinGuard {
74 *self.counts.lock().entry(version).or_insert(0) += 1;
75 PinGuard {
76 versions: Arc::clone(self),
77 version,
78 }
79 }
80
81 pub fn min_pinned(&self) -> Option<u64> {
83 self.counts.lock().keys().next().copied()
84 }
85
86 pub fn distinct_pinned(&self) -> usize {
88 self.counts.lock().len()
89 }
90}
91
92#[derive(Debug)]
94pub struct PinGuard {
95 versions: Arc<PinnedVersions>,
96 version: u64,
97}
98
99impl Drop for PinGuard {
100 fn drop(&mut self) {
101 let mut counts = self.versions.counts.lock();
102 if let Some(n) = counts.get_mut(&self.version) {
103 *n -= 1;
104 if *n == 0 {
105 counts.remove(&self.version);
106 }
107 }
108 }
109}
110
111fn dedup_entries_by_eid(entries: &mut Vec<(u64, Vid, Eid, u64)>) {
115 use std::collections::hash_map::Entry;
116
117 let mut best: HashMap<Eid, usize> = HashMap::new();
118 for (idx, (_, _, eid, ver)) in entries.iter().enumerate() {
119 match best.entry(*eid) {
120 Entry::Vacant(e) => {
121 e.insert(idx);
122 }
123 Entry::Occupied(mut e) => {
124 if *ver > entries[*e.get()].3 {
125 e.insert(idx);
126 }
127 }
128 }
129 }
130 let keep: HashSet<usize> = best.into_values().collect();
131 let mut idx = 0;
132 entries.retain(|_| {
133 let k = keep.contains(&idx);
134 idx += 1;
135 k
136 });
137}
138
139pub struct AdjacencyManager {
145 main_csr: DashMap<(u32, Direction), Arc<MainCsr>>,
148
149 active_overlay: Arc<RwLock<L0CsrSegment>>,
151
152 frozen_segments: RwLock<Vec<Arc<FrozenCsrSegment>>>,
154
155 shadow: ShadowCsr,
157 pinned_versions: Arc<PinnedVersions>,
160
161 current_bytes: AtomicUsize,
163
164 max_bytes: usize,
166
167 warm_guards: DashMap<(u32, Direction), Arc<tokio::sync::Mutex<()>>>,
170
171 compact_lock: parking_lot::Mutex<()>,
174}
175
176impl AdjacencyManager {
177 pub fn new(max_bytes: usize) -> Self {
179 Self {
180 main_csr: DashMap::new(),
181 active_overlay: Arc::new(RwLock::new(L0CsrSegment::new())),
182 frozen_segments: RwLock::new(Vec::new()),
183 shadow: ShadowCsr::new(),
184 pinned_versions: Arc::new(PinnedVersions::default()),
185 current_bytes: AtomicUsize::new(0),
186 max_bytes,
187 warm_guards: DashMap::new(),
188 compact_lock: parking_lot::Mutex::new(()),
189 }
190 }
191
192 pub fn get_neighbors(&self, vid: Vid, edge_type: u32, direction: Direction) -> Vec<(Vid, Eid)> {
197 let mut result: HashMap<Eid, Vid> = HashMap::new();
198
199 for &dir in direction.expand() {
200 if let Some(csr) = self.main_csr.get(&(edge_type, dir)) {
202 for entry in csr.get_entries(vid) {
203 result.insert(entry.eid, entry.neighbor_vid);
204 }
205 }
206
207 for segment in self.frozen_segments.read().iter() {
213 let has_inserts = segment.has_entries_for(edge_type, dir);
214 let has_tombstones = !segment.tombstones.is_empty();
215 if !has_inserts && !has_tombstones {
216 continue;
217 }
218 if has_inserts
219 && let Some(adj) = segment.inserts.get(&(edge_type, dir))
220 && let Some(neighbors) = adj.get(&vid)
221 {
222 for &(neighbor, eid, _version) in neighbors {
223 result.insert(eid, neighbor);
224 }
225 }
226 if has_tombstones {
230 result.retain(|eid, _| !segment.tombstones.contains_key(eid));
231 }
232 }
233
234 let active = self.active_overlay.read();
238 let active_has_inserts = active.has_entries_for(edge_type, dir);
239 let active_has_tombstones = !active.tombstones.is_empty();
240 if active_has_inserts
241 && let Some(adj) = active.inserts.get(&(edge_type, dir))
242 && let Some(neighbors) = adj.get(&vid)
243 {
244 for &(neighbor, eid, _version) in neighbors {
245 result.insert(eid, neighbor);
246 }
247 }
248 if active_has_tombstones {
249 result.retain(|eid, _| !active.tombstones.contains_key(eid));
250 }
251 }
252
253 result.into_iter().map(|(e, n)| (n, e)).collect()
254 }
255
256 pub fn get_neighbors_at_version(
262 &self,
263 vid: Vid,
264 edge_type: u32,
265 direction: Direction,
266 version: u64,
267 ) -> Vec<(Vid, Eid)> {
268 let mut result: HashMap<Eid, Vid> = HashMap::new();
269
270 for &dir in direction.expand() {
271 if let Some(csr) = self.main_csr.get(&(edge_type, dir)) {
273 for entry in csr.get_entries(vid) {
274 if entry.created_version <= version {
275 result.insert(entry.eid, entry.neighbor_vid);
276 }
277 }
278 }
279
280 for segment in self.frozen_segments.read().iter() {
283 let has_inserts = segment.has_entries_for(edge_type, dir);
284 let has_tombstones = !segment.tombstones.is_empty();
285 if !has_inserts && !has_tombstones {
286 continue;
287 }
288 if has_inserts
289 && let Some(adj) = segment.inserts.get(&(edge_type, dir))
290 && let Some(neighbors) = adj.get(&vid)
291 {
292 for &(neighbor, eid, ver) in neighbors {
293 if ver <= version {
294 result.insert(eid, neighbor);
295 }
296 }
297 }
298 if has_tombstones {
299 result.retain(|eid, _| {
300 segment
301 .tombstones
302 .get(eid)
303 .is_none_or(|ts| ts.version > version)
304 });
305 }
306 }
307
308 let active = self.active_overlay.read();
310 let active_has_inserts = active.has_entries_for(edge_type, dir);
311 let active_has_tombstones = !active.tombstones.is_empty();
312 if active_has_inserts
313 && let Some(adj) = active.inserts.get(&(edge_type, dir))
314 && let Some(neighbors) = adj.get(&vid)
315 {
316 for &(neighbor, eid, ver) in neighbors {
317 let not_tombstoned = active
318 .tombstones
319 .get(&eid)
320 .is_none_or(|ts| ts.version > version);
321 if ver <= version && not_tombstoned {
322 result.insert(eid, neighbor);
323 }
324 }
325 }
326 if active_has_tombstones {
327 result.retain(|eid, _| {
328 active
329 .tombstones
330 .get(eid)
331 .is_none_or(|ts| ts.version > version)
332 });
333 }
334
335 for (neighbor, eid) in self
337 .shadow
338 .get_entries_at_version(vid, edge_type, dir, version)
339 {
340 result.insert(eid, neighbor);
341 }
342 }
343
344 result.into_iter().map(|(e, n)| (n, e)).collect()
345 }
346
347 pub fn insert_edge(&self, src: Vid, dst: Vid, eid: Eid, edge_type: u32, version: u64) {
349 let active = self.active_overlay.read();
350 active.insert_edge(src, dst, eid, edge_type, version, Direction::Outgoing);
351 active.insert_edge(dst, src, eid, edge_type, version, Direction::Incoming);
352 }
353
354 pub fn add_tombstone(&self, eid: Eid, src: Vid, dst: Vid, edge_type: u32, version: u64) {
356 let active = self.active_overlay.read();
357 active.add_tombstone(eid, src, dst, edge_type, version);
358 }
359
360 pub fn set_main_csr(&self, edge_type: u32, direction: Direction, csr: MainCsr) {
364 let size = csr.memory_usage();
365 self.main_csr.insert((edge_type, direction), Arc::new(csr));
366 self.current_bytes.fetch_add(size, Ordering::Relaxed);
367 }
368
369 pub fn has_csr(&self, edge_type: u32, direction: Direction) -> bool {
371 self.main_csr.contains_key(&(edge_type, direction))
372 }
373
374 pub fn is_active_for(&self, edge_type: u32, direction: Direction) -> bool {
379 let active = self.active_overlay.read();
380 direction.expand().iter().any(|&d| {
381 self.main_csr.contains_key(&(edge_type, d)) || active.has_entries_for(edge_type, d)
382 })
383 }
384
385 #[must_use]
403 pub fn known_edge_type_ids(&self) -> Vec<u32> {
404 let mut ids: Vec<u32> = self.main_csr.iter().map(|entry| entry.key().0).collect();
405 for entry in self.active_overlay.read().inserts.iter() {
406 ids.push(entry.key().0);
407 }
408 for segment in self.frozen_segments.read().iter() {
409 ids.extend(segment.inserts.keys().map(|&(etype, _dir)| etype));
410 }
411 ids.sort_unstable();
412 ids.dedup();
413 ids
414 }
415
416 pub fn frozen_segment_count(&self) -> usize {
418 self.frozen_segments.read().len()
419 }
420
421 pub fn should_compact(&self, threshold: usize) -> bool {
423 self.frozen_segment_count() >= threshold
424 }
425
426 pub fn compact(&self) {
435 let _compact_guard = self.compact_lock.lock();
439
440 let frozen = {
442 let mut active = self.active_overlay.write();
443 let old = std::mem::take(&mut *active);
444 Arc::new(old.freeze())
445 };
446 self.frozen_segments.write().push(frozen);
447
448 let segments = self.frozen_segments.read().clone();
451
452 let mut all_keys: HashSet<(u32, Direction)> = HashSet::new();
454 for segment in &segments {
455 for key in segment.inserts.keys() {
456 all_keys.insert(*key);
457 }
458 }
459 for entry in self.main_csr.iter() {
460 all_keys.insert(*entry.key());
461 }
462
463 for (edge_type, direction) in all_keys {
465 let mut entries: Vec<(u64, Vid, Eid, u64)> = Vec::new();
466 let mut max_offset: u64 = 0;
467
468 let mut tombstoned_eids: HashSet<Eid> = HashSet::new();
470 for segment in &segments {
471 for (eid, ts) in &segment.tombstones {
472 if ts.edge_type == edge_type {
473 tombstoned_eids.insert(*eid);
474
475 let (key_vid, neighbor_vid) = if direction == Direction::Incoming {
484 (ts.dst_vid, ts.src_vid)
485 } else {
486 (ts.src_vid, ts.dst_vid)
487 };
488 self.shadow.add_deleted_edge(
489 key_vid,
490 ShadowEdge {
491 neighbor_vid,
492 eid: *eid,
493 edge_type,
494 created_version: 0, deleted_version: ts.version,
496 },
497 direction,
498 );
499 }
500 }
501 }
502
503 if let Some(old_csr) = self.main_csr.get(&(edge_type, direction)) {
505 for vid_offset in 0..old_csr.num_vertices() {
506 let vid = Vid::new(vid_offset as u64);
507 for entry in old_csr.get_entries(vid) {
508 if !tombstoned_eids.contains(&entry.eid) {
509 entries.push((
510 vid_offset as u64,
511 entry.neighbor_vid,
512 entry.eid,
513 entry.created_version,
514 ));
515 max_offset = max_offset.max(vid_offset as u64);
516 }
517 }
518 }
519 }
520
521 for segment in &segments {
523 if let Some(adj) = segment.inserts.get(&(edge_type, direction)) {
524 for (vid, neighbors) in adj {
525 for &(neighbor, eid, version) in neighbors {
526 if !tombstoned_eids.contains(&eid) {
527 let offset = vid.as_u64();
528 entries.push((offset, neighbor, eid, version));
529 max_offset = max_offset.max(offset);
530 }
531 }
532 }
533 }
534 }
535
536 dedup_entries_by_eid(&mut entries);
537
538 let new_csr = MainCsr::from_edge_entries(max_offset as usize, entries);
540 let size = new_csr.memory_usage();
541
542 if let Some(old) = self.main_csr.get(&(edge_type, direction)) {
544 self.current_bytes
545 .fetch_sub(old.memory_usage(), Ordering::Relaxed);
546 }
547
548 self.main_csr
549 .insert((edge_type, direction), Arc::new(new_csr));
550 self.current_bytes.fetch_add(size, Ordering::Relaxed);
551 }
552
553 let snapshot_ptrs: HashSet<*const FrozenCsrSegment> =
559 segments.iter().map(Arc::as_ptr).collect();
560 self.frozen_segments
561 .write()
562 .retain(|s| !snapshot_ptrs.contains(&Arc::as_ptr(s)));
563 }
564
565 pub async fn warm(
572 &self,
573 storage: &StorageManager,
574 edge_type_id: u32,
575 direction: Direction,
576 version: Option<u64>,
577 ) -> anyhow::Result<()> {
578 let schema = storage.schema_manager().schema();
579
580 let edge_type_name = schema
582 .edge_type_name_by_id_unified(edge_type_id)
583 .ok_or_else(|| anyhow::anyhow!("Edge type {} not found", edge_type_id))?;
584
585 let labels_to_load: Vec<String> = {
587 let edge_meta = schema.edge_types.get(&edge_type_name);
588 match (direction, edge_meta) {
589 (Direction::Outgoing, Some(meta)) => meta.src_labels.clone(),
590 (Direction::Incoming, Some(meta)) => meta.dst_labels.clone(),
591 (Direction::Both, Some(meta)) => {
592 let mut labels = meta.src_labels.clone();
593 labels.extend(meta.dst_labels.iter().cloned());
594 labels.sort();
595 labels.dedup();
596 labels
597 }
598 _ => Vec::new(),
599 }
600 };
601
602 use arrow_array::{ListArray, UInt8Array, UInt64Array};
603
604 let mut entries: Vec<(u64, Vid, Eid, u64)> = Vec::new();
605 let mut deleted_eids = HashSet::new();
606
607 for &read_dir in direction.expand() {
608 let dir_str = read_dir.as_str();
609 for label_name in &labels_to_load {
610 let adj_ds = storage.adjacency_dataset(&edge_type_name, label_name, dir_str);
612 let backend = storage.backend();
613
614 if let Ok(adj_ds) = adj_ds {
615 let adj_table_name = adj_ds.table_name();
616 let adj_exists = backend.table_exists(&adj_table_name).await.unwrap_or(false);
617
618 if adj_exists {
619 let mut request = crate::backend::types::ScanRequest::all(&adj_table_name);
620 if let Some(hwm) = version {
621 request = request.with_filter(
622 crate::backend::types::FilterExpr::version_at_most(hwm),
623 );
624 }
625
626 let batches: Vec<arrow_array::RecordBatch> = backend.scan(request).await?;
631
632 for batch in batches {
633 let src_col = batch
634 .column_by_name("src_vid")
635 .unwrap()
636 .as_any()
637 .downcast_ref::<UInt64Array>()
638 .unwrap();
639 let neighbors_list = batch
640 .column_by_name("neighbors")
641 .unwrap()
642 .as_any()
643 .downcast_ref::<ListArray>()
644 .unwrap();
645 let eids_list = batch
646 .column_by_name("edge_ids")
647 .unwrap()
648 .as_any()
649 .downcast_ref::<ListArray>()
650 .unwrap();
651
652 for i in 0..batch.num_rows() {
653 let src_offset = src_col.value(i);
654 let neighbors_array_ref = neighbors_list.value(i);
655 let neighbors = neighbors_array_ref
656 .as_any()
657 .downcast_ref::<UInt64Array>()
658 .unwrap();
659 let eids_array_ref = eids_list.value(i);
660 let eids = eids_array_ref
661 .as_any()
662 .downcast_ref::<UInt64Array>()
663 .unwrap();
664
665 for j in 0..neighbors.len() {
666 entries.push((
672 src_offset,
673 Vid::from(neighbors.value(j)),
674 Eid::from(eids.value(j)),
675 0,
676 ));
677 }
678 }
679 }
680 }
681 }
682 }
683
684 let delta_ds = storage.delta_dataset(&edge_type_name, dir_str)?;
686 let backend = storage.backend();
687 let delta_table_name = delta_ds.table_name();
688
689 if backend
690 .table_exists(&delta_table_name)
691 .await
692 .unwrap_or(false)
693 {
694 let mut request = crate::backend::types::ScanRequest::all(&delta_table_name);
695 if let Some(hwm) = version {
696 request = request
697 .with_filter(crate::backend::types::FilterExpr::version_at_most(hwm));
698 }
699
700 let batches = backend.scan(request).await?;
704 {
705 for batch in batches {
706 let src_col = batch
707 .column_by_name("src_vid")
708 .unwrap()
709 .as_any()
710 .downcast_ref::<UInt64Array>()
711 .unwrap();
712 let dst_col = batch
713 .column_by_name("dst_vid")
714 .unwrap()
715 .as_any()
716 .downcast_ref::<UInt64Array>()
717 .unwrap();
718 let eid_col = batch
719 .column_by_name("eid")
720 .unwrap()
721 .as_any()
722 .downcast_ref::<UInt64Array>()
723 .unwrap();
724 let op_col = batch
725 .column_by_name("op")
726 .unwrap()
727 .as_any()
728 .downcast_ref::<UInt8Array>()
729 .unwrap();
730
731 let version_col = batch
733 .column_by_name("_version")
734 .and_then(|c| c.as_any().downcast_ref::<UInt64Array>().cloned());
735
736 for i in 0..batch.num_rows() {
737 let src_vid = Vid::from(src_col.value(i));
738 let dst_vid = Vid::from(dst_col.value(i));
739 let eid = Eid::from(eid_col.value(i));
740 let op = op_col.value(i); let row_version = version_col.as_ref().map_or(0, |vc| vc.value(i));
742
743 let is_incoming = read_dir == Direction::Incoming;
746 let (key_vid, neighbor_vid) = if is_incoming {
747 (dst_vid, src_vid)
748 } else {
749 (src_vid, dst_vid)
750 };
751
752 if op == 0 {
753 entries.push((key_vid.as_u64(), neighbor_vid, eid, row_version));
754 } else {
755 deleted_eids.insert(eid);
756 self.shadow.add_deleted_edge(
757 key_vid,
758 ShadowEdge {
759 neighbor_vid,
760 eid,
761 edge_type: edge_type_id,
762 created_version: 0,
763 deleted_version: row_version,
764 },
765 read_dir,
766 );
767 }
768 }
769 }
770 }
771 }
772 }
773
774 if !deleted_eids.is_empty() {
776 entries.retain(|(_, _, eid, _)| !deleted_eids.contains(eid));
777 }
778
779 dedup_entries_by_eid(&mut entries);
780
781 let max_offset = entries.iter().map(|(o, _, _, _)| *o).max().unwrap_or(0);
783 let csr = MainCsr::from_edge_entries(max_offset as usize, entries);
784 self.set_main_csr(edge_type_id, direction, csr);
785
786 Ok(())
787 }
788
789 pub async fn warm_coalesced(
795 &self,
796 storage: &StorageManager,
797 edge_type_id: u32,
798 direction: Direction,
799 version: Option<u64>,
800 ) -> anyhow::Result<()> {
801 if self.has_csr(edge_type_id, direction) {
803 return Ok(());
804 }
805
806 let guard = self
808 .warm_guards
809 .entry((edge_type_id, direction))
810 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
811 .value()
812 .clone();
813 let _lock = guard.lock().await;
814
815 if self.has_csr(edge_type_id, direction) {
817 return Ok(());
818 }
819
820 self.warm(storage, edge_type_id, direction, version).await
821 }
822
823 pub fn memory_usage(&self) -> usize {
825 self.current_bytes.load(Ordering::Relaxed) + self.shadow.approx_bytes()
832 }
833
834 pub fn pinned_versions(&self) -> &Arc<PinnedVersions> {
836 &self.pinned_versions
837 }
838
839 pub fn gc_shadow(&self, current_version: u64) {
850 let floor = self
851 .pinned_versions
852 .min_pinned()
853 .map_or(current_version, |pinned| pinned.min(current_version));
854 self.shadow.gc(floor);
855 }
856
857 pub fn shadow_entry_count(&self) -> usize {
862 self.shadow.entry_count()
863 }
864
865 pub fn max_bytes(&self) -> usize {
867 self.max_bytes
868 }
869
870 pub fn shadow(&self) -> &ShadowCsr {
872 &self.shadow
873 }
874}
875
876impl std::fmt::Debug for AdjacencyManager {
877 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
878 f.debug_struct("AdjacencyManager")
879 .field("main_csr_count", &self.main_csr.len())
880 .field("frozen_segments", &self.frozen_segments.read().len())
881 .field("current_bytes", &self.current_bytes.load(Ordering::Relaxed))
882 .field("max_bytes", &self.max_bytes)
883 .finish()
884 }
885}
886
887#[cfg(test)]
888mod tests {
889 use super::*;
890
891 #[test]
892 fn test_insert_and_get_neighbors() {
893 let am = AdjacencyManager::new(1024 * 1024);
894 let src = Vid::new(1);
895 let dst = Vid::new(2);
896 let eid = Eid::new(100);
897
898 am.insert_edge(src, dst, eid, 1, 1);
899
900 let neighbors = am.get_neighbors(src, 1, Direction::Outgoing);
901 assert_eq!(neighbors.len(), 1);
902 assert_eq!(neighbors[0], (dst, eid));
903
904 let incoming = am.get_neighbors(dst, 1, Direction::Incoming);
906 assert_eq!(incoming.len(), 1);
907 assert_eq!(incoming[0], (src, eid));
908 }
909
910 #[test]
911 fn test_main_csr_lookup() {
912 let am = AdjacencyManager::new(1024 * 1024);
913
914 let csr = MainCsr::from_edge_entries(
915 1,
916 vec![
917 (0, Vid::new(10), Eid::new(100), 1),
918 (1, Vid::new(20), Eid::new(101), 2),
919 ],
920 );
921 am.set_main_csr(1, Direction::Outgoing, csr);
922
923 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
924 assert_eq!(n.len(), 1);
925 assert_eq!(n[0], (Vid::new(10), Eid::new(100)));
926 }
927
928 #[test]
929 fn test_overlay_on_top_of_main_csr() {
930 let am = AdjacencyManager::new(1024 * 1024);
931
932 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
934 am.set_main_csr(1, Direction::Outgoing, csr);
935
936 am.insert_edge(Vid::new(0), Vid::new(20), Eid::new(101), 1, 2);
938
939 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
940 assert_eq!(n.len(), 2);
941
942 let eids: HashSet<Eid> = n.iter().map(|(_, e)| *e).collect();
943 assert!(eids.contains(&Eid::new(100)));
944 assert!(eids.contains(&Eid::new(101)));
945 }
946
947 #[test]
948 fn test_tombstone_removes_edge() {
949 let am = AdjacencyManager::new(1024 * 1024);
950
951 am.insert_edge(Vid::new(0), Vid::new(10), Eid::new(100), 1, 1);
952 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
953
954 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
955 assert!(n.is_empty());
956 }
957
958 #[test]
959 fn test_version_filtered_query() {
960 let am = AdjacencyManager::new(1024 * 1024);
961
962 let csr = MainCsr::from_edge_entries(
964 0,
965 vec![
966 (0, Vid::new(10), Eid::new(100), 1),
967 (0, Vid::new(20), Eid::new(101), 5),
968 ],
969 );
970 am.set_main_csr(1, Direction::Outgoing, csr);
971
972 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
974 assert_eq!(n.len(), 1);
975 assert_eq!(n[0], (Vid::new(10), Eid::new(100)));
976
977 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
979 assert_eq!(n.len(), 2);
980 }
981
982 #[test]
983 fn test_shadow_csr_resurrects_deleted_edges() {
984 let am = AdjacencyManager::new(1024 * 1024);
985
986 am.shadow().add_deleted_edge(
988 Vid::new(0),
989 ShadowEdge {
990 neighbor_vid: Vid::new(10),
991 eid: Eid::new(100),
992 edge_type: 1,
993 created_version: 1,
994 deleted_version: 5,
995 },
996 Direction::Outgoing,
997 );
998
999 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
1001 assert_eq!(n.len(), 1);
1002 assert_eq!(n[0], (Vid::new(10), Eid::new(100)));
1003
1004 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
1006 assert!(n.is_empty());
1007 }
1008
1009 #[test]
1015 fn test_concurrent_compaction_conserves_edges() {
1016 let am = std::sync::Arc::new(AdjacencyManager::new(64 * 1024 * 1024));
1017 let n: u64 = 150;
1018
1019 let inserter = {
1020 let am = am.clone();
1021 std::thread::spawn(move || {
1022 for i in 1..=n {
1023 am.insert_edge(Vid::new(0), Vid::new(i), Eid::new(i), 1, i);
1024 if i % 8 == 0 {
1025 am.compact();
1026 }
1027 }
1028 })
1029 };
1030 let compactor = {
1031 let am = am.clone();
1032 std::thread::spawn(move || {
1033 for _ in 0..40 {
1034 am.compact();
1035 std::thread::yield_now();
1036 }
1037 })
1038 };
1039 inserter.join().unwrap();
1040 compactor.join().unwrap();
1041 am.compact();
1042
1043 let neighbors = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1044 let got: HashSet<u64> = neighbors.iter().map(|(v, _)| v.as_u64()).collect();
1045 for i in 1..=n {
1046 assert!(
1047 got.contains(&i),
1048 "edge to {i} was lost under concurrent compaction"
1049 );
1050 }
1051 assert_eq!(got.len(), n as usize, "no spurious or duplicate neighbors");
1052 }
1053
1054 #[test]
1055 fn test_compact_merges_into_main_csr() {
1056 let am = AdjacencyManager::new(1024 * 1024);
1057
1058 am.insert_edge(Vid::new(0), Vid::new(10), Eid::new(100), 1, 1);
1060 am.insert_edge(Vid::new(0), Vid::new(20), Eid::new(101), 1, 2);
1061
1062 am.compact();
1064
1065 assert_eq!(am.frozen_segment_count(), 0);
1067
1068 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1070 assert_eq!(n.len(), 2);
1071
1072 assert!(am.has_csr(1, Direction::Outgoing));
1073 }
1074
1075 #[test]
1076 fn test_compact_removes_tombstoned_edges() {
1077 let am = AdjacencyManager::new(1024 * 1024);
1078
1079 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1081 am.set_main_csr(1, Direction::Outgoing, csr);
1082
1083 am.insert_edge(Vid::new(0), Vid::new(20), Eid::new(101), 1, 2);
1085 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 3);
1086
1087 am.compact();
1088
1089 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1091 assert_eq!(n.len(), 1);
1092 assert_eq!(n[0], (Vid::new(20), Eid::new(101)));
1093 }
1094
1095 #[test]
1096 fn test_should_compact() {
1097 let am = AdjacencyManager::new(1024 * 1024);
1098 assert!(!am.should_compact(4));
1099
1100 for _ in 0..4 {
1102 let frozen = {
1103 let mut active = am.active_overlay.write();
1104 let old = std::mem::take(&mut *active);
1105 Arc::new(old.freeze())
1106 };
1107 am.frozen_segments.write().push(frozen);
1108 }
1109
1110 assert!(am.should_compact(4));
1111 }
1112
1113 #[test]
1114 fn test_empty_manager() {
1115 let am = AdjacencyManager::new(1024 * 1024);
1116 assert!(
1117 am.get_neighbors(Vid::new(0), 1, Direction::Outgoing)
1118 .is_empty()
1119 );
1120 assert!(!am.has_csr(1, Direction::Outgoing));
1121 }
1122
1123 #[test]
1124 fn test_overlay_tombstone_removes_main_csr_edge() {
1125 let am = AdjacencyManager::new(1024 * 1024);
1127
1128 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1130 am.set_main_csr(1, Direction::Outgoing, csr);
1131
1132 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1134 assert_eq!(n.len(), 1);
1135
1136 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
1138
1139 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1141 assert!(
1142 n.is_empty(),
1143 "Edge should be removed by overlay tombstone, got {:?}",
1144 n
1145 );
1146 }
1147
1148 #[test]
1149 fn test_overlay_tombstone_removes_main_csr_edge_versioned() {
1150 let am = AdjacencyManager::new(1024 * 1024);
1152
1153 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1154 am.set_main_csr(1, Direction::Outgoing, csr);
1155
1156 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 5);
1157
1158 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
1160 assert_eq!(n.len(), 1);
1161
1162 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
1164 assert!(
1165 n.is_empty(),
1166 "Edge should be removed by overlay tombstone at version 5"
1167 );
1168 }
1169
1170 #[test]
1171 fn test_frozen_tombstone_removes_main_csr_edge() {
1172 let am = AdjacencyManager::new(1024 * 1024);
1174
1175 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1176 am.set_main_csr(1, Direction::Outgoing, csr);
1177
1178 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
1180
1181 {
1183 let mut active = am.active_overlay.write();
1184 let old = std::mem::take(&mut *active);
1185 let frozen = std::sync::Arc::new(old.freeze());
1186 am.frozen_segments.write().push(frozen);
1187 }
1188
1189 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1191 assert!(n.is_empty(), "Frozen tombstone should remove Main CSR edge");
1192 }
1193
1194 #[test]
1195 fn test_per_edge_version_filtering() {
1196 let am = AdjacencyManager::new(1024 * 1024);
1199
1200 let src = Vid::new(0);
1201 let dst_a = Vid::new(10);
1202 let dst_b = Vid::new(20);
1203 let eid_a = Eid::new(100);
1204 let eid_b = Eid::new(200);
1205 let etype = 1;
1206
1207 am.insert_edge(src, dst_a, eid_a, etype, 3);
1209
1210 am.insert_edge(src, dst_b, eid_b, etype, 7);
1212
1213 let neighbors_v2 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 2);
1215 assert!(
1216 neighbors_v2.is_empty(),
1217 "No edges should be visible at version 2"
1218 );
1219
1220 let neighbors_v5 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 5);
1222 assert_eq!(
1223 neighbors_v5.len(),
1224 1,
1225 "Only edge A should be visible at version 5"
1226 );
1227 assert_eq!(neighbors_v5[0].0, dst_a, "Edge A destination should match");
1228 assert_eq!(neighbors_v5[0].1, eid_a, "Edge A ID should match");
1229
1230 let neighbors_v7 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 7);
1232 assert_eq!(
1233 neighbors_v7.len(),
1234 2,
1235 "Both edges should be visible at version 7"
1236 );
1237
1238 let neighbors_v10 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 10);
1240 assert_eq!(
1241 neighbors_v10.len(),
1242 2,
1243 "Both edges should be visible at version 10"
1244 );
1245 }
1246
1247 #[test]
1248 fn test_duplicate_edges_deduplicated_by_eid() {
1249 let am = AdjacencyManager::new(1024 * 1024);
1251
1252 let src = Vid::new(0);
1253 let dst = Vid::new(10);
1254 let eid = Eid::new(100);
1255 let etype = 1;
1256
1257 let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1259 am.set_main_csr(etype, Direction::Outgoing, csr);
1260
1261 am.insert_edge(src, dst, eid, etype, 3);
1263
1264 let neighbors = am.get_neighbors(src, etype, Direction::Outgoing);
1266 assert_eq!(
1267 neighbors.len(),
1268 1,
1269 "Duplicate Eid should result in single entry"
1270 );
1271 assert_eq!(neighbors[0], (dst, eid));
1272 }
1273
1274 #[test]
1275 fn test_compact_deduplicates_edges_keeps_highest_version() {
1276 let am = AdjacencyManager::new(1024 * 1024);
1280
1281 let src = Vid::new(0);
1282 let dst = Vid::new(10);
1283 let eid = Eid::new(100);
1284 let etype = 1;
1285
1286 let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1288 am.set_main_csr(etype, Direction::Outgoing, csr);
1289
1290 am.insert_edge(src, dst, eid, etype, 5);
1292
1293 am.compact();
1297
1298 let neighbors_v5 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 5);
1300 assert_eq!(neighbors_v5.len(), 1, "Edge should be visible at version 5");
1301 assert_eq!(neighbors_v5[0], (dst, eid));
1302
1303 let neighbors_v4 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 4);
1307 assert_eq!(
1308 neighbors_v4.len(),
1309 0,
1310 "After compaction, only version 5 exists; version 4 should not see it"
1311 );
1312
1313 let neighbors_v1 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 1);
1315 assert_eq!(
1316 neighbors_v1.len(),
1317 0,
1318 "Old version discarded during compaction deduplication"
1319 );
1320
1321 let neighbors_v6 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 6);
1323 assert_eq!(neighbors_v6.len(), 1, "Edge should be visible at version 6");
1324 }
1325
1326 #[test]
1329 fn test_tombstone_scan_performance() {
1330 let am = AdjacencyManager::new(1024 * 1024);
1331 let vertex_a = Vid::new(1);
1332 let vertex_b = Vid::new(2);
1333 let etype = 1;
1334
1335 let mut a_edges = Vec::new();
1337 for i in 0..5 {
1338 let dst = Vid::new(100 + i);
1339 let eid = Eid::new(1000 + i);
1340 am.insert_edge(vertex_a, dst, eid, etype, 1);
1341 a_edges.push((dst, eid));
1342 }
1343
1344 for i in 0..100 {
1346 let dst = Vid::new(200 + i);
1347 let eid = Eid::new(2000 + i);
1348 am.insert_edge(vertex_b, dst, eid, etype, 1);
1349 am.add_tombstone(eid, vertex_b, dst, etype, 2);
1350 }
1351
1352 let neighbors = am.get_neighbors(vertex_a, etype, Direction::Outgoing);
1356
1357 assert_eq!(
1359 neighbors.len(),
1360 5,
1361 "Should return all 5 edges from vertex_a"
1362 );
1363 for (dst, eid) in &a_edges {
1364 assert!(
1365 neighbors.contains(&(*dst, *eid)),
1366 "Edge {:?} should be in results",
1367 (dst, eid)
1368 );
1369 }
1370
1371 let b_neighbors = am.get_neighbors(vertex_b, etype, Direction::Outgoing);
1373 assert_eq!(
1374 b_neighbors.len(),
1375 0,
1376 "Vertex B should have no neighbors (all deleted)"
1377 );
1378 }
1379
1380 #[test]
1386 fn test_get_neighbors_skips_irrelevant_segments() {
1387 let am = AdjacencyManager::new(1024 * 1024);
1388 let participant = Vid::new(1);
1389 let session = Vid::new(2);
1390 let link_eid = Eid::new(100);
1391 let link_etype: u32 = 1;
1392 let unrelated_etype: u32 = 2;
1393
1394 for i in 0..50 {
1398 if i == 17 {
1399 am.insert_edge(participant, session, link_eid, link_etype, i as u64 + 1);
1400 } else {
1401 let src = Vid::new(1000 + i as u64);
1403 let dst = Vid::new(2000 + i as u64);
1404 let eid = Eid::new(10_000 + i as u64);
1405 am.insert_edge(src, dst, eid, unrelated_etype, i as u64 + 1);
1406 }
1407 let frozen = {
1409 let mut active = am.active_overlay.write();
1410 let old = std::mem::take(&mut *active);
1411 Arc::new(old.freeze())
1412 };
1413 am.frozen_segments.write().push(frozen);
1414 }
1415
1416 assert_eq!(am.frozen_segment_count(), 50);
1418
1419 let n = am.get_neighbors(participant, link_etype, Direction::Outgoing);
1421 assert_eq!(n.len(), 1);
1422 assert_eq!(n[0], (session, link_eid));
1423
1424 let n_at = am.get_neighbors_at_version(participant, link_etype, Direction::Outgoing, 100);
1426 assert_eq!(n_at.len(), 1);
1427 assert_eq!(n_at[0], (session, link_eid));
1428
1429 let n_before =
1432 am.get_neighbors_at_version(participant, link_etype, Direction::Outgoing, 17);
1433 assert!(n_before.is_empty());
1434
1435 let unrelated = am.get_neighbors(Vid::new(1018), unrelated_etype, Direction::Outgoing);
1440 assert_eq!(unrelated.len(), 1);
1441 }
1442
1443 #[test]
1448 fn test_tombstone_in_unrelated_segment_still_applied() {
1449 let am = AdjacencyManager::new(1024 * 1024);
1450 let src = Vid::new(0);
1451 let dst = Vid::new(10);
1452 let eid = Eid::new(100);
1453 let etype: u32 = 1;
1454
1455 let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1457 am.set_main_csr(etype, Direction::Outgoing, csr);
1458
1459 am.add_tombstone(eid, src, dst, etype, 2);
1461
1462 let frozen = {
1467 let mut active = am.active_overlay.write();
1468 let old = std::mem::take(&mut *active);
1469 Arc::new(old.freeze())
1470 };
1471 am.frozen_segments.write().push(frozen);
1472
1473 let n = am.get_neighbors(src, etype, Direction::Outgoing);
1474 assert!(
1475 n.is_empty(),
1476 "tombstone in frozen segment must still hide Main CSR edge"
1477 );
1478 }
1479}