moirai_iter/parallel/
indexed.rs1use super::{
2 ExponentialBlocks, Interleave, InterleaveShortest, IntoParallelIterator, ParallelIterator,
3 StepBy, UniformBlocks,
4};
5
6pub trait IndexedParallelIterator: ParallelIterator {
11 fn len(&self) -> usize;
13
14 fn is_empty(&self) -> bool {
16 self.len() == 0
17 }
18
19 fn collect_into_vec(self, target: &mut Vec<Self::Item>) {
25 target.clear();
26 target.extend(self.seq_items());
27 }
28
29 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 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 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 fn step_by(self, step: usize) -> StepBy<Self>
86 where
87 Self::Item: Sync + 'static,
88 {
89 StepBy::new(self, step)
90 }
91
92 fn by_exponential_blocks(self) -> ExponentialBlocks<Self>
98 where
99 Self::Item: Sync + 'static,
100 {
101 ExponentialBlocks::new(self)
102 }
103
104 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}