use crate::direction::Direction;
use crate::queue::Merge;
use crate::versions::Versions;
use bytes::Bytes;
use crossbeam_skiplist::map::Entry;
use crossbeam_skiplist::map::Range as SkipRange;
use parking_lot::RwLock;
use std::collections::btree_map::Range as TreeRange;
use std::ops::Bound;
use std::sync::Arc;
pub(crate) type SkipBounds = (Bound<Bytes>, Bound<Bytes>);
pub(crate) struct MergeQueueIter {
sources: Vec<Arc<Merge>>,
heads: Vec<Option<(Bytes, Option<Bytes>)>>,
beg: Bytes,
end: Bytes,
direction: Direction,
}
impl MergeQueueIter {
pub(crate) fn new(
sources: Vec<Arc<Merge>>,
beg: Bytes,
end: Bytes,
direction: Direction,
) -> Self {
let mut heads = Vec::with_capacity(sources.len());
for src in &sources {
heads.push(seek_in_writeset(src, direction, &beg, &end, None));
}
Self {
sources,
heads,
beg,
end,
direction,
}
}
}
fn seek_in_writeset(
src: &Arc<Merge>,
direction: Direction,
beg: &Bytes,
end: &Bytes,
after: Option<&Bytes>,
) -> Option<(Bytes, Option<Bytes>)> {
let ws = &src.writeset;
let entry = match (direction, after) {
(Direction::Forward, None) => {
ws.range::<Bytes, _>((Bound::Included(beg), Bound::Excluded(end))).next()
}
(Direction::Forward, Some(k)) => {
ws.range::<Bytes, _>((Bound::Excluded(k), Bound::Excluded(end))).next()
}
(Direction::Reverse, None) => {
ws.range::<Bytes, _>((Bound::Included(beg), Bound::Excluded(end))).next_back()
}
(Direction::Reverse, Some(k)) => {
ws.range::<Bytes, _>((Bound::Included(beg), Bound::Excluded(k))).next_back()
}
};
entry.map(|(k, v)| (k.clone(), v.clone()))
}
impl Iterator for MergeQueueIter {
type Item = (Bytes, Option<Bytes>);
fn next(&mut self) -> Option<Self::Item> {
let mut winner: Option<usize> = None;
for (i, head) in self.heads.iter().enumerate() {
let Some((k, _)) = head else {
continue;
};
match winner {
None => winner = Some(i),
Some(wi) => {
let (wk, _) = self.heads[wi].as_ref().unwrap();
let take = match self.direction {
Direction::Forward => k < wk,
Direction::Reverse => k > wk,
};
if take {
winner = Some(i);
}
}
}
}
let winner = winner?;
let (out_key, out_val) = self.heads[winner].take().unwrap();
for i in (winner + 1)..self.heads.len() {
let same = self.heads[i].as_ref().map(|(k, _)| k == &out_key).unwrap_or(false);
if same {
let new = seek_in_writeset(
&self.sources[i],
self.direction,
&self.beg,
&self.end,
Some(&out_key),
);
self.heads[i] = new;
}
}
let new = seek_in_writeset(
&self.sources[winner],
self.direction,
&self.beg,
&self.end,
Some(&out_key),
);
self.heads[winner] = new;
Some((out_key, out_val))
}
}
pub struct MergeIterator<'a> {
pub(crate) tree_iter: SkipRange<'a, Bytes, SkipBounds, Bytes, RwLock<Versions>>,
pub(crate) self_iter: TreeRange<'a, Bytes, Option<Bytes>>,
pub(crate) join_iter: Box<dyn Iterator<Item = (Bytes, Option<Bytes>)> + 'a>,
pub(crate) tree_next: Option<Entry<'a, Bytes, RwLock<Versions>>>,
pub(crate) join_next: Option<(Bytes, Option<Bytes>)>,
pub(crate) self_next: Option<(&'a Bytes, &'a Option<Bytes>)>,
pub(crate) direction: Direction,
pub(crate) version: u64,
pub(crate) skip_remaining: usize,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum KeySource {
None,
Datastore,
Committed,
Transaction,
}
impl<'a> MergeIterator<'a> {
pub fn new(
mut tree_iter: SkipRange<'a, Bytes, SkipBounds, Bytes, RwLock<Versions>>,
mut join_iter: Box<dyn Iterator<Item = (Bytes, Option<Bytes>)> + 'a>,
mut self_iter: TreeRange<'a, Bytes, Option<Bytes>>,
direction: Direction,
version: u64,
skip: usize,
) -> Self {
let tree_next = match direction {
Direction::Forward => tree_iter.next(),
Direction::Reverse => tree_iter.next_back(),
};
let self_next = match direction {
Direction::Forward => self_iter.next(),
Direction::Reverse => self_iter.next_back(),
};
let join_next = join_iter.next();
MergeIterator {
tree_iter,
self_iter,
join_iter,
tree_next,
join_next,
self_next,
direction,
version,
skip_remaining: skip,
}
}
#[inline]
fn advance_join(&mut self) {
self.join_next = self.join_iter.next();
}
pub fn next_count(&mut self) -> Option<bool> {
loop {
let mut next_key: Option<&Bytes> = None;
let mut next_source = KeySource::None;
if let Some((sk, _)) = self.self_next {
next_key = Some(sk);
next_source = KeySource::Transaction;
}
if let Some((jk, _)) = &self.join_next {
let should_use = match (next_key, &self.direction) {
(None, _) => true,
(Some(k), Direction::Forward) => jk < k,
(Some(k), Direction::Reverse) => jk > k,
};
if should_use {
next_key = Some(jk);
next_source = KeySource::Committed;
} else if next_key == Some(jk) {
next_source = KeySource::Transaction;
}
}
if let Some(t_entry) = &self.tree_next {
let tk = t_entry.key();
let should_use = match (next_key, &self.direction) {
(None, _) => true,
(Some(k), Direction::Forward) => tk < k,
(Some(k), Direction::Reverse) => tk > k,
};
if should_use {
next_source = KeySource::Datastore;
}
}
let exists = match next_source {
KeySource::Transaction => {
let (sk, sv) = self.self_next.unwrap();
let exists = sv.is_some();
self.self_next = match self.direction {
Direction::Forward => self.self_iter.next(),
Direction::Reverse => self.self_iter.next_back(),
};
if let Some((jk, _)) = &self.join_next {
if jk == sk {
self.advance_join();
}
}
if let Some(t_entry) = &self.tree_next {
if t_entry.key() == sk {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
}
exists
}
KeySource::Committed => {
let exists = self.join_next.as_ref().unwrap().1.is_some();
let should_skip_tree = if let Some(t_entry) = &self.tree_next {
if let Some((jk, _)) = &self.join_next {
t_entry.key() == jk
} else {
false
}
} else {
false
};
self.advance_join();
if should_skip_tree {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
exists
}
KeySource::Datastore => {
let t_entry = self.tree_next.as_ref().unwrap();
let tv = match t_entry.value().try_read() {
Some(guard) => guard,
None => t_entry.value().read(),
};
let exists = tv.exists_version(self.version);
drop(tv);
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
exists
}
KeySource::None => return None,
};
if exists && self.skip_remaining > 0 {
self.skip_remaining -= 1;
continue;
}
return Some(exists);
}
}
pub fn next_key(&mut self) -> Option<(Bytes, bool)> {
loop {
let mut next_key: Option<&Bytes> = None;
let mut next_source = KeySource::None;
if let Some((sk, _)) = self.self_next {
next_key = Some(sk);
next_source = KeySource::Transaction;
}
if let Some((jk, _)) = &self.join_next {
let should_use = match (next_key, &self.direction) {
(None, _) => true,
(Some(k), Direction::Forward) => jk < k,
(Some(k), Direction::Reverse) => jk > k,
};
if should_use {
next_key = Some(jk);
next_source = KeySource::Committed;
} else if next_key == Some(jk) {
next_source = KeySource::Transaction;
}
}
if let Some(t_entry) = &self.tree_next {
let tk = t_entry.key();
let should_use = match (next_key, &self.direction) {
(None, _) => true,
(Some(k), Direction::Forward) => tk < k,
(Some(k), Direction::Reverse) => tk > k,
};
if should_use {
next_source = KeySource::Datastore;
}
}
match next_source {
KeySource::Transaction => {
let (sk, sv) = self.self_next.unwrap();
let exists = sv.is_some();
let key_ref = sk;
self.self_next = match self.direction {
Direction::Forward => self.self_iter.next(),
Direction::Reverse => self.self_iter.next_back(),
};
if let Some((jk, _)) = &self.join_next {
if jk == key_ref {
self.advance_join();
}
}
if let Some(t_entry) = &self.tree_next {
if t_entry.key() == key_ref {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
}
if exists && self.skip_remaining > 0 {
self.skip_remaining -= 1;
continue;
}
return Some((key_ref.clone(), exists));
}
KeySource::Committed => {
let (jk, jv) = self.join_next.as_ref().unwrap();
if jv.is_some() && self.skip_remaining > 0 {
let should_skip_tree = if let Some(t_entry) = &self.tree_next {
t_entry.key() == jk
} else {
false
};
self.advance_join();
if should_skip_tree {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
self.skip_remaining -= 1;
continue;
}
let exists = jv.is_some();
let key = jk.clone();
self.advance_join();
if let Some(t_entry) = &self.tree_next {
if t_entry.key() == &key {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
}
return Some((key, exists));
}
KeySource::Datastore => {
let t_entry = self.tree_next.as_ref().unwrap();
let tv = match t_entry.value().try_read() {
Some(guard) => guard,
None => t_entry.value().read(),
};
let exists = tv.exists_version(self.version);
drop(tv);
if exists && self.skip_remaining > 0 {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
self.skip_remaining -= 1;
continue;
}
let tk = t_entry.key().clone();
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
return Some((tk, exists));
}
KeySource::None => return None,
}
}
}
}
impl<'a> Iterator for MergeIterator<'a> {
type Item = (Bytes, Option<Bytes>);
fn next(&mut self) -> Option<Self::Item> {
loop {
let mut next_key: Option<&Bytes> = None;
let mut next_source = KeySource::None;
if let Some((sk, _)) = self.self_next {
next_key = Some(sk);
next_source = KeySource::Transaction;
}
if let Some((jk, _)) = &self.join_next {
let should_use = match (next_key, &self.direction) {
(None, _) => true,
(Some(k), Direction::Forward) => jk < k,
(Some(k), Direction::Reverse) => jk > k,
};
if should_use {
next_key = Some(jk);
next_source = KeySource::Committed;
} else if next_key == Some(jk) {
next_source = KeySource::Transaction;
}
}
if let Some(t_entry) = &self.tree_next {
let tk = t_entry.key();
let should_use = match (next_key, &self.direction) {
(None, _) => true,
(Some(k), Direction::Forward) => tk < k,
(Some(k), Direction::Reverse) => tk > k,
};
if should_use {
next_source = KeySource::Datastore;
}
}
match next_source {
KeySource::Transaction => {
let (sk, sv) = self.self_next.unwrap();
let exists = sv.is_some();
let key_ref = sk;
self.self_next = match self.direction {
Direction::Forward => self.self_iter.next(),
Direction::Reverse => self.self_iter.next_back(),
};
if let Some((jk, _)) = &self.join_next {
if jk == key_ref {
self.advance_join();
}
}
if let Some(t_entry) = &self.tree_next {
if t_entry.key() == key_ref {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
}
if exists && self.skip_remaining > 0 {
self.skip_remaining -= 1;
continue;
}
return Some((key_ref.clone(), sv.clone()));
}
KeySource::Committed => {
let (jk, jv) = self.join_next.as_ref().unwrap();
if jv.is_some() && self.skip_remaining > 0 {
let should_skip_tree = if let Some(t_entry) = &self.tree_next {
t_entry.key() == jk
} else {
false
};
self.advance_join();
if should_skip_tree {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
self.skip_remaining -= 1;
continue;
}
let key = jk.clone();
let value_opt = jv.clone();
self.advance_join();
if let Some(t_entry) = &self.tree_next {
if t_entry.key() == &key {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
}
}
return Some((key, value_opt));
}
KeySource::Datastore => {
let t_entry = self.tree_next.as_ref().unwrap();
let tv = match t_entry.value().try_read() {
Some(guard) => guard,
None => t_entry.value().read(),
};
let value_opt = tv.fetch_version(self.version);
let exists = value_opt.is_some();
drop(tv);
if exists && self.skip_remaining > 0 {
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
self.skip_remaining -= 1;
continue;
}
let key_clone = t_entry.key().clone();
self.tree_next = match self.direction {
Direction::Forward => self.tree_iter.next(),
Direction::Reverse => self.tree_iter.next_back(),
};
return Some((key_clone, value_opt));
}
KeySource::None => return None,
}
}
}
}