subms_merge_iterator/
lib.rs1use std::cmp::Reverse;
18use std::collections::BinaryHeap;
19
20pub struct MergeIterator<T: Ord, I: Iterator<Item = T>> {
21 streams: Vec<I>,
22 heap: BinaryHeap<Reverse<(T, usize)>>,
24}
25
26impl<T: Ord, I: Iterator<Item = T>> MergeIterator<T, I> {
27 pub fn new<S: IntoIterator<Item = I>>(streams: S) -> Self {
28 let mut streams: Vec<I> = streams.into_iter().collect();
29 let mut heap = BinaryHeap::with_capacity(streams.len());
30 for (i, s) in streams.iter_mut().enumerate() {
31 if let Some(v) = s.next() {
32 heap.push(Reverse((v, i)));
33 }
34 }
35 Self { streams, heap }
36 }
37}
38
39impl<T: Ord, I: Iterator<Item = T>> Iterator for MergeIterator<T, I> {
40 type Item = T;
41 fn next(&mut self) -> Option<T> {
42 let Reverse((value, idx)) = self.heap.pop()?;
43 if let Some(next_value) = self.streams[idx].next() {
44 self.heap.push(Reverse((next_value, idx)));
45 }
46 Some(value)
47 }
48}
49
50#[cfg(feature = "harness")]
51pub mod recipe;
52
53#[cfg(any(
57 feature = "seek-to",
58 feature = "tombstones",
59 feature = "dedup",
60 feature = "priority"
61))]
62pub mod features;
63
64#[cfg(feature = "dedup")]
65pub use features::dedup::{DedupEntry, DedupMergeIterator};
66#[cfg(feature = "priority")]
67pub use features::priority::{PriorityEntry, PriorityMergeIterator, PrioritySource};
68#[cfg(feature = "seek-to")]
69pub use features::seek::SeekableMergeIterator;
70#[cfg(feature = "tombstones")]
71pub use features::tombstones::{TombstoneEntry, TombstoneMergeIterator};