use std::path::{Path, PathBuf};
use crate::common::bitvec::{BitSlice, BitVec};
use crate::common::mmap::AdviceSetting;
use crate::common::sorted_slice::SortedSlice;
use crate::common::stored_bitslice::StoredBitSlice;
use crate::common::types::PointOffsetType;
use crate::common::universal_io::{
CachedReadFs, OpenOptions, Populate, TypedStorage, UniversalRead, UniversalReadFs,
};
use super::dynamic_stored_flags::{DynamicFlagsStatus, FLAGS_FILE, status_file};
use crate::segment::common::operation_error::{OperationError, OperationResult};
#[derive(Debug)]
pub struct InMemoryBitvecFlags {
bitvec: BitVec,
count: usize,
directory: Option<PathBuf>,
}
fn bitslice_open_options(populate: Populate) -> OpenOptions {
OpenOptions {
writeable: false,
need_sequential: false,
populate,
advice: AdviceSetting::Global,
}
}
impl InMemoryBitvecFlags {
pub fn preopen(fs: &impl CachedReadFs, directory: &Path) -> OperationResult<()> {
fs.schedule_prefetch(
&status_file(directory),
Some(bitslice_open_options(Populate::PreferBackground)),
None,
)?;
fs.schedule_prefetch(
&directory.join(FLAGS_FILE),
Some(bitslice_open_options(Populate::PreferBackground)),
None,
)?;
Ok(())
}
fn persisted_len<S: UniversalRead>(
fs: &impl UniversalReadFs<File = S>,
directory: &Path,
) -> OperationResult<usize> {
let status = TypedStorage::<S, DynamicFlagsStatus>::new(fs.open(
status_file(directory),
bitslice_open_options(Populate::No),
Default::default(),
)?);
Ok(status
.read_whole()?
.first()
.map_or(0, DynamicFlagsStatus::len))
}
pub fn open<S: UniversalRead>(
fs: &impl UniversalReadFs<File = S>,
directory: &Path,
) -> OperationResult<Self> {
let len = Self::persisted_len(fs, directory)?;
let flags_path = directory.join(FLAGS_FILE);
let flags = StoredBitSlice::<S>::open(
fs,
&flags_path,
bitslice_open_options(Populate::No),
Default::default(),
)?;
let bits = flags.read_all()?;
let bitvec = bits.get(..len).map(BitVec::from_bitslice).ok_or_else(|| {
OperationError::service_error(format!(
"Flags file {} holds fewer than {len} bits",
flags_path.display(),
))
})?;
let count = bitvec.count_ones();
Ok(Self {
bitvec,
count,
directory: Some(directory.to_path_buf()),
})
}
pub fn from_bitvec(bitvec: BitVec) -> Self {
let count = bitvec.count_ones();
Self {
bitvec,
count,
directory: None,
}
}
pub fn get(&self, key: PointOffsetType) -> bool {
self.bitvec.get(key as usize).is_some_and(|bit| *bit)
}
pub fn count(&self) -> usize {
self.count
}
pub fn as_bitslice(&self) -> &BitSlice {
self.bitvec.as_bitslice()
}
pub fn insert_all(&mut self, points: &[PointOffsetType]) {
for &point in points {
let index = point as usize;
if index >= self.bitvec.len() {
self.bitvec.resize(index + 1, false);
}
if !self.bitvec.replace(index, true) {
self.count += 1;
}
}
}
pub fn reload_appended<S: UniversalRead>(
&mut self,
fs: &impl UniversalReadFs<File = S>,
new_points: &SortedSlice<'_, PointOffsetType>,
) -> OperationResult<()> {
let Some(directory) = self.directory.clone() else {
return Ok(());
};
let (Some(&first), Some(&last)) = (new_points.first(), new_points.last()) else {
return Ok(());
};
let len = Self::persisted_len(fs, &directory)? as u64;
let start = u64::from(first);
let end = u64::from(last).saturating_add(1).min(len);
if start >= end {
return Ok(());
}
let flags = StoredBitSlice::<S>::open(
fs,
&directory.join(FLAGS_FILE),
bitslice_open_options(Populate::No),
Default::default(),
)?;
let bits = flags.read_bit_range(start..end)?;
let mut deleted = Vec::new();
for &point in new_points.iter() {
let index = (u64::from(point) - start) as usize;
if bits.get(index).as_deref().copied().unwrap_or(false) {
deleted.push(point);
}
}
self.insert_all(&deleted);
Ok(())
}
}
#[allow(clippy::default_constructed_unit_structs)]
#[duplicate::duplicate_item(
tests_mod S Fs cfg_predicate;
[tests_mmap] [MmapFile] [MmapFs] [cfg(all())];
[tests_uring] [IoUringFile] [IoUringFs] [cfg(target_os = "linux")];
)]
#[cfg_predicate]
#[cfg(test)]
mod tests_mod {
use std::iter;
#[cfg_predicate]
use crate::common::universal_io::{Fs, S};
use rand::prelude::StdRng;
use rand::{RngExt, SeedableRng};
use tempfile::Builder;
use super::*;
use crate::segment::common::flags::dynamic_stored_flags::DynamicStoredFlags;
fn persist(fs: &Fs, dir: &Path, flags: &[bool]) {
let mut dynamic_flags = DynamicStoredFlags::<S>::open(fs, dir, Populate::No).unwrap();
dynamic_flags.set_len(fs, flags.len()).unwrap();
flags
.iter()
.enumerate()
.filter(|(_, flag)| **flag)
.for_each(|(i, _)| assert!(!dynamic_flags.set(i, true).unwrap()));
dynamic_flags.flusher()().unwrap();
}
#[test]
fn open_materializes_persisted_flags() {
let dir = Builder::new().prefix("storage_dir").tempdir().unwrap();
let num_flags = 5003; let mut rng = StdRng::seed_from_u64(42);
let random_flags: Vec<bool> = iter::repeat_with(|| rng.random()).take(num_flags).collect();
persist(&Fs::default(), dir.path(), &random_flags);
let flags = InMemoryBitvecFlags::open::<S>(&Fs::default(), dir.path()).unwrap();
let expected_count = random_flags.iter().filter(|flag| **flag).count();
assert_eq!(flags.count(), expected_count);
for (i, &flag) in random_flags.iter().enumerate() {
assert_eq!(flags.get(i as PointOffsetType), flag);
}
}
#[test]
fn insert_all_folds_deletion_delta() {
let dir = Builder::new().prefix("storage_dir").tempdir().unwrap();
let num_flags = 1000;
let mut rng = StdRng::seed_from_u64(7);
let random_flags: Vec<bool> = iter::repeat_with(|| rng.random()).take(num_flags).collect();
persist(&Fs::default(), dir.path(), &random_flags);
let mut flags = InMemoryBitvecFlags::open::<S>(&Fs::default(), dir.path()).unwrap();
let base_count = flags.count();
let already_set = random_flags.iter().position(|flag| *flag).unwrap();
let unset = random_flags.iter().position(|flag| !*flag).unwrap();
let beyond = num_flags as PointOffsetType + 5;
flags.insert_all(&[
already_set as PointOffsetType,
unset as PointOffsetType,
beyond,
]);
assert_eq!(flags.count(), base_count + 2);
assert!(flags.get(already_set as PointOffsetType));
assert!(flags.get(unset as PointOffsetType));
assert!(flags.get(beyond));
assert!(!flags.get(beyond + 1));
}
}