Skip to main content

moirai_iter/parallel/
indexed.rs

1use super::{
2    ExponentialBlocks, Interleave, InterleaveShortest, IntoParallelIterator, ParallelIterator,
3    StepBy, UniformBlocks,
4};
5
6/// Exact-size boundary for Moirai's bounded Rayon-style indexed source subset.
7///
8/// This trait deliberately covers source iterators with known cardinality. It
9/// does not claim Rayon's full indexed producer/consumer adapter model.
10pub trait IndexedParallelIterator: ParallelIterator {
11    /// Return the exact number of logical items in the indexed source.
12    fn len(&self) -> usize;
13
14    /// Return whether the indexed source has no logical items.
15    fn is_empty(&self) -> bool {
16        self.len() == 0
17    }
18
19    /// Move all items into caller-provided storage.
20    ///
21    /// The destination vector is cleared but keeps its allocation, matching the
22    /// bounded source contract for exact-size streams without requiring item
23    /// cloning or allocating a second output vector.
24    fn collect_into_vec(self, target: &mut Vec<Self::Item>) {
25        target.clear();
26        target.extend(self.seq_items());
27    }
28
29    /// Split pair items into caller-provided left and right storage.
30    ///
31    /// Both destination vectors are cleared but keep their allocations. Pair
32    /// values are moved into their target sides exactly once, preserving the
33    /// bounded exact-size source contract without cloning either side.
34    fn unzip_into_vecs<A, B>(self, left: &mut Vec<A>, right: &mut Vec<B>)
35    where
36        Self: ParallelIterator<Item = (A, B)>,
37        A: Send,
38        B: Send,
39    {
40        let expected_len = self.len();
41        let items = self.seq_items();
42        debug_assert_eq!(items.len(), expected_len);
43
44        left.clear();
45        right.clear();
46        left.reserve_exact(expected_len);
47        right.reserve_exact(expected_len);
48
49        for (left_item, right_item) in items {
50            left.push(left_item);
51            right.push(right_item);
52        }
53    }
54
55    /// Alternately yield items from this source and another exact-size source.
56    ///
57    /// Values are moved from both sources into one logical stream. When one
58    /// side is exhausted, remaining values from the other side are yielded.
59    fn interleave<J>(self, other: J) -> Interleave<Self, J::Iter>
60    where
61        J: IntoParallelIterator<Item = Self::Item>,
62        J::Iter: IndexedParallelIterator<Item = Self::Item>,
63        Self::Item: Sync + 'static,
64    {
65        Interleave::new(self, other.into_par_iter())
66    }
67
68    /// Alternately yield items until the shorter exact-size source is consumed.
69    ///
70    /// This matches Rayon's indexed boundary: if the left source is longer,
71    /// one trailing left item is retained after the final right item.
72    fn interleave_shortest<J>(self, other: J) -> InterleaveShortest<Self, J::Iter>
73    where
74        J: IntoParallelIterator<Item = Self::Item>,
75        J::Iter: IndexedParallelIterator<Item = Self::Item>,
76        Self::Item: Sync + 'static,
77    {
78        InterleaveShortest::new(self, other.into_par_iter())
79    }
80
81    /// Yield every `step`th item from an exact-size source.
82    ///
83    /// The step size must be non-zero. Skipped items remain owned by the
84    /// consumed source iterator and are dropped exactly once.
85    fn step_by(self, step: usize) -> StepBy<Self>
86    where
87        Self::Item: Sync + 'static,
88    {
89        StepBy::new(self, step)
90    }
91
92    /// Convert this exact-size source into value-preserving exponential blocks.
93    ///
94    /// This bounded adapter preserves logical item order and exposes Rayon's
95    /// block-hint API surface. It does not claim Rayon's full indexed
96    /// producer/consumer block-splitting scheduler model.
97    fn by_exponential_blocks(self) -> ExponentialBlocks<Self>
98    where
99        Self::Item: Sync + 'static,
100    {
101        ExponentialBlocks::new(self)
102    }
103
104    /// Convert this exact-size source into value-preserving uniform blocks.
105    ///
106    /// The block size must be non-zero. This bounded adapter validates the
107    /// block-size contract and preserves logical item order without claiming
108    /// Rayon's full block-splitting producer model.
109    fn by_uniform_blocks(self, block_size: usize) -> UniformBlocks<Self>
110    where
111        Self::Item: Sync + 'static,
112    {
113        UniformBlocks::new(self, block_size)
114    }
115}
116
117impl<I, J> IndexedParallelIterator for Interleave<I, J>
118where
119    I: IndexedParallelIterator,
120    J: IndexedParallelIterator<Item = I::Item>,
121    I::Item: Sync + 'static,
122{
123    fn len(&self) -> usize {
124        self.left_len()
125            .checked_add(self.right_len())
126            .expect("overflow")
127    }
128}
129
130impl<I, J> Interleave<I, J>
131where
132    I: IndexedParallelIterator,
133    J: IndexedParallelIterator<Item = I::Item>,
134{
135    fn left_len(&self) -> usize {
136        self.left.len()
137    }
138
139    fn right_len(&self) -> usize {
140        self.right.len()
141    }
142}
143
144impl<I, J> IndexedParallelIterator for InterleaveShortest<I, J>
145where
146    I: IndexedParallelIterator,
147    J: IndexedParallelIterator<Item = I::Item>,
148    I::Item: Sync + 'static,
149{
150    fn len(&self) -> usize {
151        if self.left.len() <= self.right.len() {
152            self.left.len().checked_mul(2).expect("overflow")
153        } else {
154            self.right
155                .len()
156                .checked_mul(2)
157                .and_then(|len| len.checked_add(1))
158                .expect("overflow")
159        }
160    }
161}
162
163impl<I> IndexedParallelIterator for StepBy<I>
164where
165    I: IndexedParallelIterator,
166    I::Item: Sync + 'static,
167{
168    fn len(&self) -> usize {
169        let len = self.base.len();
170        if len == 0 {
171            0
172        } else {
173            ((len - 1) / self.step()) + 1
174        }
175    }
176}