use std::cell::RefCell;
use std::cmp::Ordering;
use std::collections::{HashMap, HashSet};
use std::io::Read;
use std::sync::Arc;
use durability::{Directory, PersistenceError, PersistenceResult};
use segstore::{SegmentedStore, Store};
use crate::distance;
use crate::hnsw::{HNSWIndex, HNSWParams};
struct VectorBacking;
impl Store for VectorBacking {
type Id = u32;
type Item = Vec<f32>;
type Segment = Vec<(u32, Vec<f32>)>;
fn build_segment(&self, batch: &[(u32, Vec<f32>)]) -> Vec<(u32, Vec<f32>)> {
batch.to_vec()
}
fn merge_segments(
&self,
segs: &[&Vec<(u32, Vec<f32>)>],
live: &dyn Fn(&u32) -> bool,
) -> Vec<(u32, Vec<f32>)> {
segs.iter()
.flat_map(|s| s.iter())
.filter(|(id, _)| live(id))
.cloned()
.collect()
}
fn segment_len(&self, seg: &Vec<(u32, Vec<f32>)>) -> usize {
seg.len()
}
fn live_len(&self, seg: &Vec<(u32, Vec<f32>)>, live: &dyn Fn(&u32) -> bool) -> Option<usize> {
Some(seg.iter().filter(|(id, _)| live(id)).count())
}
}
struct Cache {
by_ptr: HashMap<usize, Option<HNSWIndex>>,
}
const INDEX_KIND: &str = "hnsw";
const SIDECAR_MAGIC: &[u8; 8] = b"VICHNSW1";
const SIDECAR_VERSION: u32 = 1;
pub struct UpdatableIndex {
inner: SegmentedStore<VectorBacking>,
dim: usize,
m: usize,
m_max: usize,
sidecar_recipe: String,
cache: RefCell<Cache>,
persisted: RefCell<HashSet<u64>>,
}
impl UpdatableIndex {
pub fn open(
dir: Arc<dyn Directory>,
flush_threshold: usize,
dim: usize,
m: usize,
m_max: usize,
) -> PersistenceResult<Self> {
let inner = SegmentedStore::open(dir, VectorBacking, flush_threshold)?;
Ok(Self {
inner,
dim,
m,
m_max,
sidecar_recipe: Self::make_sidecar_recipe(dim, m, m_max),
cache: RefCell::new(Cache {
by_ptr: HashMap::new(),
}),
persisted: RefCell::new(HashSet::new()),
})
}
pub fn add(&mut self, id: u32, vector: &[f32]) -> PersistenceResult<()> {
if vector.len() != self.dim {
return Err(PersistenceError::InvalidConfig(format!(
"vector dimension {} does not match index dimension {}",
vector.len(),
self.dim
)));
}
self.inner.add(id, distance::normalize(vector))?;
Ok(())
}
pub fn extend(
&mut self,
vectors: impl IntoIterator<Item = (u32, Vec<f32>)>,
) -> PersistenceResult<()> {
let dim = self.dim;
let normalized: Result<Vec<(u32, Vec<f32>)>, PersistenceError> = vectors
.into_iter()
.map(|(id, vector)| {
if vector.len() != dim {
Err(PersistenceError::InvalidConfig(format!(
"vector dimension {} does not match index dimension {}",
vector.len(),
dim
)))
} else {
Ok((id, distance::normalize(&vector)))
}
})
.collect();
self.inner.extend(normalized?)?;
Ok(())
}
pub fn delete(&mut self, id: u32) -> PersistenceResult<()> {
self.inner.delete(id)?;
let ids = self.inner.segment_ids();
let mut cache = self.cache.borrow_mut();
for (i, seg) in self.inner.segments().iter().enumerate() {
if seg.iter().any(|(sid, _)| *sid == id) {
cache.by_ptr.remove(&(Arc::as_ptr(seg) as usize));
let seg_id = ids[i];
self.persisted.borrow_mut().remove(&seg_id);
let _ = self
.inner
.dir()
.delete(&self.inner.index_name(seg_id, INDEX_KIND));
}
}
Ok(())
}
pub fn compact(&mut self) -> PersistenceResult<()> {
self.inner.compact()?;
Ok(())
}
pub fn checkpoint(&mut self) -> PersistenceResult<()> {
self.inner.checkpoint()?;
self.persist_new_segments();
Ok(())
}
pub fn compact_tiers(&mut self) -> PersistenceResult<()> {
self.inner.compact_tiers()?;
Ok(())
}
pub fn reclaim(&mut self, min_live_ratio: f64) -> PersistenceResult<()> {
self.inner.reclaim_tombstones(min_live_ratio)?;
Ok(())
}
pub fn space_amplification(&self) -> Option<f64> {
self.inner.space_amplification()
}
pub fn search(&self, query: &[f32], k: usize, ef: usize) -> Vec<(u32, f32)> {
let q = distance::normalize(query);
let mut cand: Vec<(u32, f32)> = Vec::new();
{
let segs = self.inner.segments();
let ids = self.inner.segment_ids();
let mut cache = self.cache.borrow_mut();
let current: HashSet<usize> = segs.iter().map(|a| Arc::as_ptr(a) as usize).collect();
cache.by_ptr.retain(|key, _| current.contains(key));
for (i, seg) in segs.iter().enumerate() {
let key = Arc::as_ptr(seg) as usize;
let seg_id = ids[i];
cache
.by_ptr
.entry(key)
.or_insert_with(|| self.build_or_load(&seg[..], seg_id));
}
for idx in cache.by_ptr.values().flatten() {
cand.extend(idx.search(&q, k, ef).unwrap_or_default());
}
}
let buffered = self.inner.buffer().to_vec();
if let Some(idx) = self.build_live_index(&buffered) {
cand.extend(idx.search(&q, k, ef).unwrap_or_default());
}
cand.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(Ordering::Equal));
cand.truncate(k);
cand
}
fn build_live_index(&self, batch: &[(u32, Vec<f32>)]) -> Option<HNSWIndex> {
let mut idx = match HNSWIndex::new(self.dim, self.m, self.m_max) {
Ok(i) => i,
Err(_) => return None,
};
let mut any = false;
for (id, v) in batch {
if self.inner.is_live(id) && idx.add(*id, v.clone()).is_ok() {
any = true;
}
}
if !any || idx.build().is_err() {
return None;
}
Some(idx)
}
fn build_or_load(&self, seg: &[(u32, Vec<f32>)], seg_id: u64) -> Option<HNSWIndex> {
if let Some(idx) = self.load_sidecar(seg, seg_id) {
self.persisted.borrow_mut().insert(seg_id);
return Some(idx);
}
let idx = self.build_live_index(seg)?;
self.persist_sidecar(&idx, seg_id);
Some(idx)
}
fn load_sidecar(&self, seg: &[(u32, Vec<f32>)], seg_id: u64) -> Option<HNSWIndex> {
let name = self.inner.index_name(seg_id, INDEX_KIND);
if !self.inner.dir().exists(&name) {
return None;
}
let mut bytes = Vec::new();
self.inner
.dir()
.open_file(&name)
.ok()?
.read_to_end(&mut bytes)
.ok()?;
let graph_bytes = self.decode_sidecar(&bytes)?;
let idx = HNSWIndex::from_postcard(graph_bytes).ok()?;
let mut live = HashSet::with_capacity(seg.len());
for (id, _) in seg {
if self.inner.is_live(id) {
live.insert(*id);
}
}
if idx.doc_ids.len() == live.len() && idx.doc_ids.iter().all(|id| live.contains(id)) {
Some(idx)
} else {
None
}
}
fn persist_sidecar(&self, idx: &HNSWIndex, seg_id: u64) {
if let Ok(graph) = idx.to_postcard() {
let Some(bytes) = self.encode_sidecar(&graph) else {
return;
};
if self
.inner
.dir()
.atomic_write(&self.inner.index_name(seg_id, INDEX_KIND), &bytes)
.is_ok()
{
self.persisted.borrow_mut().insert(seg_id);
}
}
}
fn make_sidecar_recipe(dim: usize, m: usize, m_max: usize) -> String {
let params = HNSWParams {
m,
m_max,
..Default::default()
};
format!(
"vicinity-store-hnsw-v1;\
dim={};m={};m_max={};m_l={:.17};ef_construction={};\
metric={:?};normalization=store-l2-on-ingest-and-query;\
seed_selection={:?};diversification={:?};seed={:?};\
codec=postcard-hnsw-v1;id_compression={}",
dim,
params.m,
params.m_max,
params.m_l,
params.ef_construction,
params.metric,
params.seed_selection,
params.neighborhood_diversification,
params.seed,
cfg!(feature = "id-compression")
)
}
fn encode_sidecar(&self, graph: &[u8]) -> Option<Vec<u8>> {
let recipe = self.sidecar_recipe.as_bytes();
let recipe_len = u32::try_from(recipe.len()).ok()?;
let mut bytes = Vec::with_capacity(16 + recipe.len() + graph.len());
bytes.extend_from_slice(SIDECAR_MAGIC);
bytes.extend_from_slice(&SIDECAR_VERSION.to_le_bytes());
bytes.extend_from_slice(&recipe_len.to_le_bytes());
bytes.extend_from_slice(recipe);
bytes.extend_from_slice(graph);
Some(bytes)
}
fn decode_sidecar<'a>(&self, bytes: &'a [u8]) -> Option<&'a [u8]> {
if bytes.len() < 16 {
return None;
}
if &bytes[..8] != SIDECAR_MAGIC {
return None;
}
let version = u32::from_le_bytes(bytes[8..12].try_into().ok()?);
if version != SIDECAR_VERSION {
return None;
}
let recipe_len = u32::from_le_bytes(bytes[12..16].try_into().ok()?) as usize;
let recipe_start = 16usize;
let recipe_end = recipe_start.checked_add(recipe_len)?;
if bytes.len() < recipe_end {
return None;
}
if &bytes[recipe_start..recipe_end] != self.sidecar_recipe.as_bytes() {
return None;
}
Some(&bytes[recipe_end..])
}
fn persist_new_segments(&self) {
let ids = self.inner.segment_ids();
let id_set: HashSet<u64> = ids.iter().copied().collect();
self.persisted.borrow_mut().retain(|id| id_set.contains(id));
for (i, seg) in self.inner.segments().iter().enumerate() {
let seg_id = ids[i];
if self.persisted.borrow().contains(&seg_id) {
continue;
}
if self.load_sidecar(&seg[..], seg_id).is_some() {
self.persisted.borrow_mut().insert(seg_id);
continue;
}
if let Some(idx) = self.build_live_index(&seg[..]) {
self.persist_sidecar(&idx, seg_id);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use durability::MemoryDirectory;
use std::io::Read;
fn read_file(dir: &Arc<dyn Directory>, name: &str) -> Vec<u8> {
let mut bytes = Vec::new();
dir.open_file(name)
.unwrap()
.read_to_end(&mut bytes)
.unwrap();
bytes
}
fn checkpointed_store(dir: Arc<dyn Directory>, m: usize, m_max: usize) -> (String, Vec<u8>) {
let mut store = UpdatableIndex::open(dir, 4, 2, m, m_max).unwrap();
for i in 0..12u32 {
let angle = i as f32 * 0.37;
store.add(i, &[angle.cos(), angle.sin()]).unwrap();
}
store.checkpoint().unwrap();
let seg_id = store.inner.segment_ids()[0];
let name = store.inner.index_name(seg_id, INDEX_KIND);
let bytes = read_file(store.inner.dir(), &name);
(name, bytes)
}
#[test]
fn add_delete_compact_recover_through_real_hnsw() {
let dir = MemoryDirectory::arc();
{
let mut store = UpdatableIndex::open(dir.clone(), 2, 2, 16, 32).unwrap();
store.add(0, &[1.0, 0.0]).unwrap();
store.add(1, &[0.0, 1.0]).unwrap(); store.add(2, &[0.7, 0.7]).unwrap();
let top: Vec<u32> = store
.search(&[0.9, 0.1], 2, 16)
.into_iter()
.map(|(id, _)| id)
.collect();
assert_eq!(top.first(), Some(&0), "nearest to the x-axis is doc 0");
let again: Vec<u32> = store
.search(&[0.9, 0.1], 2, 16)
.into_iter()
.map(|(id, _)| id)
.collect();
assert_eq!(again.first(), Some(&0), "cached query is stable");
store.delete(0).unwrap();
let top: Vec<u32> = store
.search(&[0.9, 0.1], 1, 16)
.into_iter()
.map(|(id, _)| id)
.collect();
assert_eq!(top, vec![2], "after deleting 0, nearest is doc 2");
store.compact().unwrap();
assert_eq!(
store.search(&[0.9, 0.1], 1, 16).first().map(|(id, _)| *id),
Some(2)
);
}
let store = UpdatableIndex::open(dir, 2, 2, 16, 32).unwrap();
let top: Vec<u32> = store
.search(&[0.9, 0.1], 1, 16)
.into_iter()
.map(|(id, _)| id)
.collect();
assert_eq!(top, vec![2], "recovery preserves the search");
}
#[test]
fn checkpoint_persists_sidecars_and_reopen_loads_them() {
let dir = MemoryDirectory::arc();
{
let mut store = UpdatableIndex::open(dir.clone(), 4, 3, 16, 32).unwrap();
for i in 0..12u32 {
let a = i as f32;
store.add(i, &[a.cos(), a.sin(), 1.0]).unwrap();
}
store.checkpoint().unwrap();
let ids: Vec<u64> = store.inner.segment_ids().to_vec();
assert!(
!ids.is_empty(),
"12 adds at flush 4 seal at least one segment"
);
for id in &ids {
assert!(
store
.inner
.dir()
.exists(&store.inner.index_name(*id, INDEX_KIND)),
"segment {id} must have a persisted sidecar after checkpoint"
);
}
}
let store = UpdatableIndex::open(dir, 4, 3, 16, 32).unwrap();
assert!(
!store.search(&[1.0, 0.0, 1.0], 1, 16).is_empty(),
"search over loaded sidecars returns results"
);
}
#[test]
fn hnsw_sidecar_recipe_mismatch_rebuilds() {
let dir = MemoryDirectory::arc();
let (name, before) = checkpointed_store(dir.clone(), 16, 32);
assert_eq!(
&before[..SIDECAR_MAGIC.len()],
SIDECAR_MAGIC,
"new sidecars carry the vicinity HNSW envelope"
);
let store = UpdatableIndex::open(dir.clone(), 4, 2, 8, 16).unwrap();
let seg_id = store.inner.segment_ids()[0];
assert!(
store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.is_none(),
"sidecar built with m=16/m_max=32 must not load under m=8/m_max=16"
);
assert!(
!store.search(&[1.0, 0.0], 1, 16).is_empty(),
"mismatched sidecar falls back to rebuild"
);
let after = read_file(store.inner.dir(), &name);
assert_ne!(before, after, "rebuild overwrites the stale-recipe sidecar");
assert!(
store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.is_some(),
"rebuilt sidecar now matches the current recipe"
);
}
#[test]
fn hnsw_sidecar_envelope_rejects_corrupt_headers() {
let store = UpdatableIndex::open(MemoryDirectory::arc(), 4, 2, 16, 32).unwrap();
let graph = b"graph-bytes";
let bytes = store.encode_sidecar(graph).unwrap();
assert_eq!(store.decode_sidecar(&bytes), Some(graph.as_slice()));
assert!(store.decode_sidecar(&bytes[..8]).is_none());
let mut bad_magic = bytes.clone();
bad_magic[0] ^= 0xFF;
assert!(store.decode_sidecar(&bad_magic).is_none());
let mut bad_version = bytes.clone();
bad_version[8..12].copy_from_slice(&(SIDECAR_VERSION + 1).to_le_bytes());
assert!(store.decode_sidecar(&bad_version).is_none());
let mut bad_recipe_len = bytes.clone();
bad_recipe_len[12..16].copy_from_slice(&u32::MAX.to_le_bytes());
assert!(store.decode_sidecar(&bad_recipe_len).is_none());
let mut bad_recipe = bytes.clone();
bad_recipe[16] ^= 0x01;
assert!(store.decode_sidecar(&bad_recipe).is_none());
}
#[test]
fn hnsw_sidecar_invalid_graph_payload_rebuilds() {
let dir = MemoryDirectory::arc();
let (name, _) = checkpointed_store(dir.clone(), 16, 32);
{
let store = UpdatableIndex::open(dir.clone(), 4, 2, 16, 32).unwrap();
let corrupt = store.encode_sidecar(b"not-a-postcard-hnsw-graph").unwrap();
store.inner.dir().atomic_write(&name, &corrupt).unwrap();
}
let store = UpdatableIndex::open(dir.clone(), 4, 2, 16, 32).unwrap();
let seg_id = store.inner.segment_ids()[0];
assert!(
store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.is_none(),
"valid envelope with invalid graph bytes is rejected"
);
assert!(
!store.search(&[1.0, 0.0], 1, 16).is_empty(),
"invalid graph payload falls back to rebuild"
);
assert!(
store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.is_some(),
"rebuilt sidecar loads after the fallback"
);
}
#[test]
fn hnsw_sidecar_query_ef_does_not_invalidate_recipe() {
let dir = MemoryDirectory::arc();
let (name, before) = checkpointed_store(dir.clone(), 16, 32);
let store = UpdatableIndex::open(dir.clone(), 4, 2, 16, 32).unwrap();
let seg_id = store.inner.segment_ids()[0];
assert!(
store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.is_some(),
"sidecar loads before any query"
);
assert!(!store.search(&[1.0, 0.0], 3, 8).is_empty());
assert!(!store.search(&[1.0, 0.0], 3, 64).is_empty());
assert_eq!(
read_file(&dir, &name),
before,
"query-time ef is not part of the sidecar recipe and must not rewrite it"
);
}
#[test]
fn deleted_id_does_not_resurface_through_a_sidecar() {
let dir = MemoryDirectory::arc();
{
let mut store = UpdatableIndex::open(dir.clone(), 2, 2, 16, 32).unwrap();
store.add(0, &[1.0, 0.0]).unwrap();
store.add(1, &[0.95, 0.05]).unwrap();
store.add(2, &[0.0, 1.0]).unwrap();
store.checkpoint().unwrap(); store.delete(0).unwrap(); store.checkpoint().unwrap(); }
let store = UpdatableIndex::open(dir, 2, 2, 16, 32).unwrap();
let top: Vec<u32> = store
.search(&[1.0, 0.0], 3, 16)
.into_iter()
.map(|(id, _)| id)
.collect();
assert!(
!top.contains(&0),
"deleted id 0 must not resurface from a persisted sidecar"
);
assert!(
top.contains(&1),
"nearest live vector to the x-axis is id 1"
);
}
#[test]
fn checkpoint_after_replayed_delete_rewrites_stale_sidecar() {
let dir = MemoryDirectory::arc();
let (name, stale_bytes) = {
let mut store = UpdatableIndex::open(dir.clone(), 2, 2, 16, 32).unwrap();
store.add(0, &[1.0, 0.0]).unwrap();
store.add(1, &[0.95, 0.05]).unwrap();
store.add(2, &[0.0, 1.0]).unwrap();
store.checkpoint().unwrap();
let seg_id = store.inner.segment_ids()[0];
let name = store.inner.index_name(seg_id, INDEX_KIND);
let bytes = read_file(store.inner.dir(), &name);
store.inner.delete(0).unwrap();
(name, bytes)
};
let mut store = UpdatableIndex::open(dir.clone(), 2, 2, 16, 32).unwrap();
let seg_id = store.inner.segment_ids()[0];
assert!(
store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.is_none(),
"replayed tombstone must make the old sidecar stale"
);
store.checkpoint().unwrap();
let rewritten = read_file(&dir, &name);
assert_ne!(
rewritten, stale_bytes,
"checkpoint should rewrite stale sidecars even before search"
);
let idx = store
.load_sidecar(&store.inner.segments()[0][..], seg_id)
.expect("rewritten sidecar should be valid");
assert!(
!idx.doc_ids.contains(&0),
"rewritten sidecar must exclude the replayed delete"
);
assert!(
idx.doc_ids.contains(&1),
"rewritten sidecar should keep live ids from the segment"
);
}
}