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
use super::super::{Consumer, MapConsumer, ParallelIterator, VecParIter};
use std::ops::ControlFlow;
/// Enumerate adapter for value-semantic index pairing.
pub struct Enumerate<I> {
pub(super) base: I,
}
impl<I> Enumerate<I> {
pub(crate) fn new(base: I) -> Self {
Self { base }
}
}
impl<I> ParallelIterator for Enumerate<I>
where
I: ParallelIterator,
I::Item: Sync + 'static,
{
type Item = (usize, I::Item);
fn seq_items(self) -> Vec<Self::Item> {
self.base.seq_items().into_iter().enumerate().collect()
}
fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
self.base
.seq_items_window(skip, take)
.into_iter()
.enumerate()
.map(|(offset, item)| (skip + offset, item))
.collect()
}
/// # Why this stays sequential (the logical-index boundary)
///
/// Pairing an item with its logical index needs the count of items that
/// precede it in the whole stream, and the non-indexed consumer protocol
/// does not carry one. `Consumer::split_at` receives the *source's* split
/// point — `left.len()` at the source being divided — which equals the
/// logical offset only when nothing between the source and this adapter
/// changes the element count. A `filter` below invalidates it, and the
/// consumer cannot tell the two cases apart, so a shard handed that number
/// as a base index would silently emit wrong indices for exactly the chains
/// where it matters. No consumer in the tree reads the index today, which
/// is why the mismatch is currently latent rather than a live defect.
///
/// Supplying a true logical offset means an indexed producer boundary that
/// knows each shard's position in the logical stream — the change recorded
/// as the indexed adapter model in the Rayon adapter surface audit, not a
/// consumer this adapter can push itself into. `positions`,
/// `Map::positions`, and the borrowed position stream stay sequential for
/// this same reason, as do `take`, `skip`, and `step_by`, whose retained
/// items are a function of that same absent offset.
fn drive<C, R>(self, consumer: C) -> R
where
C: Consumer<Self::Item, Result = R> + Send + Sync,
R: Send,
{
consumer.consume(VecParIter::new(self.seq_items()))
}
}
/// Copied adapter with standard reference-copy semantics.
pub struct Copied<I> {
pub(super) base: I,
}
impl<I> Copied<I> {
pub(crate) fn new(base: I) -> Self {
Self { base }
}
}
impl<'data, I, T> ParallelIterator for Copied<I>
where
I: ParallelIterator<Item = &'data T>,
T: Copy + Send + Sync + 'data + 'static,
{
type Item = T;
fn seq_items(self) -> Vec<Self::Item> {
self.base.seq_items().into_iter().copied().collect()
}
fn seq_iter(self) -> impl Iterator<Item = Self::Item> {
self.base.seq_iter().copied()
}
fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
where
FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
{
self.base
.seq_try_fold(init, move |accumulator, item| fold_fn(accumulator, *item))
}
fn drive<C, R>(self, consumer: C) -> R
where
C: Consumer<Self::Item, Result = R> + Send + Sync,
R: Send,
{
// Push the copy into the consumer and drive the base, the way `Map`
// does. Materializing `seq_items()` first collected the whole stream
// into one vector before any split, which discarded the borrowed
// source's zero-copy split for every chain containing `copied()`.
self.base
.drive(MapConsumer::new(consumer, |item: &'data T| *item))
}
}
/// Cloned adapter with standard reference-clone semantics.
pub struct Cloned<I> {
pub(super) base: I,
}
impl<I> Cloned<I> {
pub(crate) fn new(base: I) -> Self {
Self { base }
}
}
impl<'data, I, T> ParallelIterator for Cloned<I>
where
I: ParallelIterator<Item = &'data T>,
T: Clone + Send + Sync + 'data + 'static,
{
type Item = T;
fn seq_items(self) -> Vec<Self::Item> {
self.base.seq_items().into_iter().cloned().collect()
}
fn seq_iter(self) -> impl Iterator<Item = Self::Item> {
self.base.seq_iter().cloned()
}
fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
where
FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
{
self.base.seq_try_fold(init, move |accumulator, item| {
fold_fn(accumulator, item.clone())
})
}
fn drive<C, R>(self, consumer: C) -> R
where
C: Consumer<Self::Item, Result = R> + Send + Sync,
R: Send,
{
// Push the clone into the consumer, as `Copied` does above.
self.base
.drive(MapConsumer::new(consumer, |item: &'data T| item.clone()))
}
}