Skip to main content

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/// Deduplicate `(offset, neighbor, eid, version)` entries by `Eid`, keeping the
112/// entry with the highest version for each. Multiple versions of the same edge
113/// can coexist in L2+L1, across L1 runs, and across frozen overlay segments.
114fn 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
139/// Unified adjacency manager for the dual-CSR architecture.
140///
141/// Orchestrates Main CSR (packed alive edges), L0-csr overlay
142/// (in-memory mutations), and Shadow CSR (deleted edges for time-travel).
143/// Data flush never invalidates or rebuilds the CSR.
144pub struct AdjacencyManager {
145    /// Main CSR per `(edge_type, direction)` — all alive edges.
146    /// Edge type is u32 with bit 31 = 0 for schema'd, 1 for schemaless.
147    main_csr: DashMap<(u32, Direction), Arc<MainCsr>>,
148
149    /// Active L0-csr segment (current writes go here).
150    active_overlay: Arc<RwLock<L0CsrSegment>>,
151
152    /// Frozen segments awaiting compaction (oldest first).
153    frozen_segments: RwLock<Vec<Arc<FrozenCsrSegment>>>,
154
155    /// Shadow CSR for time-travel deleted edge tracking.
156    shadow: ShadowCsr,
157    /// Versions pinned by in-flight `pinned_at_version` readers; the floor for
158    /// shadow GC. See [`PinnedVersions`].
159    pinned_versions: Arc<PinnedVersions>,
160
161    /// Current approximate memory usage in bytes.
162    current_bytes: AtomicUsize,
163
164    /// Maximum memory budget in bytes.
165    max_bytes: usize,
166
167    /// Coalescing locks for warm() operations — prevents cache stampede.
168    /// Key: (edge_type_id, Direction), Value: Mutex guard for that warm operation.
169    warm_guards: DashMap<(u32, Direction), Arc<tokio::sync::Mutex<()>>>,
170
171    /// Serializes `compact()` so two compactions can't interleave their
172    /// freeze→snapshot→clear sequences and lose a frozen segment. (review H12)
173    compact_lock: parking_lot::Mutex<()>,
174}
175
176impl AdjacencyManager {
177    /// Creates a new adjacency manager with the given memory budget.
178    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    /// Returns neighbors for the current state (hot path, no version filtering).
193    ///
194    /// Reads Main CSR + frozen segments + active overlay, minus tombstones.
195    /// Tombstones from any layer remove edges from all lower layers.
196    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            // 1. Main CSR
201            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            // 2. Frozen segments (oldest first) — add inserts, then remove tombstones.
208            // Skip whole segment when it has no inserts for this (edge_type, dir)
209            // AND no tombstones at all — both is_empty checks are O(1) on plain HashMaps,
210            // so this is a strict speedup that scales away with frozen-segment count
211            // when most segments don't touch the queried edge type. See issue #55.
212            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                // Apply tombstones against ALL prior results (Main CSR + older segments).
227                // Skip the retain pass entirely when there are no tombstones — it would
228                // be a no-op but still O(result_size) due to the closure call.
229                if has_tombstones {
230                    result.retain(|eid, _| !segment.tombstones.contains_key(eid));
231                }
232            }
233
234            // 3. Active overlay — add inserts, then remove tombstones.
235            // Same short-circuits as the frozen branch, plus we hold the read lock
236            // for the whole branch so DashMap atomic-op cost is paid once, not twice.
237            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    /// Returns neighbors visible at a specific snapshot version.
257    ///
258    /// Filters Main CSR entries by `created_version`, applies frozen/active
259    /// overlay with version filtering, and resurrects Shadow CSR entries
260    /// that were alive at the given version.
261    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            // 1. Main CSR — filter by created_version
272            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            // 2. Frozen segments — filter inserts by version, apply tombstones.
281            // Same skip-irrelevant-segment short-circuit as get_neighbors. See issue #55.
282            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            // 3. Active overlay — add version-filtered inserts, then apply tombstones
309            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            // 4. Shadow CSR — resurrect edges alive at version
336            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    /// Records an edge insertion into the L0-csr overlay (both directions).
348    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    /// Records a tombstone for a deleted edge in the L0-csr overlay.
355    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    /// Sets the Main CSR for a specific edge type and direction.
361    ///
362    /// Used by `warm()` to install a freshly built CSR from storage.
363    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    /// Checks whether a Main CSR exists for the given edge type and direction.
370    pub fn has_csr(&self, edge_type: u32, direction: Direction) -> bool {
371        self.main_csr.contains_key(&(edge_type, direction))
372    }
373
374    /// Checks whether this manager has been activated for the given edge type.
375    ///
376    /// Returns `true` if a Main CSR exists or the overlay has entries for
377    /// this edge type and direction.
378    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    /// Returns the distinct edge type ids known to this manager.
386    ///
387    /// Spans the Main CSR plus the active and frozen overlay segments, so it
388    /// covers both warmed (L1/L2-loaded) edge types and live overlay-resident
389    /// types (e.g. recently committed edges that a flush moved out of L0 but kept
390    /// in the dual-write overlay). Used when an endpoint resolver knows an edge
391    /// id but not its type: this is the small set of types this query has touched,
392    /// so it bounds an eid-orientation probe to a short candidate list rather than
393    /// the whole schema.
394    ///
395    /// # Examples
396    ///
397    /// ```ignore
398    /// for etype in adjacency_manager.known_edge_type_ids() {
399    ///     // probe etype for the edge of interest
400    /// }
401    /// ```
402    #[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    /// Returns the number of frozen segments awaiting compaction.
417    pub fn frozen_segment_count(&self) -> usize {
418        self.frozen_segments.read().len()
419    }
420
421    /// Returns whether compaction should be triggered based on segment count.
422    pub fn should_compact(&self, threshold: usize) -> bool {
423        self.frozen_segment_count() >= threshold
424    }
425
426    /// Compacts frozen overlay segments into the Main CSR.
427    ///
428    /// Freezes the active overlay, merges all frozen segments with the
429    /// existing Main CSR, moves tombstoned edges to Shadow CSR, and
430    /// atomically swaps in the new Main CSR.
431    ///
432    /// CRITICAL: Frozen segments remain readable until the new CSR is installed,
433    /// eliminating the visibility gap where edges would be invisible.
434    pub fn compact(&self) {
435        // Serialize compaction: two concurrent compacts would each freeze, take
436        // their own snapshot, then clear — the second clear wiping segments the
437        // first had not yet merged. (review H12)
438        let _compact_guard = self.compact_lock.lock();
439
440        // Step 1: Freeze active overlay and push to frozen list
441        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        // Step 2: CLONE frozen segments for building (DON'T drain yet)
449        // This ensures they remain readable during CSR construction
450        let segments = self.frozen_segments.read().clone();
451
452        // Step 3: Collect all (edge_type, direction) keys from segments + existing CSRs
453        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        // Step 4: For each key, merge
464        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            // Collect all tombstone EIDs
469            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                        // Move to shadow CSR. ShadowCsr is keyed by the queried
476                        // vid for the direction, exactly like the CSRs themselves
477                        // (Incoming is keyed by dst with neighbor src — see the
478                        // insert at `insert_edge(dst, src, .., Incoming)` and the
479                        // get_neighbors_at_version swap). Key by src for Outgoing,
480                        // by dst for Incoming — otherwise a time-travel read after
481                        // compaction looks up the wrong vid and the tombstone is
482                        // invisible (deleted edge resurrected).
483                        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, // unknown; overlay tombstones don't track creation version
495                                deleted_version: ts.version,
496                            },
497                            direction,
498                        );
499                    }
500                }
501            }
502
503            // Add entries from old Main CSR
504            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            // Overlay frozen segments (oldest first)
522            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            // Build new Main CSR and install
539            let new_csr = MainCsr::from_edge_entries(max_offset as usize, entries);
540            let size = new_csr.memory_usage();
541
542            // Remove old size, add new
543            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        // Step 5: drain EXACTLY the segments we snapshotted in Step 2 — not a
554        // blanket clear(). A concurrent `freeze()` may have pushed a new frozen
555        // segment after the snapshot; that segment was NOT merged into the new
556        // CSR, so clearing it would silently lose its topology. Retain anything
557        // not in the snapshot (compared by Arc identity). (review H12)
558        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    /// Warms the Main CSR from storage (L2 adjacency + L1 delta) for a specific edge type and direction.
566    ///
567    /// Reads L2 adjacency datasets and L1 delta entries from Lance,
568    /// builds a [`MainCsr`] with version metadata, and populates the
569    /// [`ShadowCsr`] with L1 tombstones. Called once at startup or
570    /// lazily on first access per edge type.
571    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        // Use unified lookup to support both schema'd and schemaless edge types
581        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        // Determine which labels to load adjacency for based on edge type metadata
586        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                // 1. Read L2 (Adjacency Dataset)
611                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                        // Fail closed: a transient scan error must abort the warm,
627                        // not `unwrap_or_default()` into an empty L2 read that then
628                        // gets cached as the adjacency CSR — that silently drops
629                        // every base edge for this type until restart (review #3b).
630                        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                                    // L2 adjacency rows don't carry per-edge _version.
667                                    // Version 0 means "from base storage" — the `_version <= hwm` filter on
668                                    // the query already ensures we only load rows within the snapshot window.
669                                    // At query time, get_neighbors_at_version() uses created_version to filter,
670                                    // so version=0 edges are always visible (which is correct for compacted L2 data).
671                                    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            // 2. Read L1 (Delta)
685            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                // Fail closed: propagate a delta scan error rather than silently
701                // skipping it, which would drop unflushed edges from the cached
702                // adjacency CSR until restart (review #3b).
703                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                        // Optionally read _version column
732                        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); // 0=Insert, 1=Delete
741                            let row_version = version_col.as_ref().map_or(0, |vc| vc.value(i));
742
743                            // For incoming edges, the CSR key is dst (the vertex
744                            // receiving the edge) and the neighbor is src.
745                            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        // Filter out deleted edges
775        if !deleted_eids.is_empty() {
776            entries.retain(|(_, _, eid, _)| !deleted_eids.contains(eid));
777        }
778
779        dedup_entries_by_eid(&mut entries);
780
781        // Build MainCsr
782        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    /// Coalesced warm() operation to prevent cache stampede (Issue #13).
790    ///
791    /// Uses double-checked locking: fast-path checks if CSR already loaded,
792    /// then acquires per-(edge_type, direction) lock to ensure only one concurrent
793    /// warm() per adjacency key. Other readers wait for the first warm() to complete.
794    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        // Fast path: already loaded
802        if self.has_csr(edge_type_id, direction) {
803            return Ok(());
804        }
805
806        // Coalesce: only one concurrent warm per (type, dir)
807        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        // Double-check after acquiring lock
816        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    /// Returns the current approximate memory usage in bytes.
824    pub fn memory_usage(&self) -> usize {
825        // `current_bytes` accounts only for the main CSR — every mutation of it
826        // is on a `main_csr` path. Shadow retention was therefore invisible to
827        // the budget and could never trip `max_bytes`, which is how an
828        // unbounded leak there went unnoticed. Counted approximately rather
829        // than tracked incrementally: the shadow is small when healthy, and an
830        // exact counter would need hooks on every retain in `gc`.
831        self.current_bytes.load(Ordering::Relaxed) + self.shadow.approx_bytes()
832    }
833
834    /// The pinned-version registry backing shadow GC.
835    pub fn pinned_versions(&self) -> &Arc<PinnedVersions> {
836        &self.pinned_versions
837    }
838
839    /// Reclaim shadow entries no in-flight reader can reach.
840    ///
841    /// The floor is the minimum pinned version, or `current_version` when
842    /// nothing is pinned — a reader starting now pins at the current version,
843    /// so entries deleted at or below it are unreachable. `current_version` is
844    /// passed in because the manager does not track it; the writer does.
845    ///
846    /// Called after compaction. Safe to call at any time: it only ever removes
847    /// entries whose `deleted_version` is at or below the floor, which is
848    /// exactly the set `get_entries_at_version` can no longer return.
849    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    /// Shadow-CSR entries currently retained.
858    ///
859    /// Exposed for retention tests and diagnostics; see
860    /// [`ShadowCsr::add_deleted_edge`] for why this can grow.
861    pub fn shadow_entry_count(&self) -> usize {
862        self.shadow.entry_count()
863    }
864
865    /// Returns the maximum memory budget in bytes.
866    pub fn max_bytes(&self) -> usize {
867        self.max_bytes
868    }
869
870    /// Provides access to the shadow CSR for time-travel queries.
871    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        // Incoming direction
905        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        // Main CSR has one edge
933        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        // Overlay adds another
937        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        // Main CSR with two edges at different versions
963        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        // At version 3: only first edge visible
973        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        // At version 5: both visible
978        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        // Add a deleted edge to shadow: created at v1, deleted at v5
987        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        // At version 3: shadow edge should be visible
1000        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        // At version 5: deleted, not visible
1005        let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 5);
1006        assert!(n.is_empty());
1007    }
1008
1009    /// H12: concurrent compaction must never drop an edge. `compact()` is the
1010    /// only writer of `frozen_segments`, now serialized under `compact_lock`,
1011    /// and Step 5 drains exactly the snapshotted segments rather than clearing
1012    /// all — so a segment frozen by an interleaving compact survives. Asserts
1013    /// edge conservation under many concurrent compacts racing inserts.
1014    #[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        // Insert edges into overlay
1059        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        // Compact: overlay → Main CSR
1063        am.compact();
1064
1065        // Frozen segments should be empty after compaction
1066        assert_eq!(am.frozen_segment_count(), 0);
1067
1068        // Edges should still be accessible via Main CSR
1069        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        // Set up Main CSR with one edge
1080        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        // Add new edge + tombstone for old edge in overlay
1084        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        // Only the new edge should remain
1090        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        // Manually freeze the active overlay multiple times
1101        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        // Simulates: insert edge → flush/compact into Main CSR → delete edge (tombstone in overlay)
1126        let am = AdjacencyManager::new(1024 * 1024);
1127
1128        // Edge already compacted into Main CSR
1129        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        // Verify edge is visible before deletion
1133        let n = am.get_neighbors(Vid::new(0), 1, Direction::Outgoing);
1134        assert_eq!(n.len(), 1);
1135
1136        // Delete via overlay tombstone (simulates Writer::delete_edge dual-write)
1137        am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
1138
1139        // Tombstone in overlay must remove edge from Main CSR results
1140        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        // Same scenario but via get_neighbors_at_version
1151        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        // At version 3: edge created at v1, tombstone at v5 → visible
1159        let n = am.get_neighbors_at_version(Vid::new(0), 1, Direction::Outgoing, 3);
1160        assert_eq!(n.len(), 1);
1161
1162        // At version 5: tombstone applies → not visible
1163        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        // Edge in Main CSR, tombstone in a frozen segment
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        // Add tombstone to active overlay, then compact to freeze it
1179        am.add_tombstone(Eid::new(100), Vid::new(0), Vid::new(10), 1, 2);
1180
1181        // Freeze the overlay manually
1182        {
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        // The frozen segment's tombstone should remove the Main CSR edge
1190        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        // Test that edges inserted at different versions are correctly filtered
1197        // by get_neighbors_at_version()
1198        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        // Insert edge A at version 3
1208        am.insert_edge(src, dst_a, eid_a, etype, 3);
1209
1210        // Insert edge B at version 7
1211        am.insert_edge(src, dst_b, eid_b, etype, 7);
1212
1213        // Query at version 2 → neither edge visible
1214        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        // Query at version 5 → only edge A visible
1221        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        // Query at version 7 → both edges visible
1231        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        // Query at version 10 → both edges visible
1239        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        // Test Issue #41: Same Eid in MainCsr (v1) and overlay (v3) → only 1 result from get_neighbors
1250        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        // Set up Main CSR with edge at version 1
1258        let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1259        am.set_main_csr(etype, Direction::Outgoing, csr);
1260
1261        // Insert same Eid into overlay at version 3 (update scenario)
1262        am.insert_edge(src, dst, eid, etype, 3);
1263
1264        // get_neighbors should return only 1 edge (HashMap<Eid, Vid> deduplicates)
1265        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        // Test Issue #41: Same Eid at v1 in CSR and v5 in overlay
1277        // After compact: get_neighbors_at_version(v5) → visible
1278        //               get_neighbors_at_version(v1) → NOT visible (compaction kept v5)
1279        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        // Set up Main CSR with edge at version 1
1287        let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1288        am.set_main_csr(etype, Direction::Outgoing, csr);
1289
1290        // Insert same Eid into overlay at version 5 (newer version)
1291        am.insert_edge(src, dst, eid, etype, 5);
1292
1293        // Before compact: both versions exist in different layers
1294        // After compact: only highest version (v5) should remain
1295
1296        am.compact();
1297
1298        // At version 5: edge should be visible (highest version kept)
1299        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        // At version 4: edge should still be visible (v5 edge has created_version=5)
1304        // Actually, the edge at v5 replaces v1, so the edge has version 5
1305        // So at version 4, we should NOT see it
1306        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        // At version 1: edge should NOT be visible (old version discarded)
1314        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        // At version 6: edge should be visible (v5 edge still exists)
1322        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 that tombstone filtering is O(result_size), not O(tombstone_count).
1327    /// This verifies fix for issue #140 (inverted tombstone scan).
1328    #[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        // Create 5 edges from vertex_a
1336        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        // Create 100 deleted edges from vertex_b (creates 100 tombstones)
1345        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        // Query neighbors of vertex_a
1353        // With O(T) scan, this would iterate 100 tombstones
1354        // With O(result) scan, this only checks 5 edges against tombstone map
1355        let neighbors = am.get_neighbors(vertex_a, etype, Direction::Outgoing);
1356
1357        // Verify all 5 edges are returned correctly
1358        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        // Verify vertex_b has no neighbors (all tombstoned)
1372        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    /// Verify that the irrelevant-segment short-circuit (issue #55) doesn't
1381    /// change observable behavior: with many frozen segments where only one
1382    /// holds the queried edge, `get_neighbors` returns exactly that edge.
1383    ///
1384    /// Also covers `get_neighbors_at_version` with the same short-circuit.
1385    #[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        // Build up 50 frozen segments. Only segment #17 carries the LINK
1395        // edge from `participant`. The others are populated with unrelated
1396        // edges that share neither edge_type nor vid with the query.
1397        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                // Unrelated traffic: different edge_type, different vids.
1402                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            // Freeze the active overlay into a new frozen segment.
1408            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        // Sanity: 50 frozen segments accumulated, none compacted yet.
1417        assert_eq!(am.frozen_segment_count(), 50);
1418
1419        // Hot path: returns exactly the one LINK edge.
1420        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        // Snapshot path: same answer at a version that includes segment #17.
1425        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        // Snapshot path: at a version BEFORE segment #17 was created (#17's
1430        // version is 18), the edge must not be visible.
1431        let n_before =
1432            am.get_neighbors_at_version(participant, link_etype, Direction::Outgoing, 17);
1433        assert!(n_before.is_empty());
1434
1435        // The unrelated `unrelated_etype` edges must still be reachable —
1436        // short-circuiting must not have hidden them from their own queries.
1437        // i=18 was an unrelated insert (i=17 was the LINK), so Vid(1018) is
1438        // a valid source for an unrelated edge.
1439        let unrelated = am.get_neighbors(Vid::new(1018), unrelated_etype, Direction::Outgoing);
1440        assert_eq!(unrelated.len(), 1);
1441    }
1442
1443    /// Verify that the tombstone short-circuit (issue #55) doesn't drop
1444    /// transitively-shadowed edges: a frozen segment that has a tombstone
1445    /// for an edge present in Main CSR must still apply that tombstone,
1446    /// even if the segment has no inserts of its own for the queried type.
1447    #[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        // Edge exists in Main CSR.
1456        let csr = MainCsr::from_edge_entries(0, vec![(0, dst, eid, 1)]);
1457        am.set_main_csr(etype, Direction::Outgoing, csr);
1458
1459        // Add a tombstone in the active overlay deleting the Main CSR edge.
1460        am.add_tombstone(eid, src, dst, etype, 2);
1461
1462        // Freeze the active overlay so the tombstone now lives in a frozen
1463        // segment whose `inserts` is empty for `etype`. The short-circuit
1464        // for "no inserts AND no tombstones" must NOT skip this segment —
1465        // it has a tombstone we still need to honour.
1466        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}