use super::sstable::SSTableIterator;
use super::{Key, Value};
use crate::Result;
use std::cmp::{Ordering, Reverse};
use std::collections::BinaryHeap;
use std::sync::Arc;
type KVIterator = Box<dyn Iterator<Item = Result<(Key, Value)>> + Send>;
#[derive(Debug, Clone)]
struct HeapItem {
key: Key,
value: Value,
source_id: usize, }
impl PartialEq for HeapItem {
fn eq(&self, other: &Self) -> bool {
self.key == other.key && self.value.timestamp == other.value.timestamp
}
}
impl Eq for HeapItem {}
impl PartialOrd for HeapItem {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl Ord for HeapItem {
fn cmp(&self, other: &Self) -> Ordering {
self.key
.cmp(&other.key)
.then(other.value.timestamp.cmp(&self.value.timestamp)) .then(self.source_id.cmp(&other.source_id))
}
}
pub struct MergingIterator {
heap: BinaryHeap<Reverse<HeapItem>>,
sources: Vec<KVIterator>,
last_key: Option<Key>,
finished: bool,
first_error: Option<crate::StorageError>,
single_source: bool,
raw_sst: Option<SSTableIterator>,
}
impl MergingIterator {
pub fn new(sources: Vec<KVIterator>) -> Self {
let single = sources.len() == 1;
let mut iter = Self {
heap: BinaryHeap::new(),
sources,
last_key: None,
finished: false,
first_error: None,
single_source: single,
raw_sst: None,
};
if !single {
iter.fill_heap();
}
iter
}
pub fn new_raw_sst(sst: SSTableIterator) -> Self {
Self {
heap: BinaryHeap::new(),
sources: Vec::new(),
last_key: None,
finished: false,
first_error: None,
single_source: true,
raw_sst: Some(sst),
}
}
pub fn next_raw(&mut self) -> Option<Result<(Key, u64, bool, super::sstable::ValueBytes)>> {
if let Some(ref mut sst) = self.raw_sst {
return sst
.next_raw()
.map(|(key, ts, del, vb)| Ok((key, ts, del, vb)));
}
match self.next() {
Some(Ok((key, value))) => {
let vb = match &value.data {
super::ValueData::Inline(arc_vec) => {
let len = arc_vec.len();
super::sstable::ValueBytes {
block: Arc::clone(arc_vec),
start: 0,
len,
}
}
super::ValueData::Blob(_) => super::sstable::ValueBytes {
block: Arc::new(Vec::new()),
start: 0,
len: 0,
},
};
Some(Ok((key, value.timestamp, value.deleted, vb)))
}
Some(Err(e)) => Some(Err(e)),
None => None,
}
}
pub fn has_raw_sst(&self) -> bool {
self.raw_sst.is_some()
}
fn fill_heap(&mut self) {
for (source_id, source) in self.sources.iter_mut().enumerate() {
match source.next() {
Some(Ok((key, value))) => {
self.heap.push(Reverse(HeapItem {
key,
value,
source_id,
}));
}
Some(Err(e)) => {
if self.first_error.is_none() {
self.first_error = Some(e);
}
}
None => {}
}
}
}
fn refill_from_source(&mut self, source_id: usize) {
if let Some(source) = self.sources.get_mut(source_id) {
match source.next() {
Some(Ok((key, value))) => {
self.heap.push(Reverse(HeapItem {
key,
value,
source_id,
}));
}
Some(Err(e)) => {
if self.first_error.is_none() {
self.first_error = Some(e);
}
}
None => {}
}
}
}
}
impl Iterator for MergingIterator {
type Item = Result<(Key, Value)>;
fn next(&mut self) -> Option<Self::Item> {
if self.finished {
return None;
}
if self.single_source {
let source = match self.sources.get_mut(0) {
Some(s) => s,
None => {
self.finished = true;
return None;
}
};
loop {
match source.next() {
Some(Ok((key, value))) => {
if let Some(last_key) = self.last_key {
if key == last_key {
continue;
}
}
if value.deleted {
self.last_key = Some(key);
continue;
}
self.last_key = Some(key);
return Some(Ok((key, value)));
}
Some(Err(e)) => {
self.finished = true;
return Some(Err(e));
}
None => {
self.finished = true;
return None;
}
}
}
}
loop {
let Reverse(item) = match self.heap.pop() {
Some(item) => item,
None => {
self.finished = true;
if let Some(e) = self.first_error.take() {
return Some(Err(e));
}
return None;
}
};
self.refill_from_source(item.source_id);
if let Some(last_key) = self.last_key {
if item.key == last_key {
continue;
}
}
if item.value.deleted {
self.last_key = Some(item.key);
continue;
}
self.last_key = Some(item.key);
return Some(Ok((item.key, item.value)));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::lsm::ValueData;
type BoxedIter = Box<dyn Iterator<Item = Result<(Key, Value)>> + Send>;
#[test]
fn test_merging_iterator_basic() {
let source1: Vec<Result<(Key, Value)>> = vec![
Ok((1, Value::new(vec![1], 100))),
Ok((3, Value::new(vec![3], 100))),
Ok((5, Value::new(vec![5], 100))),
];
let source2: Vec<Result<(Key, Value)>> = vec![
Ok((2, Value::new(vec![2], 100))),
Ok((4, Value::new(vec![4], 100))),
Ok((6, Value::new(vec![6], 100))),
];
let sources: Vec<BoxedIter> =
vec![Box::new(source1.into_iter()), Box::new(source2.into_iter())];
let iter = MergingIterator::new(sources);
let keys: Vec<Key> = iter.map(|r| r.unwrap().0).collect();
assert_eq!(keys, vec![1, 2, 3, 4, 5, 6]);
}
#[test]
fn test_merging_iterator_mvcc() {
let source1: Vec<Result<(Key, Value)>> = vec![Ok((
1,
Value {
data: ValueData::Inline(std::sync::Arc::new(vec![1, 0, 0])), timestamp: 300,
deleted: false,
},
))];
let source2: Vec<Result<(Key, Value)>> = vec![Ok((
1,
Value {
data: ValueData::Inline(std::sync::Arc::new(vec![1, 0])), timestamp: 200,
deleted: false,
},
))];
let source3: Vec<Result<(Key, Value)>> = vec![Ok((
1,
Value {
data: ValueData::Inline(std::sync::Arc::new(vec![1])), timestamp: 100,
deleted: false,
},
))];
let sources: Vec<BoxedIter> = vec![
Box::new(source1.into_iter()),
Box::new(source2.into_iter()),
Box::new(source3.into_iter()),
];
let iter = MergingIterator::new(sources);
let results: Vec<(Key, Vec<u8>)> = iter
.map(|r| {
let (k, v) = r.unwrap();
match v.data {
ValueData::Inline(data) => (k, data.to_vec()),
_ => panic!("Expected inline data"),
}
})
.collect();
assert_eq!(results.len(), 1);
assert_eq!(results[0].0, 1);
assert_eq!(results[0].1, vec![1, 0, 0]); }
#[test]
fn test_merging_iterator_tombstone() {
let source1: Vec<Result<(Key, Value)>> = vec![
Ok((1, Value::new(vec![1], 100))),
Ok((
2,
Value {
data: ValueData::Inline(std::sync::Arc::new(vec![])),
timestamp: 200,
deleted: true, },
)),
Ok((3, Value::new(vec![3], 100))),
];
let sources: Vec<BoxedIter> = vec![Box::new(source1.into_iter())];
let iter = MergingIterator::new(sources);
let keys: Vec<Key> = iter.map(|r| r.unwrap().0).collect();
assert_eq!(keys, vec![1, 3]);
}
}