fst_incremental 1.0.0

A thread-safe, updatable finite state set: dynamic insertions, deletions and queries over an immutable fst::Set fronted by a compact mutation buffer with amortized rebuilds.
Documentation
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(())
  }
}