Skip to main content

moirai_iter/parallel/adapters/
ref_ops.rs

1use super::super::{Consumer, MapConsumer, ParallelIterator, VecParIter};
2use std::ops::ControlFlow;
3
4/// Enumerate adapter for value-semantic index pairing.
5pub struct Enumerate<I> {
6    pub(super) base: I,
7}
8
9impl<I> Enumerate<I> {
10    pub(crate) fn new(base: I) -> Self {
11        Self { base }
12    }
13}
14
15impl<I> ParallelIterator for Enumerate<I>
16where
17    I: ParallelIterator,
18    I::Item: Sync + 'static,
19{
20    type Item = (usize, I::Item);
21
22    fn seq_items(self) -> Vec<Self::Item> {
23        self.base.seq_items().into_iter().enumerate().collect()
24    }
25
26    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
27        self.base
28            .seq_items_window(skip, take)
29            .into_iter()
30            .enumerate()
31            .map(|(offset, item)| (skip + offset, item))
32            .collect()
33    }
34
35    /// # Why this stays sequential (the logical-index boundary)
36    ///
37    /// Pairing an item with its logical index needs the count of items that
38    /// precede it in the whole stream, and the non-indexed consumer protocol
39    /// does not carry one. `Consumer::split_at` receives the *source's* split
40    /// point — `left.len()` at the source being divided — which equals the
41    /// logical offset only when nothing between the source and this adapter
42    /// changes the element count. A `filter` below invalidates it, and the
43    /// consumer cannot tell the two cases apart, so a shard handed that number
44    /// as a base index would silently emit wrong indices for exactly the chains
45    /// where it matters. No consumer in the tree reads the index today, which
46    /// is why the mismatch is currently latent rather than a live defect.
47    ///
48    /// Supplying a true logical offset means an indexed producer boundary that
49    /// knows each shard's position in the logical stream — the change recorded
50    /// as the indexed adapter model in the Rayon adapter surface audit, not a
51    /// consumer this adapter can push itself into. `positions`,
52    /// `Map::positions`, and the borrowed position stream stay sequential for
53    /// this same reason, as do `take`, `skip`, and `step_by`, whose retained
54    /// items are a function of that same absent offset.
55    fn drive<C, R>(self, consumer: C) -> R
56    where
57        C: Consumer<Self::Item, Result = R> + Send + Sync,
58        R: Send,
59    {
60        consumer.consume(VecParIter::new(self.seq_items()))
61    }
62}
63
64/// Copied adapter with standard reference-copy semantics.
65pub struct Copied<I> {
66    pub(super) base: I,
67}
68
69impl<I> Copied<I> {
70    pub(crate) fn new(base: I) -> Self {
71        Self { base }
72    }
73}
74
75impl<'data, I, T> ParallelIterator for Copied<I>
76where
77    I: ParallelIterator<Item = &'data T>,
78    T: Copy + Send + Sync + 'data + 'static,
79{
80    type Item = T;
81
82    fn seq_items(self) -> Vec<Self::Item> {
83        self.base.seq_items().into_iter().copied().collect()
84    }
85
86    fn seq_iter(self) -> impl Iterator<Item = Self::Item> {
87        self.base.seq_iter().copied()
88    }
89
90    fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
91    where
92        FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
93    {
94        self.base
95            .seq_try_fold(init, move |accumulator, item| fold_fn(accumulator, *item))
96    }
97
98    fn drive<C, R>(self, consumer: C) -> R
99    where
100        C: Consumer<Self::Item, Result = R> + Send + Sync,
101        R: Send,
102    {
103        // Push the copy into the consumer and drive the base, the way `Map`
104        // does. Materializing `seq_items()` first collected the whole stream
105        // into one vector before any split, which discarded the borrowed
106        // source's zero-copy split for every chain containing `copied()`.
107        self.base
108            .drive(MapConsumer::new(consumer, |item: &'data T| *item))
109    }
110}
111
112/// Cloned adapter with standard reference-clone semantics.
113pub struct Cloned<I> {
114    pub(super) base: I,
115}
116
117impl<I> Cloned<I> {
118    pub(crate) fn new(base: I) -> Self {
119        Self { base }
120    }
121}
122
123impl<'data, I, T> ParallelIterator for Cloned<I>
124where
125    I: ParallelIterator<Item = &'data T>,
126    T: Clone + Send + Sync + 'data + 'static,
127{
128    type Item = T;
129
130    fn seq_items(self) -> Vec<Self::Item> {
131        self.base.seq_items().into_iter().cloned().collect()
132    }
133
134    fn seq_iter(self) -> impl Iterator<Item = Self::Item> {
135        self.base.seq_iter().cloned()
136    }
137
138    fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
139    where
140        FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
141    {
142        self.base.seq_try_fold(init, move |accumulator, item| {
143            fold_fn(accumulator, item.clone())
144        })
145    }
146
147    fn drive<C, R>(self, consumer: C) -> R
148    where
149        C: Consumer<Self::Item, Result = R> + Send + Sync,
150        R: Send,
151    {
152        // Push the clone into the consumer, as `Copied` does above.
153        self.base
154            .drive(MapConsumer::new(consumer, |item: &'data T| item.clone()))
155    }
156}