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::{arena::ArenaRawIter, inner::IncrementalFstSetInner};

use std::ops::Deref;

use fst::{Streamer, automaton::AlwaysMatch, set::Stream as FstSetStreamGeneric};
use ouroboros::self_referencing;
use parking_lot::RwLockReadGuard;

pub struct PeekableFst<'s> {
  stream: FstSetStreamGeneric<'s, AlwaysMatch>,
  peeked: Option<Vec<u8>>,
  done: bool,
}

impl<'s> PeekableFst<'s> {
  pub fn new(stream: FstSetStreamGeneric<'s, AlwaysMatch>) -> Self {
    PeekableFst {
      stream,
      peeked: None,
      done: false,
    }
  }

  fn ensure_peeked(&mut self) {
    if self.peeked.is_none() && !self.done {
      if let Some(slice) = self.stream.next() {
        self.peeked = Some(slice.to_vec());
      } else {
        self.done = true;
      }
    }
  }

  pub fn peek(&mut self) -> Option<&[u8]> {
    self.ensure_peeked();
    self.peeked.as_deref()
  }

  pub fn consume(&mut self) -> Option<Vec<u8>> {
    self.ensure_peeked();
    if self.done && self.peeked.is_none() {
      return None;
    }
    self.peeked.take().or_else(|| {
      if self.done {
        None
      } else {
        if let Some(slice) = self.stream.next() {
          Some(slice.to_vec())
        } else {
          self.done = true;
          None
        }
      }
    })
  }
}

pub struct PeekableArenaRawIter<'s> {
  iter: ArenaRawIter<'s>,
  peeked: Option<(&'s [u8], bool)>, // (Key, IsTombstone)
  done: bool,
}

impl<'s> PeekableArenaRawIter<'s> {
  pub fn new(iter: ArenaRawIter<'s>) -> Self {
    PeekableArenaRawIter {
      iter,
      peeked: None,
      done: false,
    }
  }

  fn ensure_peeked(&mut self) {
    if self.peeked.is_none() && !self.done {
      if let Some(item) = self.iter.next() {
        self.peeked = Some(item);
      } else {
        self.done = true;
      }
    }
  }

  pub fn peek(&mut self) -> Option<(&'s [u8], bool)> {
    self.ensure_peeked();
    self.peeked
  }

  pub fn consume(&mut self) -> Option<(&'s [u8], bool)> {
    self.ensure_peeked();
    if self.done && self.peeked.is_none() {
      return None;
    }
    self.peeked.take().or_else(|| {
      if self.done {
        None
      } else {
        self.iter.next().or_else(|| {
          self.done = true;
          None
        })
      }
    })
  }
}

/// A sorted, deduplicated view over the persisted FST merged with the mutation buffer.
///
/// Holds a read guard on the set for as long as it lives, so no writer can proceed until
/// the stream is dropped. Implements [`fst::Streamer`]; bring that trait into scope to
/// call `next()`.
#[self_referencing]
pub struct MergedSetStreamOwner<'a> {
  guard: RwLockReadGuard<'a, IncrementalFstSetInner>,

  #[borrows(guard)]
  #[covariant]
  persisted_source: PeekableFst<'this>,

  #[borrows(guard)]
  #[covariant]
  mutation_source: PeekableArenaRawIter<'this>,
}

impl<'a> Streamer<'_> for MergedSetStreamOwner<'a> {
  type Item = Vec<u8>;

  fn next(&mut self) -> Option<Self::Item> {
    loop {
      let p_peek_opt: Option<Vec<u8>> = self.with_persisted_source_mut(|ps| ps.peek().map(|s| s.to_vec()));
      let m_peek_opt: Option<(Vec<u8>, bool)> =
        self.with_mutation_source_mut(|ms| ms.peek().map(|(k, t)| (k.to_vec(), t)));

      match (p_peek_opt.as_deref(), m_peek_opt.as_ref()) {
        (Some(p_slice), Some((m_slice, is_tombstone))) => {
          if p_slice < m_slice {
            return self.with_persisted_source_mut(|ps| ps.consume());
          } else if m_slice.as_slice() < p_slice {
            let _ = self.with_mutation_source_mut(|ms| ms.consume());
            if !*is_tombstone {
              return Some(m_slice.clone());
            }
          } else {
            let _ = self.with_persisted_source_mut(|ps| ps.consume());
            let _ = self.with_mutation_source_mut(|ms| ms.consume());

            if !*is_tombstone {
              return Some(m_slice.clone());
            }
          }
        }
        (Some(_), None) => {
          return self.with_persisted_source_mut(|ps| ps.consume());
        }
        (None, Some((m_slice, is_tombstone))) => {
          let _ = self.with_mutation_source_mut(|ms| ms.consume());
          if !*is_tombstone {
            return Some(m_slice.clone());
          }
        }
        (None, None) => {
          return None;
        }
      }
    }
  }
}

pub(crate) fn new_merged_set_stream_owner<'a>(
  guard_param: RwLockReadGuard<'a, IncrementalFstSetInner>,
) -> crate::error::Result<MergedSetStreamOwner<'a>> {
  MergedSetStreamOwner::try_new(
    guard_param,
    |guard_field_ref| {
      let inner_ref = guard_field_ref.deref();
      Ok(PeekableFst::new(inner_ref.get_persisted_set()?.stream()))
    },
    |guard_field_ref| {
      let inner_ref = guard_field_ref.deref();
      Ok(PeekableArenaRawIter::new(inner_ref.mutation_buffer.iter_raw()))
    },
  )
}