use std::cmp::Reverse;
use std::collections::BinaryHeap;
pub struct MergeIterator<T: Ord, I: Iterator<Item = T>> {
streams: Vec<I>,
heap: BinaryHeap<Reverse<(T, usize)>>,
}
impl<T: Ord, I: Iterator<Item = T>> MergeIterator<T, I> {
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(v) = s.next() {
heap.push(Reverse((v, i)));
}
}
Self { streams, heap }
}
pub fn peek(&self) -> Option<&T> {
self.heap.peek().map(|Reverse((value, _))| value)
}
pub fn live_streams(&self) -> usize {
self.heap.len()
}
pub fn num_streams(&self) -> usize {
self.streams.len()
}
}
impl<T: Ord, I: Iterator<Item = T>> Iterator for MergeIterator<T, I> {
type Item = T;
fn next(&mut self) -> Option<T> {
let Reverse((value, idx)) = self.heap.pop()?;
if let Some(next_value) = self.streams[idx].next() {
self.heap.push(Reverse((next_value, idx)));
}
Some(value)
}
}
#[cfg(feature = "harness")]
pub mod recipe;
#[cfg(any(
feature = "seek-to",
feature = "reverse",
feature = "tombstones",
feature = "dedup",
feature = "priority"
))]
pub mod features;
#[cfg(feature = "dedup")]
pub use features::dedup::{DedupEntry, DedupMergeIterator};
#[cfg(feature = "priority")]
pub use features::priority::{PriorityEntry, PriorityMergeIterator, PrioritySource};
#[cfg(feature = "reverse")]
pub use features::reverse::ReverseMergeIterator;
#[cfg(feature = "seek-to")]
pub use features::seek::SeekableMergeIterator;
#[cfg(feature = "tombstones")]
pub use features::tombstones::{TombstoneEntry, TombstoneMergeIterator};
#[cfg(test)]
#[path = "merge_tests.rs"]
mod merge_tests;
#[cfg(test)]
#[path = "sample_app_tests.rs"]
mod sample_app_tests;