Skip to main content

moirai_iter/parallel/adapters/
map.rs

1use super::super::{Consumer, MapConsumer, ParallelIterator, TryStreamItem, VecParIter, fallible};
2use std::ops::ControlFlow;
3
4/// Map adapter for parallel iterators.
5pub struct Map<I, F> {
6    pub(super) base: I,
7    pub(super) map_fn: F,
8}
9
10impl<I, F> Map<I, F> {
11    pub(crate) fn new(base: I, map_fn: F) -> Self {
12        Self { base, map_fn }
13    }
14
15    /// Reduce a mapped fallible stream without materializing mapped items first.
16    pub fn try_reduce_with<ReduceFn, R>(self, reduce_fn: ReduceFn) -> Option<R>
17    where
18        I: ParallelIterator,
19        F: Fn(I::Item) -> R + Send + Sync + Clone,
20        R: TryStreamItem,
21        ReduceFn: Fn(<R as TryStreamItem>::Output, <R as TryStreamItem>::Output) -> R
22            + Send
23            + Sync
24            + Clone,
25    {
26        fallible::try_reduce_with_items(
27            self.base.seq_items().into_iter().map(self.map_fn),
28            reduce_fn,
29        )
30    }
31
32    /// Return mapped logical indices without materializing the mapped stream.
33    pub fn positions<Predicate, R>(
34        self,
35        predicate: Predicate,
36    ) -> super::position::MapPositions<I, F, Predicate>
37    where
38        I: ParallelIterator,
39        F: Fn(I::Item) -> R + Send + Sync + Clone,
40        Predicate: Fn(R) -> bool + Send + Sync + Clone,
41        R: Send,
42    {
43        super::position::MapPositions::new(self.base, self.map_fn, predicate)
44    }
45}
46
47impl<I, F, R> ParallelIterator for Map<I, F>
48where
49    I: ParallelIterator,
50    F: Fn(I::Item) -> R + Send + Sync + Clone,
51    R: Send,
52{
53    type Item = R;
54
55    fn seq_items(self) -> Vec<Self::Item> {
56        self.base.seq_items().into_iter().map(self.map_fn).collect()
57    }
58
59    fn seq_iter(self) -> impl Iterator<Item = Self::Item> {
60        self.base.seq_iter().map(self.map_fn)
61    }
62
63    fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
64    where
65        FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
66    {
67        let map_fn = self.map_fn;
68        self.base.seq_try_fold(init, move |accumulator, item| {
69            fold_fn(accumulator, map_fn(item))
70        })
71    }
72
73    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
74        self.base
75            .seq_items_window(skip, take)
76            .into_iter()
77            .map(self.map_fn)
78            .collect()
79    }
80
81    fn seq_items_reversed(self) -> Vec<Self::Item> {
82        self.base
83            .seq_items_reversed()
84            .into_iter()
85            .map(self.map_fn)
86            .collect()
87    }
88
89    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
90        self.base
91            .seq_items_reversed_prefix(count)
92            .into_iter()
93            .map(self.map_fn)
94            .collect()
95    }
96
97    fn drive<C, R2>(self, consumer: C) -> R2
98    where
99        C: Consumer<Self::Item, Result = R2> + Send + Sync,
100        R2: Send,
101    {
102        self.base.drive(MapConsumer::new(consumer, self.map_fn))
103    }
104
105    fn position_first<P>(self, predicate: P) -> Option<usize>
106    where
107        P: Fn(Self::Item) -> bool + Send + Sync + Clone,
108    {
109        self.base
110            .seq_items()
111            .into_iter()
112            .map(self.map_fn)
113            .position(predicate)
114    }
115
116    fn position_any<P>(self, predicate: P) -> Option<usize>
117    where
118        P: Fn(Self::Item) -> bool + Send + Sync + Clone,
119    {
120        self.position_first(predicate)
121    }
122
123    fn position_last<P>(self, predicate: P) -> Option<usize>
124    where
125        P: Fn(Self::Item) -> bool + Send + Sync + Clone,
126    {
127        self.base
128            .seq_items()
129            .into_iter()
130            .map(self.map_fn)
131            .rposition(predicate)
132    }
133
134    fn try_reduce<Identity, ReduceFn, T, E>(
135        self,
136        identity: Identity,
137        reduce_fn: ReduceFn,
138    ) -> Result<T, E>
139    where
140        Identity: Fn() -> T + Send + Sync + Clone,
141        ReduceFn: Fn(T, T) -> Result<T, E> + Send + Sync + Clone,
142        R: Into<Result<T, E>>,
143        T: Send,
144        E: Send,
145    {
146        let mut accumulator = identity();
147        for item in self.base.seq_items().into_iter().map(self.map_fn) {
148            accumulator = reduce_fn(accumulator, item.into()?)?;
149        }
150        Ok(accumulator)
151    }
152}
153
154/// Map adapter with cloned per-operation state.
155pub struct MapWith<I, T, F> {
156    pub(super) base: I,
157    pub(super) init: T,
158    pub(super) map_fn: F,
159}
160
161impl<I, T, F> MapWith<I, T, F> {
162    pub(crate) fn new(base: I, init: T, map_fn: F) -> Self {
163        Self { base, init, map_fn }
164    }
165}
166
167impl<I, T, F, R> ParallelIterator for MapWith<I, T, F>
168where
169    I: ParallelIterator,
170    T: Send + Clone,
171    F: Fn(&mut T, I::Item) -> R + Send + Sync + Clone,
172    R: Send + Sync + 'static,
173{
174    type Item = R;
175
176    fn seq_items(self) -> Vec<Self::Item> {
177        let mut state = self.init;
178        self.base
179            .seq_items()
180            .into_iter()
181            .map(|item| (self.map_fn)(&mut state, item))
182            .collect()
183    }
184
185    /// # Why this stays sequential
186    ///
187    /// One state value threads through the whole stream, so `map_fn` observes
188    /// every prior item's effect on it. Giving each shard its own clone is a
189    /// different contract, not a parallelisation of this one — the same reason
190    /// `for_each_with` and `try_for_each_with` stay sequential.
191    fn drive<C, R2>(self, consumer: C) -> R2
192    where
193        C: Consumer<Self::Item, Result = R2> + Send + Sync,
194        R2: Send,
195    {
196        consumer.consume(VecParIter::new(self.seq_items()))
197    }
198}
199
200/// Map adapter with lazily initialized state.
201pub struct MapInit<I, Init, F> {
202    pub(super) base: I,
203    pub(super) init: Init,
204    pub(super) map_fn: F,
205}
206
207impl<I, Init, F> MapInit<I, Init, F> {
208    pub(crate) fn new(base: I, init: Init, map_fn: F) -> Self {
209        Self { base, init, map_fn }
210    }
211}
212
213impl<I, Init, T, F, R> ParallelIterator for MapInit<I, Init, F>
214where
215    I: ParallelIterator,
216    Init: Fn() -> T + Send + Sync + Clone,
217    T: Send,
218    F: Fn(&mut T, I::Item) -> R + Send + Sync + Clone,
219    R: Send + Sync + 'static,
220{
221    type Item = R;
222
223    fn seq_items(self) -> Vec<Self::Item> {
224        let mut state = (self.init)();
225        self.base
226            .seq_items()
227            .into_iter()
228            .map(|item| (self.map_fn)(&mut state, item))
229            .collect()
230    }
231
232    /// # Why this stays sequential
233    ///
234    /// Threaded state, for the reason given on [`MapWith`].
235    fn drive<C, R2>(self, consumer: C) -> R2
236    where
237        C: Consumer<Self::Item, Result = R2> + Send + Sync,
238        R2: Send,
239    {
240        consumer.consume(VecParIter::new(self.seq_items()))
241    }
242}
243
244/// Update adapter that mutates each item before yielding it.
245pub struct Update<I, F> {
246    pub(super) base: I,
247    pub(super) update_fn: F,
248}
249
250impl<I, F> Update<I, F> {
251    pub(crate) fn new(base: I, update_fn: F) -> Self {
252        Self { base, update_fn }
253    }
254}
255
256impl<I, F> ParallelIterator for Update<I, F>
257where
258    I: ParallelIterator,
259    F: Fn(&mut I::Item) + Send + Sync + Clone,
260    I::Item: Sync + 'static,
261{
262    type Item = I::Item;
263
264    fn seq_items(self) -> Vec<Self::Item> {
265        self.base
266            .seq_items()
267            .into_iter()
268            .map(|mut item| {
269                (self.update_fn)(&mut item);
270                item
271            })
272            .collect()
273    }
274
275    fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
276    where
277        FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
278    {
279        let update_fn = self.update_fn;
280        self.base.seq_try_fold(init, move |accumulator, mut item| {
281            update_fn(&mut item);
282            fold_fn(accumulator, item)
283        })
284    }
285
286    fn drive<C, R>(self, consumer: C) -> R
287    where
288        C: Consumer<Self::Item, Result = R> + Send + Sync,
289        R: Send,
290    {
291        // One input, one output, no state between items: the same push `Map`
292        // makes. Materializing `seq_items()` first discarded the source's shards
293        // Every chain containing `update()` lost its source shards.
294        let update_fn = self.update_fn;
295        self.base
296            .drive(MapConsumer::new(consumer, move |mut item: I::Item| {
297                update_fn(&mut item);
298                item
299            }))
300    }
301}