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