Skip to main content

subms_merge_iterator/
lib.rs

1//! K-way merge iterator. Min-heap of stream heads; pop the minimum, advance
2//! that stream, push its next value if any.
3//!
4//! Streams must be sorted ascending. Output is the global sorted union.
5//!
6//! ```
7//! use subms_merge_iterator::MergeIterator;
8//! let streams: Vec<Box<dyn Iterator<Item = i32>>> = vec![
9//!     Box::new(vec![1, 4, 7].into_iter()),
10//!     Box::new(vec![2, 5, 8].into_iter()),
11//!     Box::new(vec![3, 6, 9].into_iter()),
12//! ];
13//! let merged: Vec<_> = MergeIterator::new(streams).collect();
14//! assert_eq!(merged, (1..=9).collect::<Vec<_>>());
15//! ```
16
17use std::cmp::Reverse;
18use std::collections::BinaryHeap;
19
20pub struct MergeIterator<T: Ord, I: Iterator<Item = T>> {
21    streams: Vec<I>,
22    /// Min-heap on (value, stream index). `Reverse` flips to min ordering.
23    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// Opt-in feature modules. Each is independent of the base merge iterator
54// and gated by its own Cargo feature; the bare `cargo add
55// subms-merge-iterator` shape stays zero-dep + std-only.
56#[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};