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)>, 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
})
}
})
}
}
#[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()))
},
)
}