use crate::IncrementalFstError;
use crate::arena::CompactArenaSet;
use crate::error::Result;
use crate::metrics::FstMetrics;
use crate::types::FstDataHolder;
use fst::{Set, SetBuilder, Streamer};
use memmap2::Mmap;
use once_cell::sync::OnceCell;
use std::io::{BufWriter, Cursor, Write};
use std::sync::Arc;
use std::sync::atomic::Ordering as AtomicOrdering;
use std::time::Duration as StdDuration;
use std::time::Instant;
#[derive(Clone)]
pub(crate) struct FstDataSource(Arc<FstDataHolder>);
impl FstDataSource {
pub(crate) fn new_from_vec(data: Arc<Vec<u8>>) -> Self {
Self(Arc::new(FstDataHolder::Memory(data)))
}
pub(crate) fn new_from_mmap(mmap: Mmap) -> Self {
Self(Arc::new(FstDataHolder::Mapped(Arc::new(mmap))))
}
}
impl AsRef<[u8]> for FstDataSource {
fn as_ref(&self) -> &[u8] {
self.0.as_ref().as_ref()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RebuildOutcome {
Rebuilt,
RebuildDeferred,
NotNeeded,
}
pub(crate) struct IncrementalFstSetInner {
pub(crate) persisted_fst_data_arc: FstDataSource,
pub(crate) mutation_buffer: CompactArenaSet,
pub(crate) rebuild_threshold_adds_count: usize,
pub(crate) rebuild_threshold_dels_ratio: f32,
pub(crate) cached_persisted_set_cell: OnceCell<Set<FstDataSource>>,
pub(crate) cached_add_buffer_fst_cell: OnceCell<Set<FstDataSource>>,
pub(crate) add_buffer_changed_since_last_cache: bool,
pub(crate) metrics: FstMetrics,
pub(crate) last_rebuild_time: Option<Instant>,
pub(crate) min_rebuild_interval: Option<StdDuration>,
pub(crate) rebuild_pending_due_to_interval: bool,
}
impl IncrementalFstSetInner {
pub(crate) fn new(
data_source: FstDataSource,
initial_buffer: CompactArenaSet,
is_empty_init: bool,
rebuild_threshold_adds_count: usize,
rebuild_threshold_dels_ratio: f32,
min_rebuild_interval: Option<StdDuration>,
) -> Result<Self> {
let persisted_cell = OnceCell::new();
if is_empty_init {
match Set::new(data_source.clone()) {
Ok(set_instance) => {
let _ = persisted_cell.set(set_instance);
}
Err(e) => {
return Err(IncrementalFstError::Fst(e));
}
}
}
let buffer_not_empty = !initial_buffer.is_empty();
Ok(Self {
persisted_fst_data_arc: data_source,
mutation_buffer: initial_buffer,
rebuild_threshold_adds_count,
rebuild_threshold_dels_ratio,
cached_persisted_set_cell: persisted_cell,
cached_add_buffer_fst_cell: OnceCell::new(),
add_buffer_changed_since_last_cache: buffer_not_empty,
metrics: FstMetrics::default(),
last_rebuild_time: None,
min_rebuild_interval,
rebuild_pending_due_to_interval: false,
})
}
pub(crate) fn get_persisted_set(&self) -> Result<&Set<FstDataSource>> {
self
.cached_persisted_set_cell
.get_or_try_init(|| Set::new(self.persisted_fst_data_arc.clone()).map_err(IncrementalFstError::from))
}
pub(crate) fn get_mutation_overlay_fst(&mut self) -> Result<Option<&Set<FstDataSource>>> {
if self.mutation_buffer.is_empty() {
if self.cached_add_buffer_fst_cell.get().is_some() {
self.cached_add_buffer_fst_cell = OnceCell::new();
}
self.add_buffer_changed_since_last_cache = false;
return Ok(None);
}
if self.add_buffer_changed_since_last_cache || self.cached_add_buffer_fst_cell.get().is_none() {
self.cached_add_buffer_fst_cell = OnceCell::new();
}
let result_opt = self.cached_add_buffer_fst_cell.get_or_try_init(|| {
let mut add_fst_data_vec = Vec::new();
let mut cursor = Cursor::new(&mut add_fst_data_vec);
let mut add_builder = SetBuilder::new(&mut cursor)?;
for key_vec in self.mutation_buffer.iter() {
add_builder.insert(key_vec)?;
}
add_builder.finish()?;
let add_fst_ds = FstDataSource::new_from_vec(Arc::new(add_fst_data_vec));
Set::new(add_fst_ds).map_err(IncrementalFstError::from)
});
match result_opt {
Ok(set_ref) => {
self.add_buffer_changed_since_last_cache = false;
Ok(Some(set_ref))
}
Err(e) => Err(e),
}
}
pub(crate) fn check_and_trigger_rebuild(&mut self, force_rebuild: bool) -> Result<RebuildOutcome> {
let num_persisted_keys = match self.get_persisted_set() {
Ok(set) => set.len(),
Err(_) => 0,
};
let active_mutations = self.mutation_buffer.len();
let tombstone_count = self.mutation_buffer.len_raw() - active_mutations;
let mut should_rebuild_now = (active_mutations >= self.rebuild_threshold_adds_count)
|| (num_persisted_keys > 0
&& tombstone_count > 0
&& (tombstone_count as f32 / num_persisted_keys as f32) >= self.rebuild_threshold_dels_ratio);
if self.rebuild_pending_due_to_interval {
should_rebuild_now = true;
}
if force_rebuild {
should_rebuild_now = true;
}
if should_rebuild_now {
if let Some(interval) = self.min_rebuild_interval {
let now = Instant::now();
if self
.last_rebuild_time
.is_none_or(|last| now.duration_since(last) >= interval)
{
self.rebuild_pending_due_to_interval = false;
} else {
self.rebuild_pending_due_to_interval = true;
return Ok(RebuildOutcome::RebuildDeferred);
}
}
let start_time = Instant::now();
self
.metrics
.persisted_keys_at_rebuild_sum
.fetch_add(num_persisted_keys, AtomicOrdering::Relaxed);
self
.metrics
.add_buffer_items_at_rebuild_sum
.fetch_add(active_mutations, AtomicOrdering::Relaxed);
self
.metrics
.del_buffer_items_at_rebuild_sum
.fetch_add(tombstone_count, AtomicOrdering::Relaxed);
self.rebuild_persisted_fst_impl()?;
self.metrics.num_rebuilds.fetch_add(1, AtomicOrdering::Relaxed);
let duration = start_time.elapsed();
self
.metrics
.rebuild_duration_micros_sum
.fetch_add(duration.as_micros() as u64, AtomicOrdering::Relaxed);
self.last_rebuild_time = Some(Instant::now());
return Ok(RebuildOutcome::Rebuilt);
}
Ok(RebuildOutcome::NotNeeded)
}
fn rebuild_persisted_fst_impl(&mut self) -> Result<()> {
let base_fst_ref = self.get_persisted_set()?;
let mut base_stream = base_fst_ref.stream();
let mut mutation_iter = self.mutation_buffer.iter_raw();
let mut temp_file = tempfile::tempfile().map_err(IncrementalFstError::Io)?;
let mut writer = BufWriter::new(&mut temp_file);
let mut builder = SetBuilder::new(&mut writer)?;
let mut next_base_key_opt = base_stream.next();
let mut next_mut_key_opt = mutation_iter.next();
loop {
let base_slice_opt = next_base_key_opt;
let mut_entry_opt = next_mut_key_opt;
match (base_slice_opt, mut_entry_opt) {
(Some(base_k), Some((mut_k, is_tombstone))) => {
if base_k < mut_k {
builder.insert(base_k)?;
next_base_key_opt = base_stream.next();
} else if mut_k < base_k {
if !is_tombstone {
builder.insert(mut_k)?;
}
next_mut_key_opt = mutation_iter.next();
} else {
if !is_tombstone {
builder.insert(mut_k)?;
}
next_base_key_opt = base_stream.next();
next_mut_key_opt = mutation_iter.next();
}
}
(Some(base_k), None) => {
builder.insert(base_k)?;
next_base_key_opt = base_stream.next();
}
(None, Some((mut_k, is_tombstone))) => {
if !is_tombstone {
builder.insert(mut_k)?;
}
next_mut_key_opt = mutation_iter.next();
}
(None, None) => break,
}
}
builder.finish()?;
writer.flush().map_err(IncrementalFstError::Io)?;
drop(writer);
let mmap = unsafe { Mmap::map(&temp_file).map_err(IncrementalFstError::Io)? };
self.persisted_fst_data_arc = FstDataSource::new_from_mmap(mmap);
self.cached_persisted_set_cell = OnceCell::new();
self.mutation_buffer.clear();
self.cached_add_buffer_fst_cell = OnceCell::new();
self.add_buffer_changed_since_last_cache = false;
Ok(())
}
}