use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use parking_lot::RwLock;
use crate::embedding::embedder::Embedder;
use crate::error::{LaurusError, Result};
use crate::maintenance::deletion::DeletionBitmap;
use crate::storage::Storage;
use crate::storage::structured::{StructReader, StructWriter};
use crate::vector::core::vector::Vector;
use crate::vector::index::config::FlatIndexConfig;
use crate::vector::index::flat::reader::FlatVectorIndexReader;
use crate::vector::index::flat::searcher::FlatVectorSearcher;
use crate::vector::index::flat::segment::merge_engine::MergeEngine;
use crate::vector::index::flat::writer::FlatIndexWriter;
use crate::vector::index::segment::fanout::{SegmentFanoutSearcher, SegmentedReaderFacade};
use crate::vector::index::segment::manager::{
ManagedSegmentInfo, MergeCandidate, SegmentManager, SegmentManagerConfig,
};
use crate::vector::index::segment::merge::MergeConfig;
use crate::vector::index::segment::reader_cache::SegmentedReaderCache;
use crate::vector::index::{VectorIndex, VectorIndexStats};
use crate::vector::reader::VectorIndexReader;
use crate::vector::search::searcher::VectorIndexSearcher;
use crate::vector::store::embedding_writer::EmbeddingVectorIndexWriter;
use crate::vector::writer::{VectorIndexWriter, VectorIndexWriterConfig};
#[derive(Debug)]
struct SegmentedShared {
name: String,
storage: Arc<dyn Storage>,
manager: Arc<SegmentManager>,
reader_cache: Arc<SegmentedReaderCache<FlatVectorIndexReader>>,
deletion: RwLock<Option<Arc<DeletionBitmap>>>,
bitmap_epoch: std::sync::atomic::AtomicU64,
pending_wal_seq: std::sync::atomic::AtomicU64,
}
impl SegmentedShared {
fn delmap_file_name(&self) -> String {
format!("{}.delmap", self.name)
}
fn load_or_get_bitmap(&self, create_if_missing: bool) -> Result<Option<Arc<DeletionBitmap>>> {
if let Some(bitmap) = self.deletion.read().as_ref() {
return Ok(Some(bitmap.clone()));
}
let mut guard = self.deletion.write();
if let Some(bitmap) = guard.as_ref() {
return Ok(Some(bitmap.clone()));
}
let file = self.delmap_file_name();
if self.storage.file_exists(&file) {
let input = self.storage.open_input(&file)?;
let mut reader = StructReader::new(input)?;
let bitmap = Arc::new(DeletionBitmap::read_from_storage(&mut reader)?);
*guard = Some(bitmap.clone());
return Ok(Some(bitmap));
}
if create_if_missing {
let bitmap = Arc::new(DeletionBitmap::new(self.name.clone(), 0, u64::MAX - 1));
*guard = Some(bitmap.clone());
self.reader_cache.clear();
self.bitmap_epoch.fetch_add(1, Ordering::Release);
return Ok(Some(bitmap));
}
Ok(None)
}
fn sealed_readers_newest_first(
&self,
config: &FlatIndexConfig,
) -> Result<Vec<Arc<FlatVectorIndexReader>>> {
let mut segments = self.manager.list_segments();
segments.sort_by_key(|s| std::cmp::Reverse(s.generation));
loop {
let epoch = self.bitmap_epoch.load(Ordering::Acquire);
let bitmap = self.load_or_get_bitmap(false)?;
let mut readers = Vec::with_capacity(segments.len());
for info in &segments {
let storage = self.storage.clone();
let metric = config.distance_metric;
let bitmap_for_loader = bitmap.clone();
let segment_id = info.segment_id.clone();
let reader = self.reader_cache.get_or_load(&segment_id, || {
let mut r = FlatVectorIndexReader::load(storage, &segment_id, metric)?;
if let Some(b) = bitmap_for_loader {
r.set_deletion_bitmap(b);
}
Ok(r)
})?;
readers.push(reader);
}
if self.bitmap_epoch.load(Ordering::Acquire) == epoch {
return Ok(readers);
}
self.reader_cache.clear();
}
}
}
#[derive(Debug)]
pub struct SegmentedFlatIndex {
shared: Arc<SegmentedShared>,
config: FlatIndexConfig,
closed: AtomicBool,
}
impl SegmentedFlatIndex {
pub fn open_or_create(
storage: Arc<dyn Storage>,
name: &str,
config: FlatIndexConfig,
) -> Result<Self> {
if let Some(ordinal) = name.strip_prefix("segment_")
&& !ordinal.is_empty()
&& ordinal.bytes().all(|b| b.is_ascii_digit())
{
return Err(LaurusError::invalid_config(format!(
"index name '{name}' collides with the reserved segment-id \
namespace (segment_<digits>)"
)));
}
let legacy_file = format!("{name}.flat");
let migrate = storage.file_exists(&legacy_file) && !storage.file_exists("segments.json");
let manager = Arc::new(SegmentManager::new(
SegmentManagerConfig::default(),
storage.clone(),
crate::vector::index::flat::segment::LAYOUT,
)?);
if migrate {
let (vector_count, dimension) = {
use std::io::Read;
let mut input = storage.open_input(&legacy_file)?;
let mut count_buf = [0u8; 4];
input.read_exact(&mut count_buf)?;
let mut dim_buf = [0u8; 4];
input.read_exact(&mut dim_buf)?;
(
u32::from_le_bytes(count_buf) as u64,
u32::from_le_bytes(dim_buf) as usize,
)
};
if dimension != config.dimension {
return Err(LaurusError::index(format!(
"Dimension mismatch during migration: stored {dimension}, config {}",
config.dimension
)));
}
let mut info = ManagedSegmentInfo::new(name.to_string(), vector_count, 0, 0);
info.has_deletions = storage.file_exists(&format!("{name}.delmap"));
manager.add_segment(info)?;
let _ = storage.delete_file("metadata.json");
}
Ok(Self {
shared: Arc::new(SegmentedShared {
name: name.to_string(),
storage,
manager,
reader_cache: Arc::new(SegmentedReaderCache::new()),
deletion: RwLock::new(None),
bitmap_epoch: std::sync::atomic::AtomicU64::new(0),
pending_wal_seq: std::sync::atomic::AtomicU64::new(0),
}),
config,
closed: AtomicBool::new(false),
})
}
fn check_closed(&self) -> Result<()> {
if self.closed.load(Ordering::SeqCst) {
return Err(LaurusError::InvalidOperation("Index is closed".to_string()));
}
Ok(())
}
fn merge_once(&self) -> Result<bool> {
use crate::vector::index::segment::merge_policy::TieredMergePolicy;
let Some(candidate) = self.shared.manager.check_merge(&TieredMergePolicy::new()) else {
return Ok(false);
};
let mut engine = MergeEngine::new(
MergeConfig::default(),
self.shared.storage.clone(),
self.config.clone(),
VectorIndexWriterConfig::default(),
);
if let Some(bitmap) = self.shared.load_or_get_bitmap(false)? {
engine.set_deletion_bitmap(bitmap);
}
let new_segment_id = self.shared.manager.generate_segment_id();
let result = engine.merge_segments(candidate.segments.clone(), new_segment_id)?;
let source_ids: Vec<String> = candidate
.segments
.iter()
.map(|s| s.segment_id.clone())
.collect();
self.shared
.manager
.apply_merge(candidate, result.merged_segment)?;
for id in &source_ids {
self.shared.reader_cache.invalidate(id);
}
Ok(true)
}
fn clear_deletions(&self) -> Result<()> {
*self.shared.deletion.write() = None;
let file = self.shared.delmap_file_name();
let _ = self.shared.storage.delete_file(&file);
Ok(())
}
}
impl VectorIndex for SegmentedFlatIndex {
fn reader(&self) -> Result<Arc<dyn VectorIndexReader>> {
self.check_closed()?;
let readers = self.shared.sealed_readers_newest_first(&self.config)?;
let readers: Vec<Arc<dyn VectorIndexReader>> = readers
.into_iter()
.map(|r| r as Arc<dyn VectorIndexReader>)
.collect();
let bitmap = self.shared.load_or_get_bitmap(false)?;
Ok(Arc::new(SegmentedReaderFacade::new(
readers,
bitmap,
self.config.dimension,
self.config.distance_metric,
)))
}
fn writer(&self) -> Result<Box<dyn VectorIndexWriter>> {
self.check_closed()?;
let segment_id = self.shared.manager.generate_segment_id();
let inner = FlatIndexWriter::with_storage(
self.config.clone(),
VectorIndexWriterConfig::default(),
&segment_id,
self.shared.storage.clone(),
)?;
let writer = SegmentedFlatWriter {
shared: self.shared.clone(),
segment_id,
inner,
sealed_len: None,
};
let embedder = self.embedder();
Ok(Box::new(EmbeddingVectorIndexWriter::new(
Box::new(writer),
embedder,
)))
}
fn storage(&self) -> &Arc<dyn Storage> {
&self.shared.storage
}
fn close(&self) -> Result<()> {
self.closed.store(true, Ordering::SeqCst);
Ok(())
}
fn is_closed(&self) -> bool {
self.closed.load(Ordering::SeqCst)
}
fn stats(&self) -> Result<VectorIndexStats> {
self.check_closed()?;
let total = self.shared.manager.total_vectors();
let deleted = match self.shared.load_or_get_bitmap(false)? {
Some(bitmap) => bitmap.deleted_count.load(Ordering::Relaxed),
None => 0,
};
let stats = self.shared.manager.stats();
Ok(VectorIndexStats {
vector_count: total.saturating_sub(deleted),
dimension: self.config.dimension,
total_size: stats.total_size,
deleted_count: deleted,
last_modified: 0,
})
}
fn retain_writer_after_commit(&self) -> bool {
false
}
fn optimize(&self) -> Result<()> {
self.check_closed()?;
let segments = self.shared.manager.list_segments();
if segments.is_empty() {
return Ok(());
}
let mut engine = MergeEngine::new(
MergeConfig::default(),
self.shared.storage.clone(),
self.config.clone(),
VectorIndexWriterConfig::default(),
);
if let Some(bitmap) = self.shared.load_or_get_bitmap(false)? {
engine.set_deletion_bitmap(bitmap);
}
let new_segment_id = self.shared.manager.generate_segment_id();
let result = engine.merge_segments(segments.clone(), new_segment_id)?;
let candidate = MergeCandidate {
total_vectors: segments.iter().map(|s| s.vector_count).sum(),
total_size: segments.iter().map(|s| s.size_bytes).sum(),
segments,
};
let source_ids: Vec<String> = candidate
.segments
.iter()
.map(|s| s.segment_id.clone())
.collect();
self.shared
.manager
.apply_merge(candidate, result.merged_segment)?;
for id in &source_ids {
self.shared.reader_cache.invalidate(id);
}
self.clear_deletions()?;
Ok(())
}
fn searcher(&self) -> Result<Box<dyn VectorIndexSearcher>> {
self.check_closed()?;
let readers = self.shared.sealed_readers_newest_first(&self.config)?;
let readers: Vec<Arc<dyn VectorIndexReader>> = readers
.into_iter()
.map(|r| r as Arc<dyn VectorIndexReader>)
.collect();
let bitmap = self.shared.load_or_get_bitmap(false)?;
Ok(Box::new(SegmentFanoutSearcher::new(
readers,
bitmap,
move |reader| {
Ok(Box::new(FlatVectorSearcher::new(reader)?) as Box<dyn VectorIndexSearcher>)
},
)))
}
fn embedder(&self) -> Arc<dyn Embedder> {
Arc::clone(&self.config.embedder)
}
fn last_wal_seq(&self) -> u64 {
self.shared.manager.last_wal_seq()
}
fn set_last_wal_seq(&self, seq: u64) -> Result<()> {
self.shared
.pending_wal_seq
.fetch_max(seq, Ordering::Release);
Ok(())
}
fn supports_soft_delete(&self) -> bool {
true
}
fn soft_delete_document(&self, doc_id: u64) -> Result<()> {
self.check_closed()?;
let bitmap = self
.shared
.load_or_get_bitmap(true)?
.ok_or_else(|| LaurusError::internal("deletion bitmap unexpectedly missing"))?;
let newly_deleted = bitmap.delete_document(doc_id)?;
if newly_deleted && !self.shared.manager.list_segments().is_empty() {
self.shared.manager.mark_all_has_deletions()?;
}
Ok(())
}
fn persist_deletions(&self) -> Result<()> {
let guard = self.shared.deletion.read();
if let Some(bitmap) = guard.as_ref() {
let file = self.shared.delmap_file_name();
if bitmap.deleted_count.load(Ordering::Relaxed) > 0 {
let tmp = format!("{file}.tmp");
let output = self.shared.storage.create_output(&tmp)?;
let mut writer = StructWriter::new(output);
bitmap.write_to_storage(&mut writer)?;
writer.close()?;
self.shared.storage.rename_file(&tmp, &file)?;
} else if self.shared.storage.file_exists(&file) {
self.shared.storage.delete_file(&file)?;
}
}
drop(guard);
let pending = self.shared.pending_wal_seq.load(Ordering::Acquire);
if pending > self.shared.manager.last_wal_seq() {
self.shared.manager.set_last_wal_seq(pending);
self.shared.manager.save_state()?;
}
Ok(())
}
fn maybe_auto_compact(&self) -> Result<bool> {
if self.merge_once()? {
return Ok(true);
}
if !self.config.auto_compaction {
return Ok(false);
}
let deleted = match self.shared.load_or_get_bitmap(false)? {
Some(bitmap) => bitmap.deleted_count.load(Ordering::Relaxed),
None => 0,
};
if deleted == 0 {
return Ok(false);
}
let total = self.shared.manager.total_vectors();
if total == 0 {
return Ok(false);
}
let ratio = deleted as f64 / total as f64;
if ratio >= self.config.compaction_threshold {
self.optimize()?;
return Ok(true);
}
Ok(false)
}
}
#[derive(Debug)]
struct SegmentedFlatWriter {
shared: Arc<SegmentedShared>,
segment_id: String,
inner: FlatIndexWriter,
sealed_len: Option<usize>,
}
impl VectorIndexWriter for SegmentedFlatWriter {
fn next_vector_id(&self) -> u64 {
self.inner.next_vector_id()
}
fn build(&mut self, vectors: Vec<(u64, String, Vector)>) -> Result<()> {
self.add_vectors(vectors)
}
fn add_vectors(&mut self, vectors: Vec<(u64, String, Vector)>) -> Result<()> {
if let Some(bitmap) = self.shared.load_or_get_bitmap(false)? {
for (doc_id, _, _) in &vectors {
let _ = bitmap.undelete_document(*doc_id)?;
}
}
self.inner.add_vectors(vectors)
}
fn finalize(&mut self) -> Result<()> {
self.inner.finalize()
}
fn progress(&self) -> f32 {
self.inner.progress()
}
fn estimated_memory_usage(&self) -> usize {
self.inner.estimated_memory_usage()
}
fn vectors(&self) -> &[(u64, String, Vector)] {
self.inner.vectors()
}
fn write(&self) -> Result<()> {
if self.inner.vectors().is_empty() {
return Ok(());
}
self.inner.write()
}
fn has_storage(&self) -> bool {
self.inner.has_storage()
}
fn delete_document(&mut self, doc_id: u64) -> Result<()> {
if let Some(bitmap) = self.shared.load_or_get_bitmap(true)? {
let newly_deleted = bitmap.delete_document(doc_id)?;
if newly_deleted && !self.shared.manager.list_segments().is_empty() {
self.shared.manager.mark_all_has_deletions()?;
}
}
let _ = self.inner.delete_document(doc_id);
Ok(())
}
fn has_pending_changes(&self) -> bool {
match self.sealed_len {
Some(sealed) => self.inner.vectors().len() != sealed,
None => !self.inner.vectors().is_empty(),
}
}
fn commit(&mut self) -> Result<()> {
if self.inner.vectors().is_empty() {
return Ok(());
}
if let Some(sealed) = self.sealed_len {
if self.inner.vectors().len() == sealed {
return Ok(());
}
return Err(LaurusError::InvalidOperation(format!(
"segment '{}' is already sealed; obtain a fresh writer for further changes",
self.segment_id
)));
}
self.inner.finalize()?;
self.inner.write()?;
let vector_count = self.inner.vectors().len() as u64;
let info = ManagedSegmentInfo::new(self.segment_id.clone(), vector_count, 0, 0);
self.shared.manager.add_segment(info)?;
self.sealed_len = Some(self.inner.vectors().len());
self.shared.reader_cache.invalidate(&self.segment_id);
Ok(())
}
fn rollback(&mut self) -> Result<()> {
self.shared
.pending_wal_seq
.store(self.shared.manager.last_wal_seq(), Ordering::Release);
self.inner.rollback()
}
fn pending_docs(&self) -> u64 {
self.inner.pending_docs()
}
fn close(&mut self) -> Result<()> {
if self.has_pending_changes() {
self.commit()?;
}
self.inner.close()
}
fn is_closed(&self) -> bool {
self.inner.is_closed()
}
fn build_reader(&self) -> Result<Arc<dyn VectorIndexReader>> {
self.inner.build_reader()
}
}
impl Drop for SegmentedFlatWriter {
fn drop(&mut self) {
if self.has_pending_changes() {
self.shared
.pending_wal_seq
.store(self.shared.manager.last_wal_seq(), Ordering::Release);
}
}
}