mod common;
use std::sync::Arc;
use laurus::storage::Storage;
use laurus::storage::memory::{MemoryStorage, MemoryStorageConfig};
use laurus::vector::index::VectorIndex;
use laurus::vector::index::config::HnswIndexConfig;
use laurus::vector::index::hnsw::HnswIndex;
use laurus::vector::index::hnsw::segmented::SegmentedHnswIndex;
use laurus::vector::search::searcher::{VectorIndexQuery, VectorIndexQueryParams};
use laurus::vector::{DistanceMetric, Vector};
fn doc_vec(i: u64) -> Vector {
let mut v = vec![0.0f32; 16];
let t = i as f32 * 0.001;
v[0] = t.cos();
v[1] = t.sin();
v[2] = (t * 2.0).cos();
v[3] = (t * 3.0).sin();
Vector::new(v)
}
fn config(segmented: bool) -> HnswIndexConfig {
HnswIndexConfig {
dimension: 16,
m: 16,
ef_construction: 100,
normalize_vectors: false,
distance_metric: DistanceMetric::Cosine,
segmented,
..Default::default()
}
}
fn storage() -> Arc<MemoryStorage> {
Arc::new(MemoryStorage::new(MemoryStorageConfig::default()))
}
fn commit_batch(index: &dyn VectorIndex, ids: std::ops::Range<u64>) {
let mut writer = index.writer().unwrap();
let vectors: Vec<_> = ids.map(|i| (i, "v".to_string(), doc_vec(i))).collect();
writer.add_vectors(vectors).unwrap();
writer.commit().unwrap();
}
fn query(id: u64, top_k: usize) -> VectorIndexQuery {
VectorIndexQuery {
query: doc_vec(id),
params: VectorIndexQueryParams {
top_k,
..Default::default()
},
field_name: Some("v".to_string()),
filter: None,
}
}
#[test]
fn one_doc_commit_writes_o_delta_bytes() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
commit_batch(&index, 0..1000);
let base_file = "segment_000000.hnsw";
let base_size = storage.file_size(base_file).unwrap();
commit_batch(&index, 1000..1001);
let delta_file = "segment_000001.hnsw";
let delta_size = storage.file_size(delta_file).unwrap();
assert!(
delta_size * 20 < base_size,
"a 1-doc commit must write O(delta) bytes, got delta={delta_size} vs base={base_size}"
);
assert_eq!(
storage.file_size(base_file).unwrap(),
base_size,
"the base segment must never be rewritten by a later commit (#634)"
);
let hnsw_files: Vec<String> = storage
.list_files()
.unwrap()
.into_iter()
.filter(|f| f.ends_with(".hnsw"))
.collect();
assert_eq!(
hnsw_files.len(),
2,
"exactly one new segment per non-empty commit, got {hnsw_files:?}"
);
}
#[test]
fn config_off_keeps_monolithic_layout() {
let storage = storage();
let index = HnswIndex::create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(false),
)
.unwrap();
commit_batch(&index, 0..100);
assert!(
!storage.file_exists("segments.json"),
"config OFF must not create a segment manifest"
);
assert!(storage.file_exists("vector_index.hnsw"));
}
#[test]
fn legacy_monolithic_index_migrates_zero_copy() {
let storage = storage();
{
let index = HnswIndex::create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(false),
)
.unwrap();
commit_batch(&index, 0..100);
}
let legacy_size = storage.file_size("vector_index.hnsw").unwrap();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
assert!(storage.file_exists("segments.json"), "manifest created");
assert_eq!(
storage.file_size("vector_index.hnsw").unwrap(),
legacy_size,
"zero-copy: the legacy file must not be rewritten (#882)"
);
let searcher = index.searcher().unwrap();
let results = searcher.search(&query(42, 1)).unwrap();
assert_eq!(results.results[0].doc_id, 42);
commit_batch(&index, 100..101);
assert_eq!(storage.file_size("vector_index.hnsw").unwrap(), legacy_size);
let searcher = index.searcher().unwrap();
assert_eq!(
searcher.search(&query(100, 1)).unwrap().results[0].doc_id,
100
);
assert_eq!(
searcher.search(&query(42, 1)).unwrap().results[0].doc_id,
42
);
drop(index);
let index = SegmentedHnswIndex::open_or_create(
storage as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
let searcher = index.searcher().unwrap();
assert_eq!(
searcher.search(&query(42, 1)).unwrap().results[0].doc_id,
42
);
assert_eq!(
searcher.search(&query(100, 1)).unwrap().results[0].doc_id,
100
);
}
#[test]
fn wal_checkpoint_publishes_only_after_persist_deletions() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
index.set_last_wal_seq(7).unwrap();
assert_eq!(
index.last_wal_seq(),
0,
"a pending seq must not be visible before persist_deletions (#882)"
);
commit_batch(&index, 0..5);
assert_eq!(
index.last_wal_seq(),
0,
"sealing must not publish the pending checkpoint (#882)"
);
index.persist_deletions().unwrap();
assert_eq!(index.last_wal_seq(), 7);
}
#[test]
fn multi_segment_self_recall_matches_monolithic() {
let n = 2000u64;
let per_commit = 400u64;
let seg_storage = storage();
let seg_index = SegmentedHnswIndex::open_or_create(
seg_storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
let mut lo = 0u64;
while lo < n {
commit_batch(&seg_index, lo..(lo + per_commit).min(n));
lo += per_commit;
}
let mono_storage = storage();
let mono_index = HnswIndex::create(
mono_storage.clone() as Arc<dyn Storage>,
"vector_index",
config(false),
)
.unwrap();
commit_batch(&mono_index, 0..n);
let recall = |index: &dyn VectorIndex| -> f32 {
let searcher = index.searcher().unwrap();
let mut hits = 0u64;
for id in 0..n {
let results = searcher.search(&query(id, 10)).unwrap();
if results.results.iter().any(|r| r.doc_id == id) {
hits += 1;
}
}
hits as f32 / n as f32
};
let mono = recall(&mono_index);
let seg = recall(&seg_index);
assert!(mono > 0.9, "monolithic self-recall sanity, got {mono:.4}");
assert!(
seg >= mono - 0.05,
"multi-segment self-recall ({seg:.4}) must match the monolithic build \
({mono:.4}) within tolerance (#881)"
);
}
#[test]
fn upsert_and_soft_delete_across_commits() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
commit_batch(&index, 1..50);
{
let mut writer = index.writer().unwrap();
writer.delete_document(1).unwrap();
writer
.add_vectors(vec![(1, "v".to_string(), doc_vec(9000))])
.unwrap();
writer.commit().unwrap();
}
let searcher = index.searcher().unwrap();
let results = searcher.search(&query(9000, 1)).unwrap();
assert_eq!(results.results[0].doc_id, 1, "newest copy must win");
let results = searcher.search(&query(1, 1)).unwrap();
assert_ne!(
results.results[0].doc_id, 1,
"the stale copy in the older segment must be masked (#880/#881)"
);
index.soft_delete_document(10).unwrap();
index.persist_deletions().unwrap();
let searcher = index.searcher().unwrap();
let results = searcher.search(&query(10, 5)).unwrap();
assert!(
results.results.iter().all(|r| r.doc_id != 10),
"a soft-deleted sealed doc must be search-invisible"
);
index.optimize().unwrap();
let hnsw_files: Vec<String> = storage
.list_files()
.unwrap()
.into_iter()
.filter(|f| f.ends_with(".hnsw"))
.collect();
assert_eq!(
hnsw_files.len(),
1,
"optimize must force-merge to one segment"
);
assert!(
!storage.file_exists("vector_index.delmap"),
"optimize must clear the persisted deletion bitmap"
);
let searcher = index.searcher().unwrap();
let results = searcher.search(&query(10, 5)).unwrap();
assert!(results.results.iter().all(|r| r.doc_id != 10));
let results = searcher.search(&query(9000, 1)).unwrap();
assert_eq!(
results.results[0].doc_id, 1,
"upserted copy survives the merge"
);
let stats = index.stats().unwrap();
assert_eq!(stats.vector_count, 48);
}
#[test]
fn undelete_to_zero_removes_stale_delmap_and_survives_reopen() {
let storage = storage();
{
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
commit_batch(&index, 1..10);
index.soft_delete_document(3).unwrap();
index.persist_deletions().unwrap();
assert!(storage.file_exists("vector_index.delmap"));
let mut writer = index.writer().unwrap();
writer.delete_document(3).unwrap();
writer
.add_vectors(vec![(3, "v".to_string(), doc_vec(9000))])
.unwrap();
writer.commit().unwrap();
index.persist_deletions().unwrap();
assert!(
!storage.file_exists("vector_index.delmap"),
"undelete-to-zero must remove the stale delmap (#881)"
);
}
let index = SegmentedHnswIndex::open_or_create(
storage as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
let searcher = index.searcher().unwrap();
let results = searcher.search(&query(9000, 1)).unwrap();
assert_eq!(
results.results[0].doc_id, 3,
"the committed upsert must survive a reopen (#881)"
);
}
#[test]
fn sealed_writer_rejects_second_commit_and_close_is_noop() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
let mut writer = index.writer().unwrap();
writer
.add_vectors(vec![(1, "v".to_string(), doc_vec(1))])
.unwrap();
writer.commit().unwrap();
assert!(!writer.has_pending_changes());
writer.close().unwrap();
let mut writer = index.writer().unwrap();
writer
.add_vectors(vec![(2, "v".to_string(), doc_vec(2))])
.unwrap();
writer.commit().unwrap();
writer
.add_vectors(vec![(3, "v".to_string(), doc_vec(3))])
.unwrap();
let err = writer.commit();
assert!(
err.is_err(),
"a sealed writer must reject a second commit with new changes (#881)"
);
}
#[test]
fn count_excludes_soft_deleted_docs() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
commit_batch(&index, 0..10);
index.soft_delete_document(4).unwrap();
let searcher = index.searcher().unwrap();
let count = searcher.count(query(0, 1)).unwrap();
assert_eq!(count, 9, "count must exclude soft-deleted docs (#881)");
}
#[test]
fn factory_open_path_migrates_legacy_index_and_stays_segmented() {
use laurus::vector::index::config::VectorIndexTypeConfig;
use laurus::vector::index::factory::VectorIndexFactory;
let storage = storage();
{
let index = HnswIndex::create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(false),
)
.unwrap();
commit_batch(&index, 0..50);
}
assert!(storage.file_exists("metadata.json"), "legacy precondition");
let index = VectorIndexFactory::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
VectorIndexTypeConfig::HNSW(config(true)),
)
.unwrap();
assert!(
storage.file_exists("segments.json"),
"the factory OPEN arm must migrate a legacy index (#882)"
);
assert!(
!storage.file_exists("metadata.json"),
"the stale monolithic metadata.json must be removed so factory \
routing can never regress to the monolithic view (#882)"
);
commit_batch(index.as_ref(), 50..51);
drop(index);
let index = VectorIndexFactory::open_or_create(
storage as Arc<dyn Storage>,
"vector_index",
VectorIndexTypeConfig::HNSW(config(true)),
)
.unwrap();
let searcher = index.searcher().unwrap();
assert_eq!(
searcher.search(&query(42, 1)).unwrap().results[0].doc_id,
42
);
assert_eq!(
searcher.search(&query(50, 1)).unwrap().results[0].doc_id,
50
);
}
#[test]
fn factory_rejects_segmented_directory_with_flag_off() {
use laurus::vector::index::config::VectorIndexTypeConfig;
use laurus::vector::index::factory::VectorIndexFactory;
let storage = storage();
{
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
commit_batch(&index, 0..10);
}
let result = VectorIndexFactory::open_or_create(
storage as Arc<dyn Storage>,
"vector_index",
VectorIndexTypeConfig::HNSW(config(false)),
);
assert!(
result.is_err(),
"a segmented directory must not open monolithically (#882)"
);
}
#[test]
fn append_only_segment_count_is_bounded_by_auto_merge() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
for i in 0..101u64 {
commit_batch(&index, i..i + 1);
}
let before = storage
.list_files()
.unwrap()
.iter()
.filter(|f| f.ends_with(".hnsw"))
.count();
assert!(before > 100, "precondition: {before} segments");
let compacted = index.maybe_auto_compact().unwrap();
assert!(compacted, "the segment-count bound must trigger a merge");
let after = storage
.list_files()
.unwrap()
.iter()
.filter(|f| f.ends_with(".hnsw"))
.count();
assert!(
after < before,
"the merge must reduce the segment count ({before} -> {after})"
);
let searcher = index.searcher().unwrap();
for id in [0u64, 50, 100] {
assert_eq!(
searcher.search(&query(id, 1)).unwrap().results[0].doc_id,
id
);
}
}
#[test]
fn segmented_flag_serde_default_and_explicit_false() {
let mut value: serde_json::Value = serde_json::to_value(HnswIndexConfig::default()).unwrap();
value
.as_object_mut()
.unwrap()
.remove("segmented")
.expect("the flag must serialize");
let config: HnswIndexConfig = serde_json::from_value(value).unwrap();
assert!(
config.segmented,
"a config serialized before the field existed must open segmented (#882)"
);
let explicit_false = serde_json::to_string(&HnswIndexConfig {
segmented: false,
..HnswIndexConfig::default()
})
.unwrap();
let config: HnswIndexConfig = serde_json::from_str(&explicit_false).unwrap();
assert!(
!config.segmented,
"an explicit `segmented: false` must be preserved (#882)"
);
}
#[test]
fn sustained_ingest_keeps_segment_count_tiered() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
for i in 0..100u64 {
commit_batch(&index, i * 5..(i + 1) * 5);
index.maybe_auto_compact().unwrap();
}
let segments = storage
.list_files()
.unwrap()
.iter()
.filter(|f| f.ends_with(".hnsw"))
.count();
assert!(
segments <= 30,
"tiered merging must keep the segment count bounded under \
sustained ingest, got {segments} after 100 commits (#883)"
);
let searcher = index.searcher().unwrap();
for id in [0u64, 250, 499] {
let results = searcher.search(&query(id, 5)).unwrap();
assert!(
results.results.iter().any(|r| r.doc_id == id),
"doc {id} must stay searchable after tiered merging"
);
}
}
#[test]
fn adaptive_refill_recovers_live_doc_behind_stale_band() {
let storage = storage();
let index = SegmentedHnswIndex::open_or_create(
storage.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap();
{
let mut writer = index.writer().unwrap();
let mut batch: Vec<_> = (1..=30u64)
.map(|i| (i, "v".to_string(), doc_vec(i)))
.collect();
batch.push((100, "v".to_string(), doc_vec(31)));
writer.add_vectors(batch).unwrap();
writer.commit().unwrap();
}
{
let mut writer = index.writer().unwrap();
for id in 1..=30u64 {
writer.delete_document(id).unwrap();
writer
.add_vectors(vec![(id, "v".to_string(), doc_vec(9000 + id))])
.unwrap();
}
writer.commit().unwrap();
}
let searcher = index.searcher().unwrap();
let results = searcher.search(&query(0, 1)).unwrap();
assert_eq!(
results.results.first().map(|r| r.doc_id),
Some(100),
"the live doc behind the deep stale band must surface via the \
expanding adaptive refill (#883), got {:?}",
results.results
);
}
#[test]
fn auto_commit_cumulative_bytes_are_bounded() {
use common::ByteCountingStorage;
let run = |segmented: bool| -> u64 {
let counting = Arc::new(ByteCountingStorage::new(Arc::new(MemoryStorage::new(
MemoryStorageConfig::default(),
))));
let written = counting.written.clone();
let index: Box<dyn VectorIndex> = if segmented {
Box::new(
SegmentedHnswIndex::open_or_create(
counting.clone() as Arc<dyn Storage>,
"vector_index",
config(true),
)
.unwrap(),
)
} else {
Box::new(
HnswIndex::create(
counting.clone() as Arc<dyn Storage>,
"vector_index",
config(false),
)
.unwrap(),
)
};
for i in 0..40u64 {
commit_batch(index.as_ref(), i * 50..(i + 1) * 50);
index.maybe_auto_compact().unwrap();
}
written.load(std::sync::atomic::Ordering::Relaxed)
};
let monolithic = run(false);
let segmented = run(true);
eprintln!(
"auto-commit cumulative bytes: monolithic={monolithic} segmented={segmented} \
ratio={:.1}x",
monolithic as f64 / segmented as f64
);
assert!(
segmented * 5 < monolithic,
"the segmented layout must write at least 5x fewer cumulative bytes \
under auto-commit ingest, got monolithic={monolithic} vs segmented={segmented} (#883)"
);
}