uni_store/storage/adjacency_manager.rs
1// SPDX-License-Identifier: Apache-2.0
2// Copyright 2024-2026 Dragonscale Team
3
4//! Unified adjacency manager orchestrating Main CSR, L0-csr overlay, and Shadow CSR.
5//!
6//! Implements a dual-CSR architecture where:
7//! - **Main CSR**: packed adjacency for all alive edges (one per edge_type + direction)
8//! - **L0-csr overlay**: concurrent insert/delete buffer that survives data flush
9//! - **Shadow CSR**: tracks deleted edges with version ranges for time-travel queries
10//!
11//! Regular queries read Main CSR + overlay with zero version filtering.
12//! Snapshot queries additionally filter by version and resurrect shadow entries.
13//!
14//! # Read-path cost model
15//!
16//! Every read consults three layers in order:
17//! 1. **Main CSR** — single `DashMap` lookup keyed by `(edge_type, direction)` →
18//! O(out-degree) entry scan.
19//! 2. **Frozen segments** — a `Vec<Arc<FrozenCsrSegment>>`. Each frozen segment is a
20//! plain `HashMap` (lock-free read). The read iterates this vec; each segment is
21//! short-circuited via [`FrozenCsrSegment::has_entries_for`] when it contributed
22//! no edges of the queried `(edge_type, direction)`, AND its tombstones map is
23//! empty. Both checks are O(1).
24//! 3. **Active overlay** — single `RwLock<L0CsrSegment>` taken once per read; same
25//! short-circuits applied as for frozen segments.
26//!
27//! Frozen segments are pure RAM. The on-disk Lance L1 delta is a write-through redo
28//! log, not consulted on the read hot path. Frozen segments are merged back into the
29//! Main CSR by [`AdjacencyManager::compact`], spawned from the writer's flush path
30//! once the segment count exceeds a small threshold.
31//!
32//! When extending the read path, preserve the invariant that any segment carrying a
33//! tombstone — even for an unrelated edge type — must run the `retain` pass against
34//! `result`, because tombstones can shadow edges materialized from layers below.
35
36use 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/// Versions currently pinned by in-flight snapshot readers.
49///
50/// Exists so shadow-CSR entries can be reclaimed without dropping any a live
51/// reader still resolves. The bound is subtle and worth stating once:
52///
53/// * `StorageManager::pinned()` and `at_fork` each build a **fresh**
54/// [`AdjacencyManager`] with its own empty `ShadowCsr`, so those readers
55/// never consult the live instance.
56/// * `StorageManager::pinned_at_version` **shares** the live manager, and is
57/// the path every read-write transaction takes.
58///
59/// So the only readers of the live shadow are in-flight `pinned_at_version`
60/// views, and the safe floor is the minimum version among them.
61/// `SnapshotManager` cannot supply this — it is a manifest reader-writer and
62/// tracks no live readers.
63///
64/// Refcounted rather than a set: several transactions routinely pin the same
65/// version, and the floor must not rise until the last of them is gone.
66#[derive(Debug, Default)]
67pub struct PinnedVersions {
68 counts: parking_lot::Mutex<std::collections::BTreeMap<u64, usize>>,
69}
70
71impl PinnedVersions {
72 /// Register `version` as pinned until the returned guard drops.
73 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 /// The lowest version any in-flight reader is pinned at.
82 pub fn min_pinned(&self) -> Option<u64> {
83 self.counts.lock().keys().next().copied()
84 }
85
86 /// Number of distinct pinned versions; for tests and diagnostics.
87 pub fn distinct_pinned(&self) -> usize {
88 self.counts.lock().len()
89 }
90}
91
92/// Releases its version from [`PinnedVersions`] on drop.
93#[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
111/// Unified adjacency manager for the dual-CSR architecture.
112///
113/// Orchestrates Main CSR (packed alive edges), L0-csr overlay
114/// (in-memory mutations), and Shadow CSR (deleted edges for time-travel).
115/// Data flush never invalidates or rebuilds the CSR.
116pub struct AdjacencyManager {
117 /// Main CSR per `(edge_type, direction)` — all alive edges.
118 /// Edge type is u32 with bit 31 = 0 for schema'd, 1 for schemaless.
119 main_csr: DashMap<(u32, Direction), Arc<MainCsr>>,
120
121 /// Active L0-csr segment (current writes go here).
122 active_overlay: Arc<RwLock<L0CsrSegment>>,
123
124 /// Frozen segments awaiting compaction (oldest first).
125 frozen_segments: RwLock<Vec<Arc<FrozenCsrSegment>>>,
126
127 /// Shadow CSR for time-travel deleted edge tracking.
128 shadow: ShadowCsr,
129 /// Versions pinned by in-flight `pinned_at_version` readers; the floor for
130 /// shadow GC. See [`PinnedVersions`].
131 pinned_versions: Arc<PinnedVersions>,
132
133 /// Current approximate memory usage in bytes.
134 current_bytes: AtomicUsize,
135
136 /// Maximum memory budget in bytes.
137 max_bytes: usize,
138
139 /// Coalescing locks for warm() operations — prevents cache stampede.
140 /// Key: (edge_type_id, Direction), Value: Mutex guard for that warm operation.
141 warm_guards: DashMap<(u32, Direction), Arc<tokio::sync::Mutex<()>>>,
142
143 /// Serializes `compact()` so two compactions can't interleave their
144 /// freeze→snapshot→clear sequences and lose a frozen segment. (review H12)
145 compact_lock: parking_lot::Mutex<()>,
146}
147
148impl AdjacencyManager {
149 /// Creates a new adjacency manager with the given memory budget.
150 pub fn new(max_bytes: usize) -> Self {
151 Self {
152 main_csr: DashMap::new(),
153 active_overlay: Arc::new(RwLock::new(L0CsrSegment::new())),
154 frozen_segments: RwLock::new(Vec::new()),
155 shadow: ShadowCsr::new(),
156 pinned_versions: Arc::new(PinnedVersions::default()),
157 current_bytes: AtomicUsize::new(0),
158 max_bytes,
159 warm_guards: DashMap::new(),
160 compact_lock: parking_lot::Mutex::new(()),
161 }
162 }
163
164 /// Returns neighbors for the current state (hot path, no version filtering).
165 ///
166 /// Reads Main CSR + frozen segments + active overlay, minus tombstones.
167 /// Tombstones from any layer remove edges from all lower layers.
168 pub fn get_neighbors(&self, vid: Vid, edge_type: u32, direction: Direction) -> Vec<(Vid, Eid)> {
169 let mut result: HashMap<Eid, Vid> = HashMap::new();
170
171 for &dir in direction.expand() {
172 // 1. Main CSR
173 if let Some(csr) = self.main_csr.get(&(edge_type, dir)) {
174 for entry in csr.get_entries(vid) {
175 result.insert(entry.eid, entry.neighbor_vid);
176 }
177 }
178
179 // 2. Frozen segments (oldest first) — add inserts, then remove tombstones.
180 // Skip whole segment when it has no inserts for this (edge_type, dir)
181 // AND no tombstones at all — both is_empty checks are O(1) on plain HashMaps,
182 // so this is a strict speedup that scales away with frozen-segment count
183 // when most segments don't touch the queried edge type. See issue #55.
184 for segment in self.frozen_segments.read().iter() {
185 let has_inserts = segment.has_entries_for(edge_type, dir);
186 let has_tombstones = !segment.tombstones.is_empty();
187 if !has_inserts && !has_tombstones {
188 continue;
189 }
190 if has_inserts
191 && let Some(adj) = segment.inserts.get(&(edge_type, dir))
192 && let Some(neighbors) = adj.get(&vid)
193 {
194 for &(neighbor, eid, _version) in neighbors {
195 result.insert(eid, neighbor);
196 }
197 }
198 // Apply tombstones against ALL prior results (Main CSR + older segments).
199 // Skip the retain pass entirely when there are no tombstones — it would
200 // be a no-op but still O(result_size) due to the closure call.
201 if has_tombstones {
202 result.retain(|eid, _| !segment.tombstones.contains_key(eid));
203 }
204 }
205
206 // 3. Active overlay — add inserts, then remove tombstones.
207 // Same short-circuits as the frozen branch, plus we hold the read lock
208 // for the whole branch so DashMap atomic-op cost is paid once, not twice.
209 let active = self.active_overlay.read();
210 let active_has_inserts = active.has_entries_for(edge_type, dir);
211 let active_has_tombstones = !active.tombstones.is_empty();
212 if active_has_inserts
213 && let Some(adj) = active.inserts.get(&(edge_type, dir))
214 && let Some(neighbors) = adj.get(&vid)
215 {
216 for &(neighbor, eid, _version) in neighbors {
217 result.insert(eid, neighbor);
218 }
219 }
220 if active_has_tombstones {
221 result.retain(|eid, _| !active.tombstones.contains_key(eid));
222 }
223 }
224
225 result.into_iter().map(|(e, n)| (n, e)).collect()
226 }
227
228 /// Returns neighbors visible at a specific snapshot version.
229 ///
230 /// Filters Main CSR entries by `created_version`, applies frozen/active
231 /// overlay with version filtering, and resurrects Shadow CSR entries
232 /// that were alive at the given version.
233 pub fn get_neighbors_at_version(
234 &self,
235 vid: Vid,
236 edge_type: u32,
237 direction: Direction,
238 version: u64,
239 ) -> Vec<(Vid, Eid)> {
240 let mut result: HashMap<Eid, Vid> = HashMap::new();
241
242 for &dir in direction.expand() {
243 // 1. Main CSR — filter by created_version
244 if let Some(csr) = self.main_csr.get(&(edge_type, dir)) {
245 for entry in csr.get_entries(vid) {
246 if entry.created_version <= version {
247 result.insert(entry.eid, entry.neighbor_vid);
248 }
249 }
250 }
251
252 // 2. Frozen segments — filter inserts by version, apply tombstones.
253 // Same skip-irrelevant-segment short-circuit as get_neighbors. See issue #55.
254 for segment in self.frozen_segments.read().iter() {
255 let has_inserts = segment.has_entries_for(edge_type, dir);
256 let has_tombstones = !segment.tombstones.is_empty();
257 if !has_inserts && !has_tombstones {
258 continue;
259 }
260 if has_inserts
261 && let Some(adj) = segment.inserts.get(&(edge_type, dir))
262 && let Some(neighbors) = adj.get(&vid)
263 {
264 for &(neighbor, eid, ver) in neighbors {
265 if ver <= version {
266 result.insert(eid, neighbor);
267 }
268 }
269 }
270 if has_tombstones {
271 result.retain(|eid, _| {
272 segment
273 .tombstones
274 .get(eid)
275 .is_none_or(|ts| ts.version > version)
276 });
277 }
278 }
279
280 // 3. Active overlay — add version-filtered inserts, then apply tombstones
281 let active = self.active_overlay.read();
282 let active_has_inserts = active.has_entries_for(edge_type, dir);
283 let active_has_tombstones = !active.tombstones.is_empty();
284 if active_has_inserts
285 && let Some(adj) = active.inserts.get(&(edge_type, dir))
286 && let Some(neighbors) = adj.get(&vid)
287 {
288 for &(neighbor, eid, ver) in neighbors {
289 let not_tombstoned = active
290 .tombstones
291 .get(&eid)
292 .is_none_or(|ts| ts.version > version);
293 if ver <= version && not_tombstoned {
294 result.insert(eid, neighbor);
295 }
296 }
297 }
298 if active_has_tombstones {
299 result.retain(|eid, _| {
300 active
301 .tombstones
302 .get(eid)
303 .is_none_or(|ts| ts.version > version)
304 });
305 }
306
307 // 4. Shadow CSR — resurrect edges alive at version
308 for (neighbor, eid) in self
309 .shadow
310 .get_entries_at_version(vid, edge_type, dir, version)
311 {
312 result.insert(eid, neighbor);
313 }
314 }
315
316 result.into_iter().map(|(e, n)| (n, e)).collect()
317 }
318
319 /// Records an edge insertion into the L0-csr overlay (both directions).
320 pub fn insert_edge(&self, src: Vid, dst: Vid, eid: Eid, edge_type: u32, version: u64) {
321 let active = self.active_overlay.read();
322 active.insert_edge(src, dst, eid, edge_type, version, Direction::Outgoing);
323 active.insert_edge(dst, src, eid, edge_type, version, Direction::Incoming);
324 }
325
326 /// Records a tombstone for a deleted edge in the L0-csr overlay.
327 pub fn add_tombstone(&self, eid: Eid, src: Vid, dst: Vid, edge_type: u32, version: u64) {
328 let active = self.active_overlay.read();
329 active.add_tombstone(eid, src, dst, edge_type, version);
330 }
331
332 /// Sets the Main CSR for a specific edge type and direction.
333 ///
334 /// Used by `warm()` to install a freshly built CSR from storage.
335 pub fn set_main_csr(&self, edge_type: u32, direction: Direction, csr: MainCsr) {
336 let size = csr.memory_usage();
337 self.main_csr.insert((edge_type, direction), Arc::new(csr));
338 self.current_bytes.fetch_add(size, Ordering::Relaxed);
339 }
340
341 /// Checks whether a Main CSR exists for the given edge type and direction.
342 pub fn has_csr(&self, edge_type: u32, direction: Direction) -> bool {
343 self.main_csr.contains_key(&(edge_type, direction))
344 }
345
346 /// Checks whether this manager has been activated for the given edge type.
347 ///
348 /// Returns `true` if a Main CSR exists or the overlay has entries for
349 /// this edge type and direction.
350 pub fn is_active_for(&self, edge_type: u32, direction: Direction) -> bool {
351 let active = self.active_overlay.read();
352 direction.expand().iter().any(|&d| {
353 self.main_csr.contains_key(&(edge_type, d)) || active.has_entries_for(edge_type, d)
354 })
355 }
356
357 /// Returns the distinct edge type ids known to this manager.
358 ///
359 /// Spans the Main CSR plus the active and frozen overlay segments, so it
360 /// covers both warmed (L1/L2-loaded) edge types and live overlay-resident
361 /// types (e.g. recently committed edges that a flush moved out of L0 but kept
362 /// in the dual-write overlay). Used when an endpoint resolver knows an edge
363 /// id but not its type: this is the small set of types this query has touched,
364 /// so it bounds an eid-orientation probe to a short candidate list rather than
365 /// the whole schema.
366 ///
367 /// # Examples
368 ///
369 /// ```ignore
370 /// for etype in adjacency_manager.known_edge_type_ids() {
371 /// // probe etype for the edge of interest
372 /// }
373 /// ```
374 #[must_use]
375 pub fn known_edge_type_ids(&self) -> Vec<u32> {
376 let mut ids: Vec<u32> = self.main_csr.iter().map(|entry| entry.key().0).collect();
377 for entry in self.active_overlay.read().inserts.iter() {
378 ids.push(entry.key().0);
379 }
380 for segment in self.frozen_segments.read().iter() {
381 ids.extend(segment.inserts.keys().map(|&(etype, _dir)| etype));
382 }
383 ids.sort_unstable();
384 ids.dedup();
385 ids
386 }
387
388 /// Returns the number of frozen segments awaiting compaction.
389 pub fn frozen_segment_count(&self) -> usize {
390 self.frozen_segments.read().len()
391 }
392
393 /// Returns whether compaction should be triggered based on segment count.
394 pub fn should_compact(&self, threshold: usize) -> bool {
395 self.frozen_segment_count() >= threshold
396 }
397
398 /// Compacts frozen overlay segments into the Main CSR.
399 ///
400 /// Freezes the active overlay, merges all frozen segments with the
401 /// existing Main CSR, moves tombstoned edges to Shadow CSR, and
402 /// atomically swaps in the new Main CSR.
403 ///
404 /// CRITICAL: Frozen segments remain readable until the new CSR is installed,
405 /// eliminating the visibility gap where edges would be invisible.
406 pub fn compact(&self) {
407 // Serialize compaction: two concurrent compacts would each freeze, take
408 // their own snapshot, then clear — the second clear wiping segments the
409 // first had not yet merged. (review H12)
410 let _compact_guard = self.compact_lock.lock();
411
412 // Step 1: Freeze active overlay and push to frozen list
413 let frozen = {
414 let mut active = self.active_overlay.write();
415 let old = std::mem::take(&mut *active);
416 Arc::new(old.freeze())
417 };
418 self.frozen_segments.write().push(frozen);
419
420 // Step 2: CLONE frozen segments for building (DON'T drain yet)
421 // This ensures they remain readable during CSR construction
422 let segments = self.frozen_segments.read().clone();
423
424 // Step 3: Collect all (edge_type, direction) keys from segments + existing CSRs
425 let mut all_keys: HashSet<(u32, Direction)> = HashSet::new();
426 for segment in &segments {
427 for key in segment.inserts.keys() {
428 all_keys.insert(*key);
429 }
430 }
431 for entry in self.main_csr.iter() {
432 all_keys.insert(*entry.key());
433 }
434
435 // Step 4: For each key, merge
436 for (edge_type, direction) in all_keys {
437 let mut entries: Vec<(u64, Vid, Eid, u64)> = Vec::new();
438 let mut max_offset: u64 = 0;
439
440 // Collect all tombstone EIDs
441 let mut tombstoned_eids: HashSet<Eid> = HashSet::new();
442 for segment in &segments {
443 for (eid, ts) in &segment.tombstones {
444 if ts.edge_type == edge_type {
445 tombstoned_eids.insert(*eid);
446
447 // Move to shadow CSR. ShadowCsr is keyed by the queried
448 // vid for the direction, exactly like the CSRs themselves
449 // (Incoming is keyed by dst with neighbor src — see the
450 // insert at `insert_edge(dst, src, .., Incoming)` and the
451 // get_neighbors_at_version swap). Key by src for Outgoing,
452 // by dst for Incoming — otherwise a time-travel read after
453 // compaction looks up the wrong vid and the tombstone is
454 // invisible (deleted edge resurrected).
455 let (key_vid, neighbor_vid) = if direction == Direction::Incoming {
456 (ts.dst_vid, ts.src_vid)
457 } else {
458 (ts.src_vid, ts.dst_vid)
459 };
460 self.shadow.add_deleted_edge(
461 key_vid,
462 ShadowEdge {
463 neighbor_vid,
464 eid: *eid,
465 edge_type,
466 created_version: 0, // unknown; overlay tombstones don't track creation version
467 deleted_version: ts.version,
468 },
469 direction,
470 );
471 }
472 }
473 }
474
475 // Add entries from old Main CSR
476 if let Some(old_csr) = self.main_csr.get(&(edge_type, direction)) {
477 for vid_offset in 0..old_csr.num_vertices() {
478 let vid = Vid::new(vid_offset as u64);
479 for entry in old_csr.get_entries(vid) {
480 if !tombstoned_eids.contains(&entry.eid) {
481 entries.push((
482 vid_offset as u64,
483 entry.neighbor_vid,
484 entry.eid,
485 entry.created_version,
486 ));
487 max_offset = max_offset.max(vid_offset as u64);
488 }
489 }
490 }
491 }
492
493 // Overlay frozen segments (oldest first)
494 for segment in &segments {
495 if let Some(adj) = segment.inserts.get(&(edge_type, direction)) {
496 for (vid, neighbors) in adj {
497 for &(neighbor, eid, version) in neighbors {
498 if !tombstoned_eids.contains(&eid) {
499 let offset = vid.as_u64();
500 entries.push((offset, neighbor, eid, version));
501 max_offset = max_offset.max(offset);
502 }
503 }
504 }
505 }
506 }
507
508 // Deduplicate by Eid — keep entry with highest version for each Eid
509 {
510 use std::collections::hash_map::Entry;
511
512 let mut best: HashMap<Eid, usize> = HashMap::new();
513 for (idx, (_, _, eid, ver)) in entries.iter().enumerate() {
514 match best.entry(*eid) {
515 Entry::Vacant(e) => {
516 e.insert(idx);
517 }
518 Entry::Occupied(mut e) => {
519 if *ver > entries[*e.get()].3 {
520 e.insert(idx);
521 }
522 }
523 }
524 }
525 let keep: HashSet<usize> = best.into_values().collect();
526 let mut idx = 0;
527 entries.retain(|_| {
528 let k = keep.contains(&idx);
529 idx += 1;
530 k
531 });
532 }
533
534 // Build new Main CSR and install
535 let new_csr = MainCsr::from_edge_entries(max_offset as usize, entries);
536 let size = new_csr.memory_usage();
537
538 // Remove old size, add new
539 if let Some(old) = self.main_csr.get(&(edge_type, direction)) {
540 self.current_bytes
541 .fetch_sub(old.memory_usage(), Ordering::Relaxed);
542 }
543
544 self.main_csr
545 .insert((edge_type, direction), Arc::new(new_csr));
546 self.current_bytes.fetch_add(size, Ordering::Relaxed);
547 }
548
549 // Step 5: drain EXACTLY the segments we snapshotted in Step 2 — not a
550 // blanket clear(). A concurrent `freeze()` may have pushed a new frozen
551 // segment after the snapshot; that segment was NOT merged into the new
552 // CSR, so clearing it would silently lose its topology. Retain anything
553 // not in the snapshot (compared by Arc identity). (review H12)
554 let snapshot_ptrs: HashSet<*const FrozenCsrSegment> =
555 segments.iter().map(Arc::as_ptr).collect();
556 self.frozen_segments
557 .write()
558 .retain(|s| !snapshot_ptrs.contains(&Arc::as_ptr(s)));
559 }
560
561 /// Warms the Main CSR from storage (L2 adjacency + L1 delta) for a specific edge type and direction.
562 ///
563 /// Reads L2 adjacency datasets and L1 delta entries from Lance,
564 /// builds a [`MainCsr`] with version metadata, and populates the
565 /// [`ShadowCsr`] with L1 tombstones. Called once at startup or
566 /// lazily on first access per edge type.
567 pub async fn warm(
568 &self,
569 storage: &StorageManager,
570 edge_type_id: u32,
571 direction: Direction,
572 version: Option<u64>,
573 ) -> anyhow::Result<()> {
574 let schema = storage.schema_manager().schema();
575
576 // Use unified lookup to support both schema'd and schemaless edge types
577 let edge_type_name = schema
578 .edge_type_name_by_id_unified(edge_type_id)
579 .ok_or_else(|| anyhow::anyhow!("Edge type {} not found", edge_type_id))?;
580
581 // Determine which labels to load adjacency for based on edge type metadata
582 let labels_to_load: Vec<String> = {
583 let edge_meta = schema.edge_types.get(&edge_type_name);
584 match (direction, edge_meta) {
585 (Direction::Outgoing, Some(meta)) => meta.src_labels.clone(),
586 (Direction::Incoming, Some(meta)) => meta.dst_labels.clone(),
587 (Direction::Both, Some(meta)) => {
588 let mut labels = meta.src_labels.clone();
589 labels.extend(meta.dst_labels.iter().cloned());
590 labels.sort();
591 labels.dedup();
592 labels
593 }
594 _ => Vec::new(),
595 }
596 };
597
598 use arrow_array::{ListArray, UInt8Array, UInt64Array};
599
600 let mut entries: Vec<(u64, Vid, Eid, u64)> = Vec::new();
601 let mut deleted_eids = HashSet::new();
602
603 for &read_dir in direction.expand() {
604 let dir_str = read_dir.as_str();
605 for label_name in &labels_to_load {
606 // 1. Read L2 (Adjacency Dataset)
607 let adj_ds = storage.adjacency_dataset(&edge_type_name, label_name, dir_str);
608 let backend = storage.backend();
609
610 if let Ok(adj_ds) = adj_ds {
611 let adj_table_name = adj_ds.table_name();
612 let adj_exists = backend.table_exists(&adj_table_name).await.unwrap_or(false);
613
614 if adj_exists {
615 let mut request = crate::backend::types::ScanRequest::all(&adj_table_name);
616 if let Some(hwm) = version {
617 request = request.with_filter(
618 crate::backend::types::FilterExpr::version_at_most(hwm),
619 );
620 }
621
622 // Fail closed: a transient scan error must abort the warm,
623 // not `unwrap_or_default()` into an empty L2 read that then
624 // gets cached as the adjacency CSR — that silently drops
625 // every base edge for this type until restart (review #3b).
626 let batches: Vec<arrow_array::RecordBatch> = backend.scan(request).await?;
627
628 for batch in batches {
629 let src_col = batch
630 .column_by_name("src_vid")
631 .unwrap()
632 .as_any()
633 .downcast_ref::<UInt64Array>()
634 .unwrap();
635 let neighbors_list = batch
636 .column_by_name("neighbors")
637 .unwrap()
638 .as_any()
639 .downcast_ref::<ListArray>()
640 .unwrap();
641 let eids_list = batch
642 .column_by_name("edge_ids")
643 .unwrap()
644 .as_any()
645 .downcast_ref::<ListArray>()
646 .unwrap();
647
648 for i in 0..batch.num_rows() {
649 let src_offset = src_col.value(i);
650 let neighbors_array_ref = neighbors_list.value(i);
651 let neighbors = neighbors_array_ref
652 .as_any()
653 .downcast_ref::<UInt64Array>()
654 .unwrap();
655 let eids_array_ref = eids_list.value(i);
656 let eids = eids_array_ref
657 .as_any()
658 .downcast_ref::<UInt64Array>()
659 .unwrap();
660
661 for j in 0..neighbors.len() {
662 // L2 adjacency rows don't carry per-edge _version.
663 // Version 0 means "from base storage" — the `_version <= hwm` filter on
664 // the query already ensures we only load rows within the snapshot window.
665 // At query time, get_neighbors_at_version() uses created_version to filter,
666 // so version=0 edges are always visible (which is correct for compacted L2 data).
667 entries.push((
668 src_offset,
669 Vid::from(neighbors.value(j)),
670 Eid::from(eids.value(j)),
671 0,
672 ));
673 }
674 }
675 }
676 }
677 }
678 }
679
680 // 2. Read L1 (Delta)
681 let delta_ds = storage.delta_dataset(&edge_type_name, dir_str)?;
682 let backend = storage.backend();
683 let delta_table_name = delta_ds.table_name();
684
685 if backend
686 .table_exists(&delta_table_name)
687 .await
688 .unwrap_or(false)
689 {
690 let mut request = crate::backend::types::ScanRequest::all(&delta_table_name);
691 if let Some(hwm) = version {
692 request = request
693 .with_filter(crate::backend::types::FilterExpr::version_at_most(hwm));
694 }
695
696 // Fail closed: propagate a delta scan error rather than silently
697 // skipping it, which would drop unflushed edges from the cached
698 // adjacency CSR until restart (review #3b).
699 let batches = backend.scan(request).await?;
700 {
701 for batch in batches {
702 let src_col = batch
703 .column_by_name("src_vid")
704 .unwrap()
705 .as_any()
706 .downcast_ref::<UInt64Array>()
707 .unwrap();
708 let dst_col = batch
709 .column_by_name("dst_vid")
710 .unwrap()
711 .as_any()
712 .downcast_ref::<UInt64Array>()
713 .unwrap();
714 let eid_col = batch
715 .column_by_name("eid")
716 .unwrap()
717 .as_any()
718 .downcast_ref::<UInt64Array>()
719 .unwrap();
720 let op_col = batch
721 .column_by_name("op")
722 .unwrap()
723 .as_any()
724 .downcast_ref::<UInt8Array>()
725 .unwrap();
726
727 // Optionally read _version column
728 let version_col = batch
729 .column_by_name("_version")
730 .and_then(|c| c.as_any().downcast_ref::<UInt64Array>().cloned());
731
732 for i in 0..batch.num_rows() {
733 let src_vid = Vid::from(src_col.value(i));
734 let dst_vid = Vid::from(dst_col.value(i));
735 let eid = Eid::from(eid_col.value(i));
736 let op = op_col.value(i); // 0=Insert, 1=Delete
737 let row_version = version_col.as_ref().map_or(0, |vc| vc.value(i));
738
739 // For incoming edges, the CSR key is dst (the vertex
740 // receiving the edge) and the neighbor is src.
741 let is_incoming = read_dir == Direction::Incoming;
742 let (key_vid, neighbor_vid) = if is_incoming {
743 (dst_vid, src_vid)
744 } else {
745 (src_vid, dst_vid)
746 };
747
748 if op == 0 {
749 entries.push((key_vid.as_u64(), neighbor_vid, eid, row_version));
750 } else {
751 deleted_eids.insert(eid);
752 self.shadow.add_deleted_edge(
753 key_vid,
754 ShadowEdge {
755 neighbor_vid,
756 eid,
757 edge_type: edge_type_id,
758 created_version: 0,
759 deleted_version: row_version,
760 },
761 read_dir,
762 );
763 }
764 }
765 }
766 }
767 }
768 }
769
770 // Filter out deleted edges
771 if !deleted_eids.is_empty() {
772 entries.retain(|(_, _, eid, _)| !deleted_eids.contains(eid));
773 }
774
775 // Deduplicate by Eid — keep entry with highest version for each Eid
776 // Multiple versions of the same edge can exist in L2+L1 or across L1 runs
777 {
778 use std::collections::hash_map::Entry;
779 use std::collections::{HashMap, HashSet};
780
781 let mut best: HashMap<Eid, usize> = HashMap::new();
782 for (idx, (_, _, eid, ver)) in entries.iter().enumerate() {
783 match best.entry(*eid) {
784 Entry::Vacant(e) => {
785 e.insert(idx);
786 }
787 Entry::Occupied(mut e) => {
788 if *ver > entries[*e.get()].3 {
789 e.insert(idx);
790 }
791 }
792 }
793 }
794 let keep: HashSet<usize> = best.into_values().collect();
795 let mut idx = 0;
796 entries.retain(|_| {
797 let k = keep.contains(&idx);
798 idx += 1;
799 k
800 });
801 }
802
803 // Build MainCsr
804 let max_offset = entries.iter().map(|(o, _, _, _)| *o).max().unwrap_or(0);
805 let csr = MainCsr::from_edge_entries(max_offset as usize, entries);
806 self.set_main_csr(edge_type_id, direction, csr);
807
808 Ok(())
809 }
810
811 /// Coalesced warm() operation to prevent cache stampede (Issue #13).
812 ///
813 /// Uses double-checked locking: fast-path checks if CSR already loaded,
814 /// then acquires per-(edge_type, direction) lock to ensure only one concurrent
815 /// warm() per adjacency key. Other readers wait for the first warm() to complete.
816 pub async fn warm_coalesced(
817 &self,
818 storage: &StorageManager,
819 edge_type_id: u32,
820 direction: Direction,
821 version: Option<u64>,
822 ) -> anyhow::Result<()> {
823 // Fast path: already loaded
824 if self.has_csr(edge_type_id, direction) {
825 return Ok(());
826 }
827
828 // Coalesce: only one concurrent warm per (type, dir)
829 let guard = self
830 .warm_guards
831 .entry((edge_type_id, direction))
832 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
833 .value()
834 .clone();
835 let _lock = guard.lock().await;
836
837 // Double-check after acquiring lock
838 if self.has_csr(edge_type_id, direction) {
839 return Ok(());
840 }
841
842 self.warm(storage, edge_type_id, direction, version).await
843 }
844
845 /// Returns the current approximate memory usage in bytes.
846 pub fn memory_usage(&self) -> usize {
847 // `current_bytes` accounts only for the main CSR — every mutation of it
848 // is on a `main_csr` path. Shadow retention was therefore invisible to
849 // the budget and could never trip `max_bytes`, which is how an
850 // unbounded leak there went unnoticed. Counted approximately rather
851 // than tracked incrementally: the shadow is small when healthy, and an
852 // exact counter would need hooks on every retain in `gc`.
853 self.current_bytes.load(Ordering::Relaxed) + self.shadow.approx_bytes()
854 }
855
856 /// The pinned-version registry backing shadow GC.
857 pub fn pinned_versions(&self) -> &Arc<PinnedVersions> {
858 &self.pinned_versions
859 }
860
861 /// Reclaim shadow entries no in-flight reader can reach.
862 ///
863 /// The floor is the minimum pinned version, or `current_version` when
864 /// nothing is pinned — a reader starting now pins at the current version,
865 /// so entries deleted at or below it are unreachable. `current_version` is
866 /// passed in because the manager does not track it; the writer does.
867 ///
868 /// Called after compaction. Safe to call at any time: it only ever removes
869 /// entries whose `deleted_version` is at or below the floor, which is
870 /// exactly the set `get_entries_at_version` can no longer return.
871 pub fn gc_shadow(&self, current_version: u64) {
872 let floor = self
873 .pinned_versions
874 .min_pinned()
875 .map_or(current_version, |pinned| pinned.min(current_version));
876 self.shadow.gc(floor);
877 }
878
879 /// Shadow-CSR entries currently retained.
880 ///
881 /// Exposed for retention tests and diagnostics; see
882 /// [`ShadowCsr::add_deleted_edge`] for why this can grow.
883 pub fn shadow_entry_count(&self) -> usize {
884 self.shadow.entry_count()
885 }
886
887 /// Returns the maximum memory budget in bytes.
888 pub fn max_bytes(&self) -> usize {
889 self.max_bytes
890 }
891
892 /// Provides access to the shadow CSR for time-travel queries.
893 pub fn shadow(&self) -> &ShadowCsr {
894 &self.shadow
895 }
896}
897
898impl std::fmt::Debug for AdjacencyManager {
899 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
900 f.debug_struct("AdjacencyManager")
901 .field("main_csr_count", &self.main_csr.len())
902 .field("frozen_segments", &self.frozen_segments.read().len())
903 .field("current_bytes", &self.current_bytes.load(Ordering::Relaxed))
904 .field("max_bytes", &self.max_bytes)
905 .finish()
906 }
907}
908
909#[cfg(test)]
910mod tests {
911 use super::*;
912
913 #[test]
914 fn test_insert_and_get_neighbors() {
915 let am = AdjacencyManager::new(1024 * 1024);
916 let src = Vid::new(1);
917 let dst = Vid::new(2);
918 let eid = Eid::new(100);
919
920 am.insert_edge(src, dst, eid, 1, 1);
921
922 let neighbors = am.get_neighbors(src, 1, Direction::Outgoing);
923 assert_eq!(neighbors.len(), 1);
924 assert_eq!(neighbors[0], (dst, eid));
925
926 // Incoming direction
927 let incoming = am.get_neighbors(dst, 1, Direction::Incoming);
928 assert_eq!(incoming.len(), 1);
929 assert_eq!(incoming[0], (src, eid));
930 }
931
932 #[test]
933 fn test_main_csr_lookup() {
934 let am = AdjacencyManager::new(1024 * 1024);
935
936 let csr = MainCsr::from_edge_entries(
937 1,
938 vec![
939 (0, Vid::new(10), Eid::new(100), 1),
940 (1, Vid::new(20), Eid::new(101), 2),
941 ],
942 );
943 am.set_main_csr(1, Direction::Outgoing, csr);
944
945 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
946 assert_eq!(n.len(), 1);
947 assert_eq!(n[0], (Vid::new(10), Eid::new(100)));
948 }
949
950 #[test]
951 fn test_overlay_on_top_of_main_csr() {
952 let am = AdjacencyManager::new(1024 * 1024);
953
954 // Main CSR has one edge
955 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
956 am.set_main_csr(1, Direction::Outgoing, csr);
957
958 // Overlay adds another
959 am.insert_edge(Vid::new(0), Vid::new(20), Eid::new(101), 1, 2);
960
961 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
962 assert_eq!(n.len(), 2);
963
964 let eids: HashSet<Eid> = n.iter().map(|(_, e)| *e).collect();
965 assert!(eids.contains(&Eid::new(100)));
966 assert!(eids.contains(&Eid::new(101)));
967 }
968
969 #[test]
970 fn test_tombstone_removes_edge() {
971 let am = AdjacencyManager::new(1024 * 1024);
972
973 am.insert_edge(Vid::new(0), Vid::new(10), Eid::new(100), 1, 1);
974 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
975
976 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
977 assert!(n.is_empty());
978 }
979
980 #[test]
981 fn test_version_filtered_query() {
982 let am = AdjacencyManager::new(1024 * 1024);
983
984 // Main CSR with two edges at different versions
985 let csr = MainCsr::from_edge_entries(
986 0,
987 vec![
988 (0, Vid::new(10), Eid::new(100), 1),
989 (0, Vid::new(20), Eid::new(101), 5),
990 ],
991 );
992 am.set_main_csr(1, Direction::Outgoing, csr);
993
994 // At version 3: only first edge visible
995 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
996 assert_eq!(n.len(), 1);
997 assert_eq!(n[0], (Vid::new(10), Eid::new(100)));
998
999 // At version 5: both visible
1000 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
1001 assert_eq!(n.len(), 2);
1002 }
1003
1004 #[test]
1005 fn test_shadow_csr_resurrects_deleted_edges() {
1006 let am = AdjacencyManager::new(1024 * 1024);
1007
1008 // Add a deleted edge to shadow: created at v1, deleted at v5
1009 am.shadow().add_deleted_edge(
1010 Vid::new(0),
1011 ShadowEdge {
1012 neighbor_vid: Vid::new(10),
1013 eid: Eid::new(100),
1014 edge_type: 1,
1015 created_version: 1,
1016 deleted_version: 5,
1017 },
1018 Direction::Outgoing,
1019 );
1020
1021 // At version 3: shadow edge should be visible
1022 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
1023 assert_eq!(n.len(), 1);
1024 assert_eq!(n[0], (Vid::new(10), Eid::new(100)));
1025
1026 // At version 5: deleted, not visible
1027 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
1028 assert!(n.is_empty());
1029 }
1030
1031 /// H12: concurrent compaction must never drop an edge. `compact()` is the
1032 /// only writer of `frozen_segments`, now serialized under `compact_lock`,
1033 /// and Step 5 drains exactly the snapshotted segments rather than clearing
1034 /// all — so a segment frozen by an interleaving compact survives. Asserts
1035 /// edge conservation under many concurrent compacts racing inserts.
1036 #[test]
1037 fn test_concurrent_compaction_conserves_edges() {
1038 let am = std::sync::Arc::new(AdjacencyManager::new(64 * 1024 * 1024));
1039 let n: u64 = 150;
1040
1041 let inserter = {
1042 let am = am.clone();
1043 std::thread::spawn(move || {
1044 for i in 1..=n {
1045 am.insert_edge(Vid::new(0), Vid::new(i), Eid::new(i), 1, i);
1046 if i % 8 == 0 {
1047 am.compact();
1048 }
1049 }
1050 })
1051 };
1052 let compactor = {
1053 let am = am.clone();
1054 std::thread::spawn(move || {
1055 for _ in 0..40 {
1056 am.compact();
1057 std::thread::yield_now();
1058 }
1059 })
1060 };
1061 inserter.join().unwrap();
1062 compactor.join().unwrap();
1063 am.compact();
1064
1065 let neighbors = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1066 let got: HashSet<u64> = neighbors.iter().map(|(v, _)| v.as_u64()).collect();
1067 for i in 1..=n {
1068 assert!(
1069 got.contains(&i),
1070 "edge to {i} was lost under concurrent compaction"
1071 );
1072 }
1073 assert_eq!(got.len(), n as usize, "no spurious or duplicate neighbors");
1074 }
1075
1076 #[test]
1077 fn test_compact_merges_into_main_csr() {
1078 let am = AdjacencyManager::new(1024 * 1024);
1079
1080 // Insert edges into overlay
1081 am.insert_edge(Vid::new(0), Vid::new(10), Eid::new(100), 1, 1);
1082 am.insert_edge(Vid::new(0), Vid::new(20), Eid::new(101), 1, 2);
1083
1084 // Compact: overlay → Main CSR
1085 am.compact();
1086
1087 // Frozen segments should be empty after compaction
1088 assert_eq!(am.frozen_segment_count(), 0);
1089
1090 // Edges should still be accessible via Main CSR
1091 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1092 assert_eq!(n.len(), 2);
1093
1094 assert!(am.has_csr(1, Direction::Outgoing));
1095 }
1096
1097 #[test]
1098 fn test_compact_removes_tombstoned_edges() {
1099 let am = AdjacencyManager::new(1024 * 1024);
1100
1101 // Set up Main CSR with one edge
1102 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1103 am.set_main_csr(1, Direction::Outgoing, csr);
1104
1105 // Add new edge + tombstone for old edge in overlay
1106 am.insert_edge(Vid::new(0), Vid::new(20), Eid::new(101), 1, 2);
1107 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 3);
1108
1109 am.compact();
1110
1111 // Only the new edge should remain
1112 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1113 assert_eq!(n.len(), 1);
1114 assert_eq!(n[0], (Vid::new(20), Eid::new(101)));
1115 }
1116
1117 #[test]
1118 fn test_should_compact() {
1119 let am = AdjacencyManager::new(1024 * 1024);
1120 assert!(!am.should_compact(4));
1121
1122 // Manually freeze the active overlay multiple times
1123 for _ in 0..4 {
1124 let frozen = {
1125 let mut active = am.active_overlay.write();
1126 let old = std::mem::take(&mut *active);
1127 Arc::new(old.freeze())
1128 };
1129 am.frozen_segments.write().push(frozen);
1130 }
1131
1132 assert!(am.should_compact(4));
1133 }
1134
1135 #[test]
1136 fn test_empty_manager() {
1137 let am = AdjacencyManager::new(1024 * 1024);
1138 assert!(
1139 am.get_neighbors(Vid::new(0), 1, Direction::Outgoing)
1140 .is_empty()
1141 );
1142 assert!(!am.has_csr(1, Direction::Outgoing));
1143 }
1144
1145 #[test]
1146 fn test_overlay_tombstone_removes_main_csr_edge() {
1147 // Simulates: insert edge → flush/compact into Main CSR → delete edge (tombstone in overlay)
1148 let am = AdjacencyManager::new(1024 * 1024);
1149
1150 // Edge already compacted into Main CSR
1151 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1152 am.set_main_csr(1, Direction::Outgoing, csr);
1153
1154 // Verify edge is visible before deletion
1155 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1156 assert_eq!(n.len(), 1);
1157
1158 // Delete via overlay tombstone (simulates Writer::delete_edge dual-write)
1159 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
1160
1161 // Tombstone in overlay must remove edge from Main CSR results
1162 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1163 assert!(
1164 n.is_empty(),
1165 "Edge should be removed by overlay tombstone, got {:?}",
1166 n
1167 );
1168 }
1169
1170 #[test]
1171 fn test_overlay_tombstone_removes_main_csr_edge_versioned() {
1172 // Same scenario but via get_neighbors_at_version
1173 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, 5);
1179
1180 // At version 3: edge created at v1, tombstone at v5 → visible
1181 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
1182 assert_eq!(n.len(), 1);
1183
1184 // At version 5: tombstone applies → not visible
1185 let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
1186 assert!(
1187 n.is_empty(),
1188 "Edge should be removed by overlay tombstone at version 5"
1189 );
1190 }
1191
1192 #[test]
1193 fn test_frozen_tombstone_removes_main_csr_edge() {
1194 // Edge in Main CSR, tombstone in a frozen segment
1195 let am = AdjacencyManager::new(1024 * 1024);
1196
1197 let csr = MainCsr::from_edge_entries(0, vec![(0, Vid::new(10), Eid::new(100), 1)]);
1198 am.set_main_csr(1, Direction::Outgoing, csr);
1199
1200 // Add tombstone to active overlay, then compact to freeze it
1201 am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
1202
1203 // Freeze the overlay manually
1204 {
1205 let mut active = am.active_overlay.write();
1206 let old = std::mem::take(&mut *active);
1207 let frozen = std::sync::Arc::new(old.freeze());
1208 am.frozen_segments.write().push(frozen);
1209 }
1210
1211 // The frozen segment's tombstone should remove the Main CSR edge
1212 let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1213 assert!(n.is_empty(), "Frozen tombstone should remove Main CSR edge");
1214 }
1215
1216 #[test]
1217 fn test_per_edge_version_filtering() {
1218 // Test that edges inserted at different versions are correctly filtered
1219 // by get_neighbors_at_version()
1220 let am = AdjacencyManager::new(1024 * 1024);
1221
1222 let src = Vid::new(0);
1223 let dst_a = Vid::new(10);
1224 let dst_b = Vid::new(20);
1225 let eid_a = Eid::new(100);
1226 let eid_b = Eid::new(200);
1227 let etype = 1;
1228
1229 // Insert edge A at version 3
1230 am.insert_edge(src, dst_a, eid_a, etype, 3);
1231
1232 // Insert edge B at version 7
1233 am.insert_edge(src, dst_b, eid_b, etype, 7);
1234
1235 // Query at version 2 → neither edge visible
1236 let neighbors_v2 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 2);
1237 assert!(
1238 neighbors_v2.is_empty(),
1239 "No edges should be visible at version 2"
1240 );
1241
1242 // Query at version 5 → only edge A visible
1243 let neighbors_v5 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 5);
1244 assert_eq!(
1245 neighbors_v5.len(),
1246 1,
1247 "Only edge A should be visible at version 5"
1248 );
1249 assert_eq!(neighbors_v5[0].0, dst_a, "Edge A destination should match");
1250 assert_eq!(neighbors_v5[0].1, eid_a, "Edge A ID should match");
1251
1252 // Query at version 7 → both edges visible
1253 let neighbors_v7 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 7);
1254 assert_eq!(
1255 neighbors_v7.len(),
1256 2,
1257 "Both edges should be visible at version 7"
1258 );
1259
1260 // Query at version 10 → both edges visible
1261 let neighbors_v10 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 10);
1262 assert_eq!(
1263 neighbors_v10.len(),
1264 2,
1265 "Both edges should be visible at version 10"
1266 );
1267 }
1268
1269 #[test]
1270 fn test_duplicate_edges_deduplicated_by_eid() {
1271 // Test Issue #41: Same Eid in MainCsr (v1) and overlay (v3) → only 1 result from get_neighbors
1272 let am = AdjacencyManager::new(1024 * 1024);
1273
1274 let src = Vid::new(0);
1275 let dst = Vid::new(10);
1276 let eid = Eid::new(100);
1277 let etype = 1;
1278
1279 // Set up Main CSR with edge at version 1
1280 let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1281 am.set_main_csr(etype, Direction::Outgoing, csr);
1282
1283 // Insert same Eid into overlay at version 3 (update scenario)
1284 am.insert_edge(src, dst, eid, etype, 3);
1285
1286 // get_neighbors should return only 1 edge (HashMap<Eid, Vid> deduplicates)
1287 let neighbors = am.get_neighbors(src, etype, Direction::Outgoing);
1288 assert_eq!(
1289 neighbors.len(),
1290 1,
1291 "Duplicate Eid should result in single entry"
1292 );
1293 assert_eq!(neighbors[0], (dst, eid));
1294 }
1295
1296 #[test]
1297 fn test_compact_deduplicates_edges_keeps_highest_version() {
1298 // Test Issue #41: Same Eid at v1 in CSR and v5 in overlay
1299 // After compact: get_neighbors_at_version(v5) → visible
1300 // get_neighbors_at_version(v1) → NOT visible (compaction kept v5)
1301 let am = AdjacencyManager::new(1024 * 1024);
1302
1303 let src = Vid::new(0);
1304 let dst = Vid::new(10);
1305 let eid = Eid::new(100);
1306 let etype = 1;
1307
1308 // Set up Main CSR with edge at version 1
1309 let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1310 am.set_main_csr(etype, Direction::Outgoing, csr);
1311
1312 // Insert same Eid into overlay at version 5 (newer version)
1313 am.insert_edge(src, dst, eid, etype, 5);
1314
1315 // Before compact: both versions exist in different layers
1316 // After compact: only highest version (v5) should remain
1317
1318 am.compact();
1319
1320 // At version 5: edge should be visible (highest version kept)
1321 let neighbors_v5 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 5);
1322 assert_eq!(neighbors_v5.len(), 1, "Edge should be visible at version 5");
1323 assert_eq!(neighbors_v5[0], (dst, eid));
1324
1325 // At version 4: edge should still be visible (v5 edge has created_version=5)
1326 // Actually, the edge at v5 replaces v1, so the edge has version 5
1327 // So at version 4, we should NOT see it
1328 let neighbors_v4 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 4);
1329 assert_eq!(
1330 neighbors_v4.len(),
1331 0,
1332 "After compaction, only version 5 exists; version 4 should not see it"
1333 );
1334
1335 // At version 1: edge should NOT be visible (old version discarded)
1336 let neighbors_v1 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 1);
1337 assert_eq!(
1338 neighbors_v1.len(),
1339 0,
1340 "Old version discarded during compaction deduplication"
1341 );
1342
1343 // At version 6: edge should be visible (v5 edge still exists)
1344 let neighbors_v6 = am.get_neighbors_at_version(src, etype, Direction::Outgoing, 6);
1345 assert_eq!(neighbors_v6.len(), 1, "Edge should be visible at version 6");
1346 }
1347
1348 /// Test that tombstone filtering is O(result_size), not O(tombstone_count).
1349 /// This verifies fix for issue #140 (inverted tombstone scan).
1350 #[test]
1351 fn test_tombstone_scan_performance() {
1352 let am = AdjacencyManager::new(1024 * 1024);
1353 let vertex_a = Vid::new(1);
1354 let vertex_b = Vid::new(2);
1355 let etype = 1;
1356
1357 // Create 5 edges from vertex_a
1358 let mut a_edges = Vec::new();
1359 for i in 0..5 {
1360 let dst = Vid::new(100 + i);
1361 let eid = Eid::new(1000 + i);
1362 am.insert_edge(vertex_a, dst, eid, etype, 1);
1363 a_edges.push((dst, eid));
1364 }
1365
1366 // Create 100 deleted edges from vertex_b (creates 100 tombstones)
1367 for i in 0..100 {
1368 let dst = Vid::new(200 + i);
1369 let eid = Eid::new(2000 + i);
1370 am.insert_edge(vertex_b, dst, eid, etype, 1);
1371 am.add_tombstone(eid, vertex_b, dst, etype, 2);
1372 }
1373
1374 // Query neighbors of vertex_a
1375 // With O(T) scan, this would iterate 100 tombstones
1376 // With O(result) scan, this only checks 5 edges against tombstone map
1377 let neighbors = am.get_neighbors(vertex_a, etype, Direction::Outgoing);
1378
1379 // Verify all 5 edges are returned correctly
1380 assert_eq!(
1381 neighbors.len(),
1382 5,
1383 "Should return all 5 edges from vertex_a"
1384 );
1385 for (dst, eid) in &a_edges {
1386 assert!(
1387 neighbors.contains(&(*dst, *eid)),
1388 "Edge {:?} should be in results",
1389 (dst, eid)
1390 );
1391 }
1392
1393 // Verify vertex_b has no neighbors (all tombstoned)
1394 let b_neighbors = am.get_neighbors(vertex_b, etype, Direction::Outgoing);
1395 assert_eq!(
1396 b_neighbors.len(),
1397 0,
1398 "Vertex B should have no neighbors (all deleted)"
1399 );
1400 }
1401
1402 /// Verify that the irrelevant-segment short-circuit (issue #55) doesn't
1403 /// change observable behavior: with many frozen segments where only one
1404 /// holds the queried edge, `get_neighbors` returns exactly that edge.
1405 ///
1406 /// Also covers `get_neighbors_at_version` with the same short-circuit.
1407 #[test]
1408 fn test_get_neighbors_skips_irrelevant_segments() {
1409 let am = AdjacencyManager::new(1024 * 1024);
1410 let participant = Vid::new(1);
1411 let session = Vid::new(2);
1412 let link_eid = Eid::new(100);
1413 let link_etype: u32 = 1;
1414 let unrelated_etype: u32 = 2;
1415
1416 // Build up 50 frozen segments. Only segment #17 carries the LINK
1417 // edge from `participant`. The others are populated with unrelated
1418 // edges that share neither edge_type nor vid with the query.
1419 for i in 0..50 {
1420 if i == 17 {
1421 am.insert_edge(participant, session, link_eid, link_etype, i as u64 + 1);
1422 } else {
1423 // Unrelated traffic: different edge_type, different vids.
1424 let src = Vid::new(1000 + i as u64);
1425 let dst = Vid::new(2000 + i as u64);
1426 let eid = Eid::new(10_000 + i as u64);
1427 am.insert_edge(src, dst, eid, unrelated_etype, i as u64 + 1);
1428 }
1429 // Freeze the active overlay into a new frozen segment.
1430 let frozen = {
1431 let mut active = am.active_overlay.write();
1432 let old = std::mem::take(&mut *active);
1433 Arc::new(old.freeze())
1434 };
1435 am.frozen_segments.write().push(frozen);
1436 }
1437
1438 // Sanity: 50 frozen segments accumulated, none compacted yet.
1439 assert_eq!(am.frozen_segment_count(), 50);
1440
1441 // Hot path: returns exactly the one LINK edge.
1442 let n = am.get_neighbors(participant, link_etype, Direction::Outgoing);
1443 assert_eq!(n.len(), 1);
1444 assert_eq!(n[0], (session, link_eid));
1445
1446 // Snapshot path: same answer at a version that includes segment #17.
1447 let n_at = am.get_neighbors_at_version(participant, link_etype, Direction::Outgoing, 100);
1448 assert_eq!(n_at.len(), 1);
1449 assert_eq!(n_at[0], (session, link_eid));
1450
1451 // Snapshot path: at a version BEFORE segment #17 was created (#17's
1452 // version is 18), the edge must not be visible.
1453 let n_before =
1454 am.get_neighbors_at_version(participant, link_etype, Direction::Outgoing, 17);
1455 assert!(n_before.is_empty());
1456
1457 // The unrelated `unrelated_etype` edges must still be reachable —
1458 // short-circuiting must not have hidden them from their own queries.
1459 // i=18 was an unrelated insert (i=17 was the LINK), so Vid(1018) is
1460 // a valid source for an unrelated edge.
1461 let unrelated = am.get_neighbors(Vid::new(1018), unrelated_etype, Direction::Outgoing);
1462 assert_eq!(unrelated.len(), 1);
1463 }
1464
1465 /// Verify that the tombstone short-circuit (issue #55) doesn't drop
1466 /// transitively-shadowed edges: a frozen segment that has a tombstone
1467 /// for an edge present in Main CSR must still apply that tombstone,
1468 /// even if the segment has no inserts of its own for the queried type.
1469 #[test]
1470 fn test_tombstone_in_unrelated_segment_still_applied() {
1471 let am = AdjacencyManager::new(1024 * 1024);
1472 let src = Vid::new(0);
1473 let dst = Vid::new(10);
1474 let eid = Eid::new(100);
1475 let etype: u32 = 1;
1476
1477 // Edge exists in Main CSR.
1478 let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1479 am.set_main_csr(etype, Direction::Outgoing, csr);
1480
1481 // Add a tombstone in the active overlay deleting the Main CSR edge.
1482 am.add_tombstone(eid, src, dst, etype, 2);
1483
1484 // Freeze the active overlay so the tombstone now lives in a frozen
1485 // segment whose `inserts` is empty for `etype`. The short-circuit
1486 // for "no inserts AND no tombstones" must NOT skip this segment —
1487 // it has a tombstone we still need to honour.
1488 let frozen = {
1489 let mut active = am.active_overlay.write();
1490 let old = std::mem::take(&mut *active);
1491 Arc::new(old.freeze())
1492 };
1493 am.frozen_segments.write().push(frozen);
1494
1495 let n = am.get_neighbors(src, etype, Direction::Outgoing);
1496 assert!(
1497 n.is_empty(),
1498 "tombstone in frozen segment must still hide Main CSR edge"
1499 );
1500 }
1501}