use std::path::PathBuf;
use std::sync::Arc;
use ahash::AHashMap;
use crate::common::counter::hardware_counter::HardwareCounterCell;
use crate::common::generic_consts::AccessPattern;
use crate::common::is_alive_lock::IsAliveLock;
use crate::common::types::PointOffsetType;
use crate::common::universal_io::MmapFile;
use parking_lot::RwLock;
use super::appendable_mmap_multi_dense_vector_storage::MultivectorMmapOffset;
use crate::segment::common::Flusher;
use crate::segment::common::operation_error::OperationResult;
use crate::segment::vector_storage::VectorOffsetType;
use crate::segment::vector_storage::chunked_vectors::ChunkedVectors;
type OffsetsStore = ChunkedVectors<MultivectorMmapOffset, MmapFile>;
#[derive(Debug)]
struct Inner {
store: OffsetsStore,
pending: AHashMap<VectorOffsetType, MultivectorMmapOffset>,
len: usize,
}
impl Inner {
fn get<P: AccessPattern>(&self, key: VectorOffsetType) -> Option<MultivectorMmapOffset> {
if let Some(entry) = self.pending.get(&key) {
return Some(*entry);
}
let &[entry] = self.store.get::<P>(key)?.as_ref() else {
unreachable!("multi-vector offsets are stored as vectors of length 1");
};
Some(entry)
}
}
#[derive(Debug)]
pub(super) struct BufferedOffsets {
inner: Arc<RwLock<Inner>>,
is_alive: IsAliveLock,
}
impl BufferedOffsets {
pub(super) fn new(store: OffsetsStore) -> Self {
let len = store.len();
Self {
inner: Arc::new(RwLock::new(Inner {
store,
pending: AHashMap::new(),
len,
})),
is_alive: IsAliveLock::new(),
}
}
pub(super) fn get<P: AccessPattern>(
&self,
key: VectorOffsetType,
) -> Option<MultivectorMmapOffset> {
self.inner.read().get::<P>(key)
}
pub(super) fn resolve_rows<P, U, I>(&self, keys: I) -> Vec<(U, PointOffsetType, u32)>
where
P: AccessPattern,
I: IntoIterator<Item = (U, PointOffsetType)>,
{
let inner = self.inner.read();
keys.into_iter()
.map(|(user_data, key)| {
let entry = inner
.get::<P>(key as VectorOffsetType)
.expect("offset entry exists");
(user_data, entry.offset, entry.count)
})
.collect()
}
pub(super) fn set(&self, key: VectorOffsetType, entry: MultivectorMmapOffset) {
let mut inner = self.inner.write();
inner.pending.insert(key, entry);
inner.len = inner.len.max(key + 1);
}
pub(super) fn len(&self) -> usize {
self.inner.read().len
}
pub(super) fn populate(&self) -> OperationResult<()> {
self.inner.read().store.populate()
}
pub(super) fn clear_cache(&self) -> OperationResult<()> {
self.inner.read().store.clear_cache()
}
pub(super) fn files(&self) -> Vec<PathBuf> {
self.inner.read().store.files()
}
pub(super) fn immutable_files(&self) -> Vec<PathBuf> {
self.inner.read().store.immutable_files()
}
pub(super) fn flusher(&self) -> Flusher {
let snapshot: AHashMap<VectorOffsetType, MultivectorMmapOffset> =
self.inner.read().pending.clone();
let inner = Arc::downgrade(&self.inner);
let is_alive = self.is_alive.handle();
Box::new(move || {
let (Some(_guard), Some(inner)) = (is_alive.lock_if_alive(), inner.upgrade()) else {
return Ok(());
};
let store_flusher = {
let mut items: Vec<(VectorOffsetType, MultivectorMmapOffset)> =
snapshot.iter().map(|(key, entry)| (*key, *entry)).collect();
items.sort_unstable_by_key(|(key, _)| *key);
let mut inner = inner.write();
let hw_counter = HardwareCounterCell::disposable();
for (key, entry) in &items {
inner.store.insert(*key, &[*entry], &hw_counter)?;
}
inner.pending.retain(|key, value| {
snapshot.get(key).is_none_or(|persisted| persisted != value)
});
inner.store.flusher()
};
store_flusher()
})
}
}
#[cfg(test)]
mod tests {
use crate::common::generic_consts::Sequential;
use crate::common::mmap::AdviceSetting;
use crate::common::universal_io::{MmapFs, Populate};
use tempfile::Builder;
use super::*;
fn offset(o: u32, c: u32) -> MultivectorMmapOffset {
MultivectorMmapOffset {
offset: o,
count: c,
capacity: c,
}
}
fn open_inner(dir: &std::path::Path) -> OffsetsStore {
ChunkedVectors::open(MmapFs, dir, 1, AdviceSetting::Global, Populate::No).unwrap()
}
#[test]
fn post_snapshot_write_does_not_leak_into_durable_store() {
let dir = Builder::new().prefix("buffered_offsets").tempdir().unwrap();
let store = BufferedOffsets::new(open_inner(dir.path()));
store.set(0, offset(0, 2));
store.set(1, offset(2, 3));
let flush = store.flusher();
store.set(1, offset(5, 5));
flush().unwrap();
drop(store);
let durable = open_inner(dir.path());
assert_eq!(durable.len(), 2, "both points are durable");
let &[e1] = durable.get::<Sequential>(1).unwrap().as_ref() else {
unreachable!()
};
assert_eq!(
e1,
offset(2, 3),
"durable entry for point 1 must be the pre-snapshot value, not the relocation",
);
let store = BufferedOffsets::new(durable);
assert_eq!(store.get::<Sequential>(1), Some(offset(2, 3)));
}
#[test]
fn deferred_write_commits_on_next_flush() {
let dir = Builder::new().prefix("buffered_offsets").tempdir().unwrap();
let store = BufferedOffsets::new(open_inner(dir.path()));
store.set(0, offset(0, 2));
let flush = store.flusher();
store.set(1, offset(2, 3)); flush().unwrap();
assert_eq!(store.get::<Sequential>(1), Some(offset(2, 3)));
assert_eq!(store.len(), 2);
store.flusher()().unwrap();
drop(store);
let durable = open_inner(dir.path());
assert_eq!(durable.len(), 2);
let &[e1] = durable.get::<Sequential>(1).unwrap().as_ref() else {
unreachable!()
};
assert_eq!(e1, offset(2, 3));
}
}