diff --git a/crates/ruvector-diskann/examples/bench_delete_recall.rs b/crates/ruvector-diskann/examples/bench_delete_recall.rs new file mode 100644 index 000000000..beb428bf2 --- /dev/null +++ b/crates/ruvector-diskann/examples/bench_delete_recall.rs @@ -0,0 +1,159 @@ +//! Measure recall and delete latency after deleting 20% of a DiskANN index. + +use rand::prelude::*; +use ruvector_diskann::{DiskAnnConfig, DiskAnnIndex}; +use std::collections::HashSet; +use std::fs; +use std::path::PathBuf; +use std::time::{Duration, Instant}; + +const DIM: usize = 32; +const K: usize = 10; +const QUERY_COUNT: usize = 50; + +fn l2_squared(a: &[f32], b: &[f32]) -> f32 { + a.iter() + .zip(b) + .map(|(left, right)| { + let delta = left - right; + delta * delta + }) + .sum() +} + +fn median_micros(samples: &mut [Duration]) -> f64 { + samples.sort_unstable(); + samples[samples.len() / 2].as_secs_f64() * 1_000_000.0 +} + +fn recall_at_10( + index: &DiskAnnIndex, + data: &[(String, Vec)], + queries: &[usize], + truth: &[HashSet], +) -> f64 { + let matches: usize = queries + .iter() + .zip(truth) + .map(|(&query_idx, expected)| { + index + .search(&data[query_idx].1, K) + .expect("search should succeed") + .iter() + .filter(|result| expected.contains(&result.id)) + .count() + }) + .sum(); + matches as f64 / (queries.len() * K) as f64 +} + +fn run_case(n: usize) { + let mut rng = StdRng::seed_from_u64(0x679D_15CA ^ n as u64); + let data: Vec<(String, Vec)> = (0..n) + .map(|idx| { + ( + format!("v{idx}"), + (0..DIM).map(|_| rng.gen::()).collect(), + ) + }) + .collect(); + let config = DiskAnnConfig { + dim: DIM, + max_degree: 32, + build_beam: 64, + search_beam: 64, + alpha: 1.2, + ..Default::default() + }; + + let mut order: Vec = (0..n).collect(); + order.shuffle(&mut rng); + let delete_count = n / 5; + let deleted: HashSet = order[..delete_count].iter().copied().collect(); + let survivor_indices: Vec = (0..n).filter(|idx| !deleted.contains(idx)).collect(); + let queries: Vec = order[delete_count..delete_count + QUERY_COUNT].to_vec(); + + let truth: Vec> = queries + .iter() + .map(|&query_idx| { + let mut exact: Vec<(usize, f32)> = survivor_indices + .iter() + .map(|&idx| (idx, l2_squared(&data[idx].1, &data[query_idx].1))) + .collect(); + exact.sort_unstable_by(|a, b| a.1.total_cmp(&b.1)); + exact + .iter() + .take(K) + .map(|(idx, _)| data[*idx].0.clone()) + .collect() + }) + .collect(); + + let dir: PathBuf = + std::env::temp_dir().join(format!("ruvector-delete-recall-{}-{n}", std::process::id())); + let _ = fs::remove_dir_all(&dir); + + let mut base = DiskAnnIndex::new(config.clone()); + base.insert_batch(data.clone()).expect("insert base data"); + base.build().expect("build base index"); + base.save(&dir).expect("save base index"); + + let mut deferred = DiskAnnIndex::load(&dir).expect("load deferred index"); + let mut repaired = DiskAnnIndex::load(&dir).expect("load repaired index"); + let mut deferred_times = Vec::with_capacity(delete_count); + let mut repair_times = Vec::with_capacity(delete_count); + + for &idx in &order[..delete_count] { + let started = Instant::now(); + deferred + .delete_deferred(&data[idx].0) + .expect("deferred delete"); + deferred_times.push(started.elapsed()); + + let started = Instant::now(); + repaired.delete(&data[idx].0).expect("repairing delete"); + repair_times.push(started.elapsed()); + } + + let survivor_data: Vec<(String, Vec)> = survivor_indices + .iter() + .map(|&idx| data[idx].clone()) + .collect(); + let mut fresh = DiskAnnIndex::new(config); + fresh + .insert_batch(survivor_data) + .expect("insert survivor data"); + fresh.build().expect("build fresh survivor index"); + + println!("RESULT n={n} mode=main recall_at_10=UNMEASURABLE note=deleted_ids_returned"); + println!( + "RESULT n={n} mode=tombstone_only recall_at_10={:.4}", + recall_at_10(&deferred, &data, &queries, &truth) + ); + println!( + "RESULT n={n} mode=tombstone_repair recall_at_10={:.4}", + recall_at_10(&repaired, &data, &queries, &truth) + ); + println!( + "RESULT n={n} mode=fresh_rebuild recall_at_10={:.4}", + recall_at_10(&fresh, &data, &queries, &truth) + ); + println!( + "LATENCY n={n} mode=tombstone_only median_us={:.3}", + median_micros(&mut deferred_times) + ); + println!( + "LATENCY n={n} mode=tombstone_repair median_us={:.3}", + median_micros(&mut repair_times) + ); + + drop((base, deferred, repaired)); + let _ = fs::remove_dir_all(dir); +} + +fn main() { + println!("BENCHMARK hardware=Apple_M2_Max threads=1 profile=release dim={DIM} k={K}"); + for n in [20_000, 100_000] { + run_case(n); + } +} diff --git a/crates/ruvector-diskann/src/distance.rs b/crates/ruvector-diskann/src/distance.rs index 5716380dc..804d9a305 100644 --- a/crates/ruvector-diskann/src/distance.rs +++ b/crates/ruvector-diskann/src/distance.rs @@ -41,15 +41,6 @@ impl FlatVectors { &self.data[start..start + self.dim] } - /// Zero out a vector (lazy deletion) - #[inline] - pub fn zero_out(&mut self, idx: usize) { - let start = idx * self.dim; - for v in &mut self.data[start..start + self.dim] { - *v = f32::NAN; - } - } - pub fn len(&self) -> usize { self.count } diff --git a/crates/ruvector-diskann/src/graph.rs b/crates/ruvector-diskann/src/graph.rs index c8d6e5bff..bff27ac2e 100644 --- a/crates/ruvector-diskann/src/graph.rs +++ b/crates/ruvector-diskann/src/graph.rs @@ -11,6 +11,12 @@ use rayon::prelude::*; use std::cmp::Ordering; use std::collections::BinaryHeap; +fn is_tombstoned(tombstones: &[u64], node: usize) -> bool { + tombstones + .get(node / 64) + .is_some_and(|word| word & (1u64 << (node % 64)) != 0) +} + #[derive(Clone)] struct Candidate { id: u32, @@ -19,7 +25,7 @@ struct Candidate { impl PartialEq for Candidate { fn eq(&self, other: &Self) -> bool { - self.distance == other.distance + self.distance.total_cmp(&other.distance) == Ordering::Equal } } impl Eq for Candidate {} @@ -30,10 +36,7 @@ impl PartialOrd for Candidate { } impl Ord for Candidate { fn cmp(&self, other: &Self) -> Ordering { - other - .distance - .partial_cmp(&self.distance) - .unwrap_or(Ordering::Equal) + other.distance.total_cmp(&self.distance) } } @@ -43,7 +46,7 @@ struct MaxCandidate { } impl PartialEq for MaxCandidate { fn eq(&self, other: &Self) -> bool { - self.distance == other.distance + self.distance.total_cmp(&other.distance) == Ordering::Equal } } impl Eq for MaxCandidate {} @@ -54,9 +57,7 @@ impl PartialOrd for MaxCandidate { } impl Ord for MaxCandidate { fn cmp(&self, other: &Self) -> Ordering { - self.distance - .partial_cmp(&other.distance) - .unwrap_or(Ordering::Equal) + self.distance.total_cmp(&other.distance) } } @@ -198,7 +199,7 @@ impl VamanaGraph { } let mut result: Vec<(u32, f32)> = best.into_iter().map(|c| (c.id, c.distance)).collect(); - result.sort_unstable_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(Ordering::Equal)); + result.sort_unstable_by(|a, b| a.1.total_cmp(&b.1)); let ids: Vec = result.into_iter().map(|(id, _)| id).collect(); (ids, visit_count) @@ -232,7 +233,7 @@ impl VamanaGraph { .filter(|&&c| c != node) .map(|&c| (c, l2_squared(vectors.get(c as usize), node_vec))) .collect(); - sorted.sort_unstable_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(Ordering::Equal)); + sorted.sort_unstable_by(|a, b| a.1.total_cmp(&b.1)); let mut result = Vec::with_capacity(self.max_degree); for (cand_id, cand_dist) in &sorted { @@ -253,6 +254,67 @@ impl VamanaGraph { result } + /// Remove a deleted node and locally repair each of its live in-neighbors. + /// + /// This index deliberately does not maintain reverse adjacency. Finding the + /// in-neighbors therefore performs a bounded scan of the graph and costs + /// O(n * degree) per deletion. Tombstoned sources and candidates are skipped. + pub(crate) fn repair_deleted( + &mut self, + vectors: &FlatVectors, + deleted: u32, + tombstones: &[u64], + ) { + let deleted_idx = deleted as usize; + if deleted_idx >= self.neighbors.len() || !is_tombstoned(tombstones, deleted_idx) { + return; + } + + let deleted_out: Vec = self.neighbors[deleted_idx] + .iter() + .copied() + .filter(|&candidate| { + candidate != deleted && !is_tombstoned(tombstones, candidate as usize) + }) + .collect(); + + if self.medoid == deleted { + if let Some(replacement) = deleted_out.first().copied().or_else(|| { + (0..self.neighbors.len()) + .find(|&node| !is_tombstoned(tombstones, node)) + .map(|idx| idx as u32) + }) { + self.medoid = replacement; + } + } + + for source in 0..self.neighbors.len() { + if source == deleted_idx || !self.neighbors[source].contains(&deleted) { + continue; + } + + if is_tombstoned(tombstones, source) { + self.neighbors[source].retain(|&neighbor| neighbor != deleted); + continue; + } + + let mut candidates: Vec = self.neighbors[source] + .iter() + .copied() + .chain(deleted_out.iter().copied()) + .filter(|&candidate| { + candidate != deleted && !is_tombstoned(tombstones, candidate as usize) + }) + .collect(); + candidates.sort_unstable(); + candidates.dedup(); + self.neighbors[source] = + self.robust_prune(vectors, source as u32, &candidates, self.alpha); + } + + self.neighbors[deleted_idx].clear(); + } + /// Parallel medoid finding using rayon fn find_medoid_parallel(&self, vectors: &FlatVectors) -> u32 { let n = vectors.len(); @@ -274,7 +336,7 @@ impl VamanaGraph { (0..n as u32) .into_par_iter() .map(|i| (i, l2_squared(vectors.get(i as usize), ¢roid))) - .min_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(Ordering::Equal)) + .min_by(|a, b| a.1.total_cmp(&b.1)) .map(|(id, _)| id) .unwrap_or(0) } diff --git a/crates/ruvector-diskann/src/index.rs b/crates/ruvector-diskann/src/index.rs index 587bcebec..578f7b098 100644 --- a/crates/ruvector-diskann/src/index.rs +++ b/crates/ruvector-diskann/src/index.rs @@ -62,6 +62,10 @@ pub struct DiskAnnIndex { id_map: Vec, /// Reverse mapping: external ID -> internal index id_reverse: HashMap, + /// Packed deletion markers. Tombstoned vectors remain available for routing. + tombstones: Vec, + /// Number of tombstoned vectors. + deleted_count: usize, /// Vamana graph graph: Option, /// Product quantizer (optional) @@ -85,6 +89,8 @@ impl DiskAnnIndex { vectors: FlatVectors::new(dim), id_map: Vec::new(), id_reverse: HashMap::new(), + tombstones: Vec::new(), + deleted_count: 0, graph: None, pq: None, pq_codes: Vec::new(), @@ -109,6 +115,9 @@ impl DiskAnnIndex { let idx = self.vectors.len() as u32; self.id_reverse.insert(id.clone(), idx); self.id_map.push(id); + if idx as usize / 64 == self.tombstones.len() { + self.tombstones.push(0); + } self.vectors.push(&vector); self.built = false; Ok(()) @@ -152,6 +161,11 @@ impl DiskAnnIndex { self.config.alpha, ); graph.build(&self.vectors)?; + for deleted in 0..n { + if self.is_tombstoned(deleted) { + graph.repair_deleted(&self.vectors, deleted as u32, &self.tombstones); + } + } self.graph = Some(graph); // Pre-allocate visited set for search @@ -185,9 +199,10 @@ impl DiskAnnIndex { // Re-rank candidates with exact distance let mut scored: Vec<(u32, f32)> = candidates .into_iter() + .filter(|&id| !self.is_tombstoned(id as usize)) .map(|id| (id, l2_squared(self.vectors.get(id as usize), query))) .collect(); - scored.sort_unstable_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal)); + scored.sort_unstable_by(|a, b| a.1.total_cmp(&b.1)); Ok(scored .into_iter() @@ -201,18 +216,45 @@ impl DiskAnnIndex { /// Get the number of vectors in the index pub fn count(&self) -> usize { - self.vectors.len() + self.vectors.len() - self.deleted_count } - /// Delete a vector by ID (marks as deleted, doesn't rebuild graph) + /// Delete a vector by ID and locally repair its live in-neighbors. + /// + /// In-neighbor discovery scans the bounded-degree adjacency, making this + /// O(n * degree) per delete. Unknown and already-deleted IDs return an error. pub fn delete(&mut self, id: &str) -> Result { - if let Some(&idx) = self.id_reverse.get(id) { - self.vectors.zero_out(idx as usize); - self.id_reverse.remove(id); - Ok(true) - } else { - Ok(false) + let idx = self.mark_deleted(id)?; + if let Some(graph) = self.graph.as_mut() { + graph.repair_deleted(&self.vectors, idx, &self.tombstones); } + Ok(true) + } + + /// Tombstone a vector without repairing graph edges. + /// + /// Searches may still traverse the preserved vector as a routing waypoint, + /// but the ID is filtered from results. Unknown and repeated deletes error. + pub fn delete_deferred(&mut self, id: &str) -> Result { + self.mark_deleted(id)?; + Ok(true) + } + + fn mark_deleted(&mut self, id: &str) -> Result { + let idx = self + .id_reverse + .remove(id) + .ok_or_else(|| DiskAnnError::NotFound(id.to_string()))?; + let idx = idx as usize; + self.tombstones[idx / 64] |= 1u64 << (idx % 64); + self.deleted_count += 1; + Ok(idx as u32) + } + + fn is_tombstoned(&self, idx: usize) -> bool { + self.tombstones + .get(idx / 64) + .is_some_and(|word| word & (1u64 << (idx % 64)) != 0) } /// Save index to disk @@ -257,6 +299,15 @@ impl DiskAnnIndex { .map_err(|e| DiskAnnError::Serialization(e.to_string()))?; fs::write(&ids_path, ids_json)?; + // Persist tombstones independently so older indexes without this file + // remain loadable and are interpreted as having no deletions. + let tombstones: Vec = self + .tombstones + .iter() + .flat_map(|word| word.to_le_bytes()) + .collect(); + fs::write(dir.join("tombstones.bin"), tombstones)?; + // Save PQ if present if let Some(ref pq) = self.pq { let pq_path = dir.join("pq.bin"); @@ -282,7 +333,8 @@ impl DiskAnnIndex { "search_beam": self.config.search_beam, "alpha": self.config.alpha, "pq_subspaces": self.config.pq_subspaces, - "count": self.vectors.len(), + "count": self.count(), + "vector_count": self.vectors.len(), "built": self.built, }); fs::write( @@ -346,9 +398,43 @@ impl DiskAnnIndex { let id_map: Vec = serde_json::from_str(&ids_json) .map_err(|e| DiskAnnError::Serialization(e.to_string()))?; + let tombstone_path = dir.join("tombstones.bin"); + let tombstones = if tombstone_path.exists() { + let bytes = fs::read(tombstone_path)?; + let word_count = n.div_ceil(64); + if bytes.len() != word_count * 8 { + return Err(DiskAnnError::Serialization( + "invalid tombstone data".to_string(), + )); + } + let words: Vec = bytes + .chunks_exact(8) + .map(|chunk| u64::from_le_bytes(chunk.try_into().unwrap())) + .collect(); + if n % 64 != 0 + && words + .last() + .is_some_and(|word| word & !((1u64 << (n % 64)) - 1) != 0) + { + return Err(DiskAnnError::Serialization( + "invalid tombstone data".to_string(), + )); + } + words + } else { + vec![0; n.div_ceil(64)] + }; + let deleted_count = tombstones + .iter() + .map(|word| word.count_ones() as usize) + .sum(); + let mut id_reverse = HashMap::new(); for (i, id) in id_map.iter().enumerate() { - id_reverse.insert(id.clone(), i as u32); + let is_tombstoned = tombstones[i / 64] & (1u64 << (i % 64)) != 0; + if !is_tombstoned { + id_reverse.insert(id.clone(), i as u32); + } } // Load graph @@ -403,6 +489,8 @@ impl DiskAnnIndex { vectors, id_map, id_reverse, + tombstones, + deleted_count, graph: Some(graph), pq, pq_codes, @@ -511,6 +599,191 @@ mod tests { assert_eq!(results[0].id, "vec-7"); } + #[test] + fn test_delete_filters_exact_match_and_preserves_vector() { + let mut index = DiskAnnIndex::new(DiskAnnConfig { + dim: 16, + max_degree: 8, + build_beam: 24, + search_beam: 24, + ..Default::default() + }); + let data = random_vectors(200, 16); + let query = data[42].1.clone(); + index.insert_batch(data).unwrap(); + index.build().unwrap(); + + let internal_id = index.id_reverse["vec-42"] as usize; + let before = index.vectors.get(internal_id).to_vec(); + index.delete("vec-42").unwrap(); + + assert_eq!(index.vectors.get(internal_id), before); + let graph = index.graph.as_ref().unwrap(); + assert!(graph.neighbors[internal_id].is_empty()); + assert!(graph + .neighbors + .iter() + .all(|neighbors| !neighbors.contains(&(internal_id as u32)))); + let first = index.search(&query, 10).unwrap(); + let second = index.search(&query, 10).unwrap(); + assert!(first.iter().all(|result| result.id != "vec-42")); + assert!(first.iter().all(|result| result.distance.is_finite())); + assert_eq!( + first + .iter() + .map(|result| (&result.id, result.distance.to_bits())) + .collect::>(), + second + .iter() + .map(|result| (&result.id, result.distance.to_bits())) + .collect::>() + ); + } + + #[test] + fn test_delete_count_and_errors() { + let mut index = DiskAnnIndex::new(DiskAnnConfig { + dim: 4, + ..Default::default() + }); + index.insert("a".to_string(), vec![0.0; 4]).unwrap(); + index.insert("b".to_string(), vec![1.0; 4]).unwrap(); + assert_eq!(index.count(), 2); + + assert!(matches!( + index.delete("missing"), + Err(DiskAnnError::NotFound(id)) if id == "missing" + )); + assert!(index.delete_deferred("a").unwrap()); + assert_eq!(index.count(), 1); + assert!(matches!( + index.delete("a"), + Err(DiskAnnError::NotFound(id)) if id == "a" + )); + } + + #[test] + fn test_tombstones_survive_save_load() { + let dir = tempdir().unwrap(); + let mut index = DiskAnnIndex::new(DiskAnnConfig { + dim: 8, + max_degree: 8, + build_beam: 16, + search_beam: 16, + ..Default::default() + }); + let data = random_vectors(100, 8); + let deleted_query = data[9].1.clone(); + index.insert_batch(data).unwrap(); + index.build().unwrap(); + index.delete_deferred("vec-9").unwrap(); + index.save(dir.path()).unwrap(); + assert_eq!( + fs::metadata(dir.path().join("tombstones.bin")) + .unwrap() + .len(), + 16 + ); + + let loaded = DiskAnnIndex::load(dir.path()).unwrap(); + assert_eq!(loaded.count(), 99); + assert!(loaded + .search(&deleted_query, 10) + .unwrap() + .iter() + .all(|result| result.id != "vec-9")); + } + + #[test] + #[ignore = "20k-vector recall regression; run explicitly for release validation"] + fn test_delete_repair_recall_regression_20k() { + use rand::prelude::*; + use std::collections::HashSet; + + fn recall_at_10( + index: &DiskAnnIndex, + data: &[(String, Vec)], + survivors: &HashSet, + queries: &[usize], + ) -> f64 { + let mut matches = 0usize; + for &query_idx in queries { + let query = &data[query_idx].1; + let mut exact: Vec<(usize, f32)> = survivors + .iter() + .map(|&idx| (idx, l2_squared(&data[idx].1, query))) + .collect(); + exact.sort_unstable_by(|a, b| a.1.total_cmp(&b.1)); + let truth: HashSet<&str> = exact + .iter() + .take(10) + .map(|(idx, _)| data[*idx].0.as_str()) + .collect(); + matches += index + .search(query, 10) + .unwrap() + .iter() + .filter(|result| truth.contains(result.id.as_str())) + .count(); + } + matches as f64 / (queries.len() * 10) as f64 + } + + let n = 20_000; + let dim = 32; + let mut rng = StdRng::seed_from_u64(0x679D_15CA); + let data: Vec<(String, Vec)> = (0..n) + .map(|idx| { + ( + format!("v{idx}"), + (0..dim).map(|_| rng.gen::()).collect(), + ) + }) + .collect(); + let config = DiskAnnConfig { + dim, + max_degree: 32, + build_beam: 64, + search_beam: 64, + alpha: 1.2, + ..Default::default() + }; + + let dir = tempdir().unwrap(); + let mut base = DiskAnnIndex::new(config.clone()); + base.insert_batch(data.clone()).unwrap(); + base.build().unwrap(); + base.save(dir.path()).unwrap(); + let mut repaired = DiskAnnIndex::load(dir.path()).unwrap(); + let mut deferred = DiskAnnIndex::load(dir.path()).unwrap(); + + let mut order: Vec = (0..n).collect(); + order.shuffle(&mut rng); + let deleted: HashSet = order[..n / 5].iter().copied().collect(); + let survivors: HashSet = (0..n).filter(|idx| !deleted.contains(idx)).collect(); + for &idx in &order[..n / 5] { + repaired.delete(&format!("v{idx}")).unwrap(); + deferred.delete_deferred(&format!("v{idx}")).unwrap(); + } + + let survivor_data: Vec<(String, Vec)> = + survivors.iter().map(|&idx| data[idx].clone()).collect(); + let mut fresh = DiskAnnIndex::new(config); + fresh.insert_batch(survivor_data).unwrap(); + fresh.build().unwrap(); + let queries: Vec = order[n / 5..n / 5 + 100].to_vec(); + let repaired_recall = recall_at_10(&repaired, &data, &survivors, &queries); + let deferred_recall = recall_at_10(&deferred, &data, &survivors, &queries); + let fresh_recall = recall_at_10(&fresh, &data, &survivors, &queries); + println!( + "delete recall@10: deferred={deferred_recall:.4} repaired={repaired_recall:.4} fresh={fresh_recall:.4}" + ); + assert!( + (repaired_recall - fresh_recall).abs() <= 0.02, + "repaired recall {repaired_recall:.4} differs from fresh recall {fresh_recall:.4} by more than 0.02" + ); + } + #[test] fn test_recall_at_10() { // Measure recall@10: what fraction of true top-10 neighbors does DiskANN find?