use crate::arena::CompactArenaSet;
use crate::error::{IncrementalFstError, Result};
use crate::inner::{FstDataSource, IncrementalFstSetInner, RebuildOutcome};
use crate::metrics::FstMetricsSnapshot;
use crate::stream::{self, MergedSetStreamOwner};
use crate::types::{FstChangeType, FstConfigOptions, FstMutationResult, SerializableFstBuffers};
use fst::{Automaton, IntoStreamer, Set, Streamer};
use memmap2::MmapOptions;
use parking_lot::RwLock;
use std::collections::BTreeSet;
use std::fs::File;
use std::io::Cursor;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
pub struct IncrementalFstSet {
inner: RwLock<IncrementalFstSetInner>,
}
impl IncrementalFstSet {
fn resolve_config_options_and_buffers(
options: Option<FstConfigOptions>,
) -> (usize, f32, Option<Duration>, CompactArenaSet) {
let opts = options.unwrap_or_default();
let adds_count = opts
.rebuild_threshold_adds_count
.unwrap_or(FstConfigOptions::DEFAULT_REBUILD_ADDS_COUNT);
let dels_ratio = opts
.rebuild_threshold_dels_ratio
.unwrap_or(FstConfigOptions::DEFAULT_REBUILD_DELS_RATIO);
let interval = opts.min_rebuild_interval_ms.map(Duration::from_millis);
let mutation_buffer = opts.initial_mutation_buffer.unwrap_or_default();
(adds_count, dels_ratio, interval, mutation_buffer)
}
pub fn new(options: Option<FstConfigOptions>) -> Result<Self> {
let (rebuild_adds, rebuild_dels_ratio, min_interval, initial_buf) =
Self::resolve_config_options_and_buffers(options);
let mut empty_data_writer = Vec::new();
let mut cursor = Cursor::new(&mut empty_data_writer);
let builder = fst::SetBuilder::new(&mut cursor)?;
builder.finish()?;
let empty_fst_data_source = FstDataSource::new_from_vec(Arc::new(empty_data_writer));
Ok(Self {
inner: RwLock::new(IncrementalFstSetInner::new(
empty_fst_data_source,
initial_buf,
true,
rebuild_adds,
rebuild_dels_ratio,
min_interval,
)?),
})
}
pub fn from_persisted_mmap(mmap_path: &Path, options: Option<FstConfigOptions>) -> Result<Self> {
let (rebuild_adds, rebuild_dels_ratio, min_interval, initial_buf) =
Self::resolve_config_options_and_buffers(options);
let fst_file = File::open(mmap_path).map_err(IncrementalFstError::Io)?;
let mmap = unsafe { MmapOptions::new().map(&fst_file).map_err(IncrementalFstError::Io)? };
let fst_data_source = FstDataSource::new_from_mmap(mmap);
Ok(Self {
inner: RwLock::new(IncrementalFstSetInner::new(
fst_data_source,
initial_buf,
false,
rebuild_adds,
rebuild_dels_ratio,
min_interval,
)?),
})
}
pub fn from_data(fst_bytes: Option<Vec<u8>>, options: Option<FstConfigOptions>) -> Result<Self> {
let (rebuild_adds, rebuild_dels_ratio, min_interval, initial_buf) =
Self::resolve_config_options_and_buffers(options);
let is_empty_init_flag;
let fst_data_source = match fst_bytes {
Some(bytes) => {
is_empty_init_flag = Set::new(&bytes)?.is_empty();
FstDataSource::new_from_vec(Arc::new(bytes))
}
None => {
let mut empty_data_writer = Vec::new();
let mut cursor = Cursor::new(&mut empty_data_writer);
let builder = fst::SetBuilder::new(&mut cursor)?;
builder.finish()?;
is_empty_init_flag = true;
FstDataSource::new_from_vec(Arc::new(empty_data_writer))
}
};
Ok(Self {
inner: RwLock::new(IncrementalFstSetInner::new(
fst_data_source,
initial_buf,
is_empty_init_flag,
rebuild_adds,
rebuild_dels_ratio,
min_interval,
)?),
})
}
pub fn get_metrics(&self) -> Result<FstMetricsSnapshot> {
Ok(self.inner.read().metrics.snapshot())
}
pub fn persisted_fst_as_bytes(&self) -> Vec<u8> {
self.inner.read().persisted_fst_data_arc.as_ref().to_vec()
}
pub fn buffers_snapshot(&self) -> SerializableFstBuffers {
let inner = self.inner.read();
SerializableFstBuffers::new(inner.mutation_buffer.clone())
}
pub fn insert(&self, key: Vec<u8>) -> Result<FstMutationResult> {
let mut inner = self.inner.write();
if inner.mutation_buffer.insert(&key) {
inner.add_buffer_changed_since_last_cache = true;
if !inner.rebuild_pending_due_to_interval {
let rebuild_outcome = IncrementalFstSetInner::check_and_trigger_rebuild(&mut inner, false)?;
match rebuild_outcome {
RebuildOutcome::Rebuilt => {
return Ok(FstMutationResult {
change_type: FstChangeType::FstRebuilt,
});
}
RebuildOutcome::RebuildDeferred => {
return Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
});
}
RebuildOutcome::NotNeeded => {
return Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
});
}
}
}
return Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
});
}
Ok(FstMutationResult {
change_type: FstChangeType::NoChange,
})
}
pub fn remove(&self, key: &[u8]) -> Result<FstMutationResult> {
let mut inner = self.inner.write();
if inner.mutation_buffer.remove(key) {
inner.add_buffer_changed_since_last_cache = true;
let rebuild_outcome = IncrementalFstSetInner::check_and_trigger_rebuild(&mut inner, false)?;
match rebuild_outcome {
RebuildOutcome::Rebuilt => {
return Ok(FstMutationResult {
change_type: FstChangeType::FstRebuilt,
});
}
RebuildOutcome::RebuildDeferred | RebuildOutcome::NotNeeded => {
return Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
});
}
}
}
Ok(FstMutationResult {
change_type: FstChangeType::NoChange,
})
}
pub fn bulk_insert<I>(&self, iter: I) -> Result<FstMutationResult>
where
I: IntoIterator<Item = Vec<u8>>,
{
let mut inner = self.inner.write();
let mut changed = false;
for key in iter {
if inner.mutation_buffer.insert(&key) {
changed = true;
}
}
if changed {
inner.add_buffer_changed_since_last_cache = true;
Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
})
} else {
Ok(FstMutationResult {
change_type: FstChangeType::NoChange,
})
}
}
pub fn finish_bulk_operations_and_rebuild_if_needed(&self) -> Result<FstMutationResult> {
let mut inner = self.inner.write();
let buffer_has_data = inner.mutation_buffer.len_raw() > 0;
let initial_pending_rebuild = inner.rebuild_pending_due_to_interval;
let should_attempt_rebuild =
inner.add_buffer_changed_since_last_cache || buffer_has_data || initial_pending_rebuild;
if should_attempt_rebuild {
let rebuild_outcome = IncrementalFstSetInner::check_and_trigger_rebuild(&mut inner, false)?;
match rebuild_outcome {
RebuildOutcome::Rebuilt => Ok(FstMutationResult {
change_type: FstChangeType::FstRebuilt,
}),
RebuildOutcome::RebuildDeferred => Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
}),
RebuildOutcome::NotNeeded => {
if inner.add_buffer_changed_since_last_cache || buffer_has_data {
Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
})
} else {
Ok(FstMutationResult {
change_type: FstChangeType::NoChange,
})
}
}
}
} else {
Ok(FstMutationResult {
change_type: FstChangeType::NoChange,
})
}
}
pub fn force_rebuild(&self) -> Result<FstMutationResult> {
let mut inner = self.inner.write();
let outcome = IncrementalFstSetInner::check_and_trigger_rebuild(&mut inner, true)?;
match outcome {
RebuildOutcome::Rebuilt => Ok(FstMutationResult {
change_type: FstChangeType::FstRebuilt,
}),
RebuildOutcome::RebuildDeferred => Ok(FstMutationResult {
change_type: FstChangeType::BuffersModified,
}),
RebuildOutcome::NotNeeded => Ok(FstMutationResult {
change_type: FstChangeType::NoChange,
}),
}
}
pub fn len(&self) -> Result<usize> {
let mut s = self.stream()?;
let mut count = 0;
while s.next().is_some() {
count += 1;
}
Ok(count)
}
pub fn contains(&self, key: &[u8]) -> Result<bool> {
let inner = self.inner.read();
match inner.mutation_buffer.get_status(key) {
Some(true) => Ok(true),
Some(false) => Ok(false),
None => {
let persisted_set = inner.get_persisted_set()?;
Ok(persisted_set.contains(key))
}
}
}
pub fn stream(&self) -> Result<MergedSetStreamOwner<'_>> {
let guard = self.inner.read();
stream::new_merged_set_stream_owner(guard)
}
pub fn is_empty(&self) -> Result<bool> {
let mut s = self.stream()?;
Ok(s.next().is_none())
}
pub fn search<A: Automaton + Clone>(&self, aut: A) -> Result<Vec<Vec<u8>>> {
let mut inner_guard = self.inner.write();
let cloned_overlay_fst_opt: Option<Set<FstDataSource>> = inner_guard.get_mutation_overlay_fst()?.cloned();
let persisted_s_ref = inner_guard.get_persisted_set()?;
let mut results_set = BTreeSet::new();
let mut persisted_stream = persisted_s_ref.search(aut.clone()).into_stream();
while let Some(key_slice) = persisted_stream.next() {
if inner_guard.mutation_buffer.get_status(key_slice) != Some(false) {
results_set.insert(key_slice.to_vec());
}
}
if let Some(overlay_fst) = &cloned_overlay_fst_opt {
let mut overlay_stream = overlay_fst.search(aut).into_stream();
while let Some(key_slice) = overlay_stream.next() {
results_set.insert(key_slice.to_vec());
}
}
Ok(results_set.into_iter().collect())
}
}
#[cfg(test)]
mod tests {
use super::*;
use fst::automaton::Str;
use fst::set::OpBuilder as FstSetOpBuilder;
use std::sync::atomic::Ordering as AtomicOrdering;
use std::thread;
#[cfg(feature = "serde")]
use crate::arena::CompactArenaSet;
#[cfg(feature = "serde")]
use std::fs;
#[cfg(feature = "serde")]
fn temp_dir() -> tempfile::TempDir {
tempfile::Builder::new().prefix("inc_fst_test_").tempdir().unwrap()
}
fn new_test_set_with_options(options: Option<FstConfigOptions>) -> Result<IncrementalFstSet> {
IncrementalFstSet::new(options)
}
fn new_test_set(adds: usize, dels_ratio: f32, interval_ms: Option<u64>) -> Result<IncrementalFstSet> {
let mut opts_builder = FstConfigOptions::new()
.with_rebuild_adds_count(adds)
.with_rebuild_dels_ratio(dels_ratio);
if let Some(ms) = interval_ms {
opts_builder = opts_builder.with_min_rebuild_interval_ms(ms);
}
IncrementalFstSet::new(Some(opts_builder))
}
#[test]
fn test_insert_and_remove_outcomes() -> Result<()> {
let config = FstConfigOptions::new()
.with_rebuild_adds_count(3)
.with_rebuild_dels_ratio(1.1);
let set = new_test_set_with_options(Some(config))?;
let res1 = set.insert(b"apple".to_vec())?;
assert_eq!(res1.change_type, FstChangeType::BuffersModified);
let res2 = set.insert(b"banana".to_vec())?;
assert_eq!(res2.change_type, FstChangeType::BuffersModified);
let res3 = set.insert(b"orange".to_vec())?;
assert_eq!(res3.change_type, FstChangeType::FstRebuilt);
let buffers_after_rebuild = set.buffers_snapshot();
assert!(buffers_after_rebuild.mutation_buffer().is_empty());
let res4 = set.insert(b"apple".to_vec())?;
assert_eq!(res4.change_type, FstChangeType::BuffersModified);
assert!(set.buffers_snapshot().mutation_buffer().contains(b"apple".as_ref()));
let res5 = set.remove(b"apple")?;
assert_eq!(res5.change_type, FstChangeType::BuffersModified);
assert!(!set.buffers_snapshot().mutation_buffer().contains(b"apple".as_ref()));
assert!(set.buffers_snapshot().mutation_buffer().len_raw() > 0);
let res6 = set.remove(b"banana")?;
assert_eq!(res6.change_type, FstChangeType::BuffersModified);
assert!(!set.buffers_snapshot().mutation_buffer().contains(b"banana".as_ref()));
let res7 = set.remove(b"grape")?;
assert_eq!(res7.change_type, FstChangeType::BuffersModified);
let res_force = set.force_rebuild()?;
assert_eq!(res_force.change_type, FstChangeType::FstRebuilt);
assert!(!set.contains(b"banana")?);
assert!(set.contains(b"orange")?);
assert!(set.buffers_snapshot().mutation_buffer().is_empty());
Ok(())
}
#[test]
fn test_fst_config_options() -> Result<()> {
let set_default = IncrementalFstSet::new(None)?;
let inner_default = set_default.inner.read();
assert_eq!(
inner_default.rebuild_threshold_adds_count,
FstConfigOptions::DEFAULT_REBUILD_ADDS_COUNT
);
assert_eq!(
inner_default.rebuild_threshold_dels_ratio,
FstConfigOptions::DEFAULT_REBUILD_DELS_RATIO
);
assert_eq!(inner_default.min_rebuild_interval, None);
drop(inner_default);
let custom_adds = 555;
let custom_ratio = 0.77f32;
let custom_interval_ms = 12345u64;
let options = FstConfigOptions::new()
.with_rebuild_adds_count(custom_adds)
.with_rebuild_dels_ratio(custom_ratio)
.with_min_rebuild_interval_ms(custom_interval_ms);
let set_custom = IncrementalFstSet::new(Some(options))?;
let inner_custom = set_custom.inner.read();
assert_eq!(inner_custom.rebuild_threshold_adds_count, custom_adds);
assert_eq!(inner_custom.rebuild_threshold_dels_ratio, custom_ratio);
assert_eq!(
inner_custom.min_rebuild_interval,
Some(Duration::from_millis(custom_interval_ms))
);
drop(inner_custom);
Ok(())
}
#[test]
#[cfg(feature = "serde")]
fn test_caller_persistence_from_data() -> Result<()> {
let dir = temp_dir();
let fst_file_path_v1 = dir.path().join("data_v1.fst");
let buffers_file_path_v1 = dir.path().join("buffers_v1.rmp");
let fst_file_path_v2 = dir.path().join("data_v2.fst");
let buffers_file_path_v2 = dir.path().join("buffers_v2.rmp");
let opts_for_set = FstConfigOptions::new().with_rebuild_adds_count(2);
let set = IncrementalFstSet::new(Some(opts_for_set.clone()))?;
let res_insert1 = set.insert(b"apple".to_vec())?;
assert_eq!(res_insert1.change_type, FstChangeType::BuffersModified);
let fst_bytes_v1 = set.persisted_fst_as_bytes();
fs::write(&fst_file_path_v1, &fst_bytes_v1).unwrap();
let buffers_v1_snapshot = set.buffers_snapshot();
let serialized_buffers_v1 = rmp_serde::to_vec_named(&buffers_v1_snapshot).unwrap();
fs::write(&buffers_file_path_v1, &serialized_buffers_v1).unwrap();
let res_insert2 = set.insert(b"banana".to_vec())?; assert_eq!(res_insert2.change_type, FstChangeType::FstRebuilt);
let fst_bytes_v2 = set.persisted_fst_as_bytes();
fs::write(&fst_file_path_v2, &fst_bytes_v2).unwrap();
let buffers_v2_snapshot = set.buffers_snapshot();
let serialized_buffers_v2 = rmp_serde::to_vec_named(&buffers_v2_snapshot).unwrap();
fs::write(&buffers_file_path_v2, &serialized_buffers_v2).unwrap();
let loaded_fst_bytes_v1_from_file = fs::read(&fst_file_path_v1).unwrap();
let loaded_serialized_buffers_v1_from_file = fs::read(&buffers_file_path_v1).unwrap();
let loaded_buffers_struct_v1: SerializableFstBuffers =
rmp_serde::from_slice(&loaded_serialized_buffers_v1_from_file).unwrap();
let options_for_load_v1 = FstConfigOptions::new()
.with_rebuild_adds_count(2)
.with_initial_buffer(loaded_buffers_struct_v1.into_mutation_buffer());
let set_v1_loaded = IncrementalFstSet::from_data(Some(loaded_fst_bytes_v1_from_file), Some(options_for_load_v1))?;
assert!(set_v1_loaded.contains(b"apple")?);
assert!(!set_v1_loaded.contains(b"banana")?);
assert_eq!(set_v1_loaded.buffers_snapshot().mutation_buffer().len(), 1);
assert!(
set_v1_loaded
.buffers_snapshot()
.mutation_buffer()
.contains(b"apple".as_ref())
);
let loaded_fst_bytes_v2_from_file = fs::read(&fst_file_path_v2).unwrap();
let loaded_serialized_buffers_v2_from_file = fs::read(&buffers_file_path_v2).unwrap();
let loaded_buffers_struct_v2: SerializableFstBuffers =
rmp_serde::from_slice(&loaded_serialized_buffers_v2_from_file).unwrap();
let options_for_load_v2 = FstConfigOptions::new()
.with_rebuild_adds_count(2)
.with_initial_buffer(loaded_buffers_struct_v2.into_mutation_buffer());
let set_v2_loaded = IncrementalFstSet::from_data(Some(loaded_fst_bytes_v2_from_file), Some(options_for_load_v2))?;
assert!(set_v2_loaded.contains(b"apple")?);
assert!(set_v2_loaded.contains(b"banana")?);
assert!(set_v2_loaded.buffers_snapshot().mutation_buffer().is_empty());
assert_eq!(set_v2_loaded.persisted_fst_as_bytes(), fst_bytes_v2.as_slice());
let mutation_buf_only: CompactArenaSet = vec![b"only_add"].into_iter().collect();
let opts_buf_only = FstConfigOptions::new().with_initial_buffer(mutation_buf_only.clone());
let set_buf_only = IncrementalFstSet::from_data(None, Some(opts_buf_only))?;
assert!(
set_buf_only.persisted_fst_as_bytes().is_empty() || Set::new(set_buf_only.persisted_fst_as_bytes())?.is_empty()
);
assert_eq!(set_buf_only.buffers_snapshot().mutation_buffer(), &mutation_buf_only);
assert!(set_buf_only.contains(b"only_add")?);
Ok(())
}
#[test]
#[cfg(feature = "serde")]
fn test_caller_persistence_from_mmap() -> Result<()> {
let dir = temp_dir();
let fst_file_path = dir.path().join("mmap_test.fst");
let buffers_file_path = dir.path().join("mmap_buffers.rmp");
let initial_opts = FstConfigOptions::new()
.with_rebuild_adds_count(5)
.with_rebuild_dels_ratio(1.0);
let setup_set = IncrementalFstSet::new(Some(initial_opts.clone()))?;
setup_set.insert(b"mmap_apple".to_vec())?;
setup_set.insert(b"mmap_banana".to_vec())?;
setup_set.force_rebuild()?;
setup_set.insert(b"mmap_cherry".to_vec())?;
setup_set.remove(b"mmap_apple".as_ref())?;
let fst_bytes_to_write = setup_set.persisted_fst_as_bytes();
fs::write(&fst_file_path, &fst_bytes_to_write).unwrap();
let buffers_to_write = setup_set.buffers_snapshot();
let serialized_buffers = rmp_serde::to_vec_named(&buffers_to_write).unwrap();
fs::write(&buffers_file_path, &serialized_buffers).unwrap();
let loaded_serialized_buffers = fs::read(&buffers_file_path).unwrap();
let loaded_buffers_struct: SerializableFstBuffers = rmp_serde::from_slice(&loaded_serialized_buffers).unwrap();
let load_options = FstConfigOptions::new()
.with_rebuild_adds_count(5)
.with_initial_buffer(loaded_buffers_struct.into_mutation_buffer());
let loaded_set = IncrementalFstSet::from_persisted_mmap(&fst_file_path, Some(load_options))?;
assert!(!loaded_set.contains(b"mmap_apple")?, "mmap_apple should be deleted");
assert!(
loaded_set.contains(b"mmap_banana")?,
"mmap_banana should be in persisted FST"
);
assert!(
loaded_set.contains(b"mmap_cherry")?,
"mmap_cherry should be in loaded buffer"
);
let loaded_buffer = loaded_set.buffers_snapshot().into_mutation_buffer();
assert!(loaded_buffer.contains(b"mmap_cherry".as_ref()));
assert!(!loaded_buffer.contains(b"mmap_apple".as_ref()));
assert_eq!(loaded_buffer.len_raw(), 2);
let mmap_fst_bytes = loaded_set.persisted_fst_as_bytes();
assert_eq!(mmap_fst_bytes, fst_bytes_to_write.as_slice());
Ok(())
}
#[test]
fn test_getters() -> Result<()> {
let set = new_test_set(2, 1.0, None)?;
assert!(set.persisted_fst_as_bytes().is_empty() || Set::new(set.persisted_fst_as_bytes())?.is_empty());
let snap0 = set.buffers_snapshot();
assert!(snap0.mutation_buffer().is_empty());
set.insert(b"key1".to_vec())?;
assert!(set.persisted_fst_as_bytes().is_empty() || Set::new(set.persisted_fst_as_bytes())?.is_empty());
let snap1 = set.buffers_snapshot();
assert!(snap1.mutation_buffer().contains(b"key1".as_ref()));
assert_eq!(snap1.mutation_buffer().len(), 1);
set.insert(b"key2".to_vec())?;
let fst_bytes_after_rebuild = set.persisted_fst_as_bytes();
let rebuilt_fst = Set::new(fst_bytes_after_rebuild)?;
assert!(rebuilt_fst.contains(b"key1"));
assert!(rebuilt_fst.contains(b"key2"));
let snap2 = set.buffers_snapshot();
assert!(snap2.mutation_buffer().is_empty());
set.remove(b"key1")?;
let snap3 = set.buffers_snapshot();
assert!(!snap3.mutation_buffer().contains(b"key1".as_ref()));
assert_eq!(snap3.mutation_buffer().len_raw(), 1);
Ok(())
}
#[test]
fn test_true_streaming_merged_set_stream() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
let _ = ufst.insert(b"b".to_vec())?;
let _ = ufst.insert(b"d".to_vec())?;
let _ = ufst.force_rebuild()?;
let _ = ufst.insert(b"a".to_vec())?;
let _ = ufst.insert(b"c".to_vec())?;
let _ = ufst.insert(b"d".to_vec())?;
let _ = ufst.remove(b"b")?;
let mut stream_owner = ufst.stream()?;
assert_eq!(stream_owner.next(), Some(b"a".to_vec()));
assert_eq!(stream_owner.next(), Some(b"c".to_vec()));
assert_eq!(stream_owner.next(), Some(b"d".to_vec()));
assert_eq!(stream_owner.next(), None);
Ok(())
}
#[test]
fn test_thread_safety_and_metrics() -> Result<()> {
let ufst_arc = Arc::new(new_test_set(2, 0.3, Some(10))?);
let mut handles = vec![];
for i in 0..20 {
let ufst_clone = ufst_arc.clone();
handles.push(thread::spawn(move || {
let key_str = format!("key{:02}", i);
let _ = ufst_clone.insert(key_str.into_bytes()).unwrap();
if i % 3 == 0 && i > 0 {
let key_to_remove = format!("key{:02}", i - 1);
let _ = ufst_clone.remove(key_to_remove.as_bytes());
}
if i % 5 == 0 {
let _ = ufst_clone.search(Str::new("key").starts_with()).unwrap();
}
let _ = ufst_clone.contains(format!("key{:02}", i).as_bytes()).unwrap_or(false);
}));
}
for _i in 0..5 {
let ufst_clone = ufst_arc.clone();
handles.push(thread::spawn(move || {
thread::sleep(Duration::from_millis(5));
if let Ok(mut stream_owner) = ufst_clone.stream() {
while stream_owner.next().is_some() { }
}
}));
}
for handle in handles {
handle.join().expect("Thread panicked");
}
let _ = ufst_arc.force_rebuild()?;
let metrics = ufst_arc.get_metrics()?;
assert!(metrics.num_rebuilds > 0);
assert!(ufst_arc.len()? > 0);
Ok(())
}
#[test]
fn test_op_facade_via_snapshot() -> Result<()> {
let ufst1 = new_test_set(5, 0.5, None)?;
let _ = ufst1.insert(b"a".to_vec())?;
let _ = ufst1.insert(b"b".to_vec())?;
let _ = ufst1.insert(b"c".to_vec())?;
let force_rebuild_result = ufst1.force_rebuild()?;
assert_eq!(force_rebuild_result.change_type, FstChangeType::FstRebuilt);
let persisted_fst1_bytes = ufst1.persisted_fst_as_bytes();
let set_from_ufst1 = Set::new(persisted_fst1_bytes)?;
let fst2_data_vec: Vec<u8> = {
let mut write_buf = Vec::new();
let mut cursor = Cursor::new(&mut write_buf);
let mut builder = fst::SetBuilder::new(&mut cursor)?;
builder.insert(b"b")?;
builder.insert(b"c")?;
builder.insert(b"d")?;
builder.finish()?;
write_buf
};
let fst2 = Set::new(fst2_data_vec)?;
let mut op = FstSetOpBuilder::new();
op.push(set_from_ufst1.stream());
op.push(fst2.stream());
let mut union_stream = op.union();
let mut results = Vec::new();
while let Some(key_slice) = union_stream.next() {
results.push(key_slice.to_vec());
}
results.sort();
assert_eq!(
results,
vec![b"a".to_vec(), b"b".to_vec(), b"c".to_vec(), b"d".to_vec()]
);
Ok(())
}
#[test]
fn test_bulk_insert_and_finish() -> Result<()> {
let ufst = new_test_set(3, 0.1, None)?; let keys_to_bulk_insert: Vec<Vec<u8>> = (0..5).map(|i| format!("bulk{}", i).into_bytes()).collect();
let bulk_res = ufst.bulk_insert(keys_to_bulk_insert.clone())?;
assert_eq!(bulk_res.change_type, FstChangeType::BuffersModified);
let inner_read = ufst.inner.read();
assert_eq!(inner_read.mutation_buffer.len(), 5);
let num_rebuilds_before = inner_read.metrics.num_rebuilds.load(AtomicOrdering::Relaxed);
assert!(inner_read.add_buffer_changed_since_last_cache);
drop(inner_read);
let finish_res = ufst.finish_bulk_operations_and_rebuild_if_needed()?;
assert_eq!(finish_res.change_type, FstChangeType::FstRebuilt);
let inner_read_after = ufst.inner.read();
assert_eq!(
inner_read_after.metrics.num_rebuilds.load(AtomicOrdering::Relaxed),
num_rebuilds_before + 1
);
assert!(inner_read_after.mutation_buffer.is_empty());
assert!(!inner_read_after.add_buffer_changed_since_last_cache);
drop(inner_read_after);
for key_vec in keys_to_bulk_insert {
assert!(ufst.contains(&key_vec)?);
}
Ok(())
}
fn assert_stream_contents(stream: &mut MergedSetStreamOwner<'_>, expected_slices: Vec<&[u8]>) {
let mut actual_items = Vec::new();
while let Some(item) = stream.next() {
actual_items.push(item);
}
let expected_owned_items: Vec<Vec<u8>> = expected_slices.iter().map(|s| s.to_vec()).collect();
assert_eq!(actual_items, expected_owned_items, "Stream contents mismatch");
}
#[test]
fn test_stream_empty_set_direct() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
let mut stream_owner = ufst.stream()?;
assert_stream_contents(&mut stream_owner, vec![]);
assert!(ufst.is_empty()?);
assert_eq!(ufst.len()?, 0);
Ok(())
}
#[test]
fn test_stream_persisted_only_simple() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
ufst.insert(b"apple".to_vec())?;
ufst.insert(b"banana".to_vec())?;
ufst.force_rebuild()?;
let mut stream_owner = ufst.stream()?;
assert_stream_contents(&mut stream_owner, vec![b"apple", b"banana"]);
Ok(())
}
#[test]
fn test_stream_add_buffer_only_out_of_order() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?; ufst.insert(b"zebra".to_vec())?;
ufst.insert(b"apple".to_vec())?;
ufst.insert(b"monkey".to_vec())?;
let mut stream_owner = ufst.stream()?;
assert_stream_contents(&mut stream_owner, vec![b"apple", b"monkey", b"zebra"]);
Ok(())
}
#[test]
fn test_stream_persisted_and_add_buffer_interleaved() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
ufst.insert(b"b_persisted".to_vec())?;
ufst.insert(b"d_persisted".to_vec())?;
ufst.force_rebuild()?;
ufst.insert(b"a_added".to_vec())?;
ufst.insert(b"c_added".to_vec())?;
ufst.insert(b"e_added".to_vec())?;
let mut stream_owner = ufst.stream()?;
assert_stream_contents(
&mut stream_owner,
vec![b"a_added", b"b_persisted", b"c_added", b"d_persisted", b"e_added"],
);
Ok(())
}
#[test]
fn test_stream_with_deletes_affecting_persisted() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
ufst.insert(b"item1".to_vec())?;
ufst.insert(b"item2_delete_me".to_vec())?;
ufst.insert(b"item3".to_vec())?;
ufst.force_rebuild()?;
ufst.insert(b"item4_added".to_vec())?;
ufst.remove(b"item2_delete_me")?;
let mut stream_owner = ufst.stream()?;
assert_stream_contents(&mut stream_owner, vec![b"item1", b"item3", b"item4_added"]);
Ok(())
}
#[test]
fn test_stream_add_overwrites_persisted_key_and_delete_original() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
ufst.insert(b"key1_persisted".to_vec())?;
ufst.force_rebuild()?;
ufst.insert(b"key1_persisted".to_vec())?;
ufst.insert(b"key2_added".to_vec())?;
ufst.remove(b"key1_persisted")?;
let mut stream_owner = ufst.stream()?;
assert_stream_contents(&mut stream_owner, vec![b"key2_added"]);
Ok(())
}
#[test]
fn test_delete_item_from_persisted_when_also_in_add_buffer_no_rebuild() -> Result<()> {
let ufst = new_test_set(100, 0.5, None)?;
ufst.insert(b"key1".to_vec())?;
ufst.force_rebuild()?;
ufst.insert(b"key1".to_vec())?;
ufst.insert(b"key2_added".to_vec())?;
ufst.remove(b"key1")?;
let mut stream = ufst.stream()?;
assert_stream_contents(&mut stream, vec![b"key2_added"]);
Ok(())
}
#[test]
fn test_reinsert_deleted_item_before_rebuild_typical() -> Result<()> {
let ufst = new_test_set(100, 0.5, None)?;
ufst.insert(b"a".to_vec())?;
ufst.insert(b"b_to_delete_and_readd".to_vec())?;
ufst.insert(b"c".to_vec())?;
ufst.force_rebuild()?;
ufst.remove(b"b_to_delete_and_readd")?;
ufst.insert(b"b_to_delete_and_readd".to_vec())?;
let mut stream = ufst.stream()?;
assert_stream_contents(&mut stream, vec![b"a", b"b_to_delete_and_readd", b"c"]);
Ok(())
}
#[test]
fn test_stream_with_empty_string_key() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
ufst.insert(b"".to_vec())?;
ufst.insert(b"alpha".to_vec())?;
ufst.force_rebuild()?;
ufst.insert(b"beta".to_vec())?;
ufst.insert(b"".to_vec())?;
ufst.remove(b"alpha")?;
let mut stream = ufst.stream()?;
assert_stream_contents(&mut stream, vec![b"", b"beta"]);
Ok(())
}
#[test]
fn test_multiple_independent_streams_on_same_state() -> Result<()> {
let ufst = new_test_set(10, 0.1, None)?;
ufst.insert(b"persisted1".to_vec())?;
ufst.force_rebuild()?;
ufst.insert(b"added1".to_vec())?;
ufst.remove(b"persisted1")?;
ufst.insert(b"added2".to_vec())?;
let expected_vec = vec![b"added1".to_vec(), b"added2".to_vec()];
let mut stream1 = ufst.stream()?;
let mut items1 = Vec::new();
while let Some(item) = stream1.next() {
items1.push(item);
}
assert_eq!(items1, expected_vec);
assert!(stream1.next().is_none(), "Stream1 should be exhausted");
let mut stream2 = ufst.stream()?;
let mut items2 = Vec::new();
while let Some(item) = stream2.next() {
items2.push(item);
}
assert_eq!(items2, expected_vec);
assert!(stream2.next().is_none(), "Stream2 should be exhausted");
Ok(())
}
}