use std::cmp::Reverse;
use std::collections::BinaryHeap;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TombstoneEntry<K, V> {
pub key: K,
pub value: Option<V>,
}
impl<K, V> TombstoneEntry<K, V> {
pub fn live(key: K, value: V) -> Self {
Self {
key,
value: Some(value),
}
}
pub fn tombstone(key: K) -> Self {
Self { key, value: None }
}
pub fn is_tombstone(&self) -> bool {
self.value.is_none()
}
}
pub struct TombstoneMergeIterator<K, V, I>
where
K: Ord,
I: Iterator<Item = TombstoneEntry<K, V>>,
{
streams: Vec<I>,
heap: BinaryHeap<Reverse<HeapItem<K, V>>>,
}
struct HeapItem<K, V> {
key: K,
source: usize,
value: Option<V>,
}
impl<K: Ord, V> PartialEq for HeapItem<K, V> {
fn eq(&self, other: &Self) -> bool {
self.key == other.key && self.source == other.source
}
}
impl<K: Ord, V> Eq for HeapItem<K, V> {}
impl<K: Ord, V> PartialOrd for HeapItem<K, V> {
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl<K: Ord, V> Ord for HeapItem<K, V> {
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
self.key
.cmp(&other.key)
.then(other.source.cmp(&self.source))
}
}
impl<K, V, I> TombstoneMergeIterator<K, V, I>
where
K: Ord,
I: Iterator<Item = TombstoneEntry<K, V>>,
{
pub fn new<S: IntoIterator<Item = I>>(streams: S) -> Self {
let mut streams: Vec<I> = streams.into_iter().collect();
let mut heap = BinaryHeap::with_capacity(streams.len());
for (i, s) in streams.iter_mut().enumerate() {
if let Some(e) = s.next() {
heap.push(Reverse(HeapItem {
key: e.key,
source: i,
value: e.value,
}));
}
}
Self { streams, heap }
}
}
impl<K, V, I> Iterator for TombstoneMergeIterator<K, V, I>
where
K: Ord,
I: Iterator<Item = TombstoneEntry<K, V>>,
{
type Item = TombstoneEntry<K, V>;
fn next(&mut self) -> Option<TombstoneEntry<K, V>> {
loop {
let Reverse(HeapItem {
key: winning_key,
source,
value: winning_value,
}) = self.heap.pop()?;
self.advance(source);
while let Some(Reverse(item)) = self.heap.peek() {
if item.key == winning_key {
let Reverse(item) = self.heap.pop().unwrap();
self.advance(item.source);
} else {
break;
}
}
if winning_value.is_none() {
continue;
}
return Some(TombstoneEntry {
key: winning_key,
value: winning_value,
});
}
}
}
impl<K, V, I> TombstoneMergeIterator<K, V, I>
where
K: Ord,
I: Iterator<Item = TombstoneEntry<K, V>>,
{
fn advance(&mut self, source: usize) {
if let Some(e) = self.streams[source].next() {
self.heap.push(Reverse(HeapItem {
key: e.key,
source,
value: e.value,
}));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn live(k: &str, v: &str) -> TombstoneEntry<String, String> {
TombstoneEntry::live(k.to_string(), v.to_string())
}
fn tomb(k: &str) -> TombstoneEntry<String, String> {
TombstoneEntry::tombstone(k.to_string())
}
#[test]
fn live_entries_passthrough_when_no_tombstones() {
let s0 = vec![live("a", "1"), live("c", "3")];
let s1 = vec![live("b", "2"), live("d", "4")];
let merged: Vec<_> =
TombstoneMergeIterator::new([s0.into_iter(), s1.into_iter()]).collect();
assert_eq!(
merged,
vec![
live("a", "1"),
live("b", "2"),
live("c", "3"),
live("d", "4")
]
);
}
#[test]
fn later_source_tombstone_hides_earlier_value() {
let older = vec![live("k", "v")];
let newer = vec![tomb("k")];
let merged: Vec<_> =
TombstoneMergeIterator::new([older.into_iter(), newer.into_iter()]).collect();
assert!(merged.is_empty(), "tombstone in newer source must hide k");
}
#[test]
fn later_source_live_overwrites_earlier_value() {
let older = vec![live("k", "old")];
let newer = vec![live("k", "new")];
let merged: Vec<_> =
TombstoneMergeIterator::new([older.into_iter(), newer.into_iter()]).collect();
assert_eq!(merged, vec![live("k", "new")]);
}
#[test]
fn earlier_source_tombstone_does_not_hide_later_value() {
let older = vec![tomb("k")];
let newer = vec![live("k", "v")];
let merged: Vec<_> =
TombstoneMergeIterator::new([older.into_iter(), newer.into_iter()]).collect();
assert_eq!(merged, vec![live("k", "v")]);
}
#[test]
fn tombstone_shadowing_spans_three_sources() {
let s0 = vec![live("a", "1"), live("b", "2")];
let s1 = vec![tomb("a")];
let s2 = vec![live("c", "3")];
let merged: Vec<_> =
TombstoneMergeIterator::new([s0.into_iter(), s1.into_iter(), s2.into_iter()]).collect();
assert_eq!(merged, vec![live("b", "2"), live("c", "3")]);
}
#[test]
fn all_sources_empty_yields_nothing() {
let s0: Vec<TombstoneEntry<String, String>> = vec![];
let s1: Vec<TombstoneEntry<String, String>> = vec![];
let merged: Vec<_> =
TombstoneMergeIterator::new([s0.into_iter(), s1.into_iter()]).collect();
assert!(merged.is_empty());
}
#[test]
fn all_tombstones_yields_nothing() {
let s0 = vec![tomb("a"), tomb("b")];
let s1 = vec![tomb("a"), tomb("c")];
let merged: Vec<_> =
TombstoneMergeIterator::new([s0.into_iter(), s1.into_iter()]).collect();
assert!(merged.is_empty());
}
#[test]
fn interleaved_tombstones_and_live_resolve_per_key() {
let s0 = vec![live("a", "1"), live("b", "2"), live("c", "3")];
let s1 = vec![tomb("a"), live("b", "new")];
let s2 = vec![tomb("b")];
let merged: Vec<_> =
TombstoneMergeIterator::new([s0.into_iter(), s1.into_iter(), s2.into_iter()]).collect();
assert_eq!(merged, vec![live("c", "3")]);
}
}