use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use laurus::storage::{FileMetadata, LoadingMode, Storage, StorageInput, StorageOutput};
use laurus::{Document, Result};
use laurus::lexical::{
InvertedIndexWriter, InvertedIndexWriterConfig, LexicalIndexConfig, LexicalSearchRequest,
LexicalStore, TermQuery,
};
use laurus::storage::memory::{MemoryStorage, MemoryStorageConfig};
#[derive(Debug)]
struct CountingStorage {
inner: MemoryStorage,
list_files_calls: AtomicUsize,
fail_list_files: std::sync::atomic::AtomicBool,
}
impl CountingStorage {
fn new() -> Self {
Self {
inner: MemoryStorage::new(MemoryStorageConfig::default()),
list_files_calls: AtomicUsize::new(0),
fail_list_files: std::sync::atomic::AtomicBool::new(false),
}
}
fn list_files_count(&self) -> usize {
self.list_files_calls.load(Ordering::SeqCst)
}
fn set_fail_list_files(&self, fail: bool) {
self.fail_list_files.store(fail, Ordering::SeqCst);
}
}
impl Storage for CountingStorage {
fn loading_mode(&self) -> LoadingMode {
self.inner.loading_mode()
}
fn open_input(&self, name: &str) -> Result<Box<dyn StorageInput>> {
self.inner.open_input(name)
}
fn create_output(&self, name: &str) -> Result<Box<dyn StorageOutput>> {
self.inner.create_output(name)
}
fn create_output_append(&self, name: &str) -> Result<Box<dyn StorageOutput>> {
self.inner.create_output_append(name)
}
fn file_exists(&self, name: &str) -> bool {
self.inner.file_exists(name)
}
fn delete_file(&self, name: &str) -> Result<()> {
self.inner.delete_file(name)
}
fn list_files(&self) -> Result<Vec<String>> {
if self.fail_list_files.load(Ordering::SeqCst) {
return Err(laurus::LaurusError::storage("injected list_files failure"));
}
self.list_files_calls.fetch_add(1, Ordering::SeqCst);
self.inner.list_files()
}
fn file_size(&self, name: &str) -> Result<u64> {
self.inner.file_size(name)
}
fn metadata(&self, name: &str) -> Result<FileMetadata> {
self.inner.metadata(name)
}
fn rename_file(&self, old_name: &str, new_name: &str) -> Result<()> {
self.inner.rename_file(old_name, new_name)
}
fn create_temp_output(&self, prefix: &str) -> Result<(String, Box<dyn StorageOutput>)> {
self.inner.create_temp_output(prefix)
}
fn sync(&self) -> Result<()> {
self.inner.sync()
}
fn close(&mut self) -> Result<()> {
self.inner.close()
}
}
fn doc(title: &str) -> Document {
Document::builder().add_text("title", title).build()
}
fn seed_committed_segment(storage: &Arc<dyn Storage>, ids: &[u64]) {
let mut writer =
InvertedIndexWriter::new(storage.clone(), InvertedIndexWriterConfig::default()).unwrap();
for &id in ids {
writer
.upsert_document(id, doc(&format!("seed{id}")))
.unwrap();
}
writer.commit().unwrap();
}
#[test]
fn fresh_id_upserts_never_rescan_meta_files() {
let counting = Arc::new(CountingStorage::new());
let storage: Arc<dyn Storage> = counting.clone();
seed_committed_segment(&storage, &[1, 2, 3]);
let before_construction = counting.list_files_count();
let mut writer =
InvertedIndexWriter::new(storage.clone(), InvertedIndexWriterConfig::default()).unwrap();
assert_eq!(
counting.list_files_count(),
before_construction + 1,
"the constructor performs exactly one recovery scan"
);
let after_construction = counting.list_files_count();
for id in 10..60u64 {
writer.upsert_document(id, doc(&format!("t{id}"))).unwrap();
}
assert_eq!(
counting.list_files_count(),
after_construction,
"50 fresh-id upserts must not list the storage at all \
(pre-#864 behavior: one full .meta list+parse per upsert)"
);
}
#[test]
fn overwrite_committed_docs_reuses_cache_and_manager() {
let counting = Arc::new(CountingStorage::new());
let storage: Arc<dyn Storage> = counting.clone();
seed_committed_segment(&storage, &[1, 2, 3]);
let mut writer =
InvertedIndexWriter::new(storage.clone(), InvertedIndexWriterConfig::default()).unwrap();
let after_construction = counting.list_files_count();
writer.upsert_document(1, doc("one-v2")).unwrap();
assert_eq!(
counting.list_files_count(),
after_construction + 1,
"the first overwrite pays exactly the DeletionManager's one-time \
bitmap-loading scan"
);
writer.upsert_document(2, doc("two-v2")).unwrap();
writer.upsert_document(3, doc("three-v2")).unwrap();
assert_eq!(
counting.list_files_count(),
after_construction + 1,
"subsequent overwrites must reuse the cached manager and ranges \
(pre-#864: fresh manager + full .delmap reload per overwrite)"
);
writer.flush_deletions().unwrap();
let delmaps: Vec<String> = storage
.list_files()
.unwrap()
.into_iter()
.filter(|f| f.ends_with(".delmap"))
.collect();
assert_eq!(
delmaps.len(),
1,
"the overwrites must have marked deletions in the seeded segment: {delmaps:?}"
);
}
#[test]
fn mid_life_flush_extends_cache() {
let counting = Arc::new(CountingStorage::new());
let storage: Arc<dyn Storage> = counting.clone();
let config = InvertedIndexWriterConfig {
max_buffered_docs: 2, ..Default::default()
};
let mut writer = InvertedIndexWriter::new(storage.clone(), config).unwrap();
writer.upsert_document(1, doc("alpha")).unwrap();
writer.upsert_document(2, doc("bravo")).unwrap();
let after_flush = counting.list_files_count();
writer.upsert_document(1, doc("alpha-v2")).unwrap();
assert_eq!(
counting.list_files_count(),
after_flush + 1,
"the overwrite must resolve the mid-life flushed segment from the \
cache, paying only the one-time DeletionManager construction"
);
writer.flush_deletions().unwrap();
let delmaps: Vec<String> = storage
.list_files()
.unwrap()
.into_iter()
.filter(|f| f.ends_with(".delmap"))
.collect();
assert_eq!(
delmaps.len(),
1,
"the overwrite must have marked a deletion in the flushed segment: {delmaps:?}"
);
}
#[test]
fn optimize_rebuilds_live_writer_cache() {
let storage: Arc<dyn Storage> = Arc::new(MemoryStorage::new(MemoryStorageConfig::default()));
let store = LexicalStore::new(storage.clone(), LexicalIndexConfig::default()).unwrap();
let hits = |field: &str, term: &str| -> usize {
let query = Box::new(TermQuery::new(field, term));
store
.search(LexicalSearchRequest::new(query))
.unwrap()
.hits
.len()
};
store.upsert_document(1, doc("alpha")).unwrap();
store.commit().unwrap();
store.upsert_document(2, doc("bravo")).unwrap();
store.commit().unwrap();
store.upsert_document(3, doc("charlie")).unwrap();
store.optimize().unwrap();
store.upsert_document(1, doc("alphav2")).unwrap();
store.commit().unwrap();
assert_eq!(
hits("title", "alpha"),
0,
"the pre-merge version must be dead — a stale segment cache would \
have marked the deletion in a ghost segment and left it alive"
);
assert_eq!(hits("title", "alphav2"), 1, "the overwrite must be live");
assert_eq!(
hits("title", "bravo"),
1,
"untouched doc survives the merge"
);
assert_eq!(
hits("title", "charlie"),
1,
"buffered doc survives the merge"
);
let delmaps: Vec<String> = storage
.list_files()
.unwrap()
.into_iter()
.filter(|f| f.ends_with(".delmap"))
.collect();
assert!(
delmaps.iter().all(|f| f.starts_with("merged_")),
"deletions must land in the merged segment only: {delmaps:?}"
);
}
#[test]
fn invalidate_failure_preserves_old_cache() {
let counting = Arc::new(CountingStorage::new());
let storage: Arc<dyn Storage> = counting.clone();
seed_committed_segment(&storage, &[1, 2, 3]);
let mut writer =
InvertedIndexWriter::new(storage.clone(), InvertedIndexWriterConfig::default()).unwrap();
counting.set_fail_list_files(true);
writer
.invalidate_segment_cache()
.expect_err("the injected list_files failure must propagate");
counting.set_fail_list_files(false);
writer.upsert_document(1, doc("one-v2")).unwrap();
writer.flush_deletions().unwrap();
let delmaps: Vec<String> = storage
.list_files()
.unwrap()
.into_iter()
.filter(|f| f.ends_with(".delmap"))
.collect();
assert_eq!(
delmaps.len(),
1,
"the overwrite must still mark the deletion from the preserved \
cache — an emptied cache would have skipped it silently: {delmaps:?}"
);
}