Skip to main content

moirai_iter/parallel/adapters/
flat.rs

1use super::super::{Consumer, FlatMapConsumer, ParallelIterator};
2use std::ops::ControlFlow;
3
4/// Flat-map adapter with standard left-to-right flattening semantics.
5pub struct FlatMap<I, F> {
6    pub(super) base: I,
7    pub(super) flat_map_fn: F,
8}
9
10impl<I, F> FlatMap<I, F> {
11    pub(crate) fn new(base: I, flat_map_fn: F) -> Self {
12        Self { base, flat_map_fn }
13    }
14}
15
16impl<I, F, U> ParallelIterator for FlatMap<I, F>
17where
18    I: ParallelIterator,
19    F: Fn(I::Item) -> U + Send + Sync + Clone,
20    U: IntoIterator,
21    U::Item: Send + Sync + 'static,
22{
23    type Item = U::Item;
24
25    fn seq_items(self) -> Vec<Self::Item> {
26        self.base
27            .seq_items()
28            .into_iter()
29            .flat_map(self.flat_map_fn)
30            .collect()
31    }
32
33    fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
34    where
35        FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
36    {
37        let flat_map_fn = self.flat_map_fn;
38        self.base.seq_try_fold(init, move |accumulator, item| {
39            flat_map_fn(item)
40                .into_iter()
41                .try_fold(accumulator, &mut fold_fn)
42        })
43    }
44
45    fn drive<C, R>(self, consumer: C) -> R
46    where
47        C: Consumer<Self::Item, Result = R> + Send + Sync,
48        R: Send,
49    {
50        // Push the expansion into the consumer and drive the base, the way `Map`
51        // does. Materializing `seq_items()` first collected the whole flattened
52        // stream into one vector before any split, discarding the source's
53        // shards for every chain containing `flat_map()`. One input expanding to
54        // many outputs does not block the push: each expansion depends on its
55        // own item alone, so a shard produces exactly the sub-sequence a
56        // sequential pass over its range would, and shards combine in logical
57        // order.
58        self.base
59            .drive(FlatMapConsumer::new(consumer, self.flat_map_fn))
60    }
61}
62
63/// Flatten adapter with standard left-to-right nested stream semantics.
64pub struct Flatten<I> {
65    pub(super) base: I,
66}
67
68impl<I> Flatten<I> {
69    pub(crate) fn new(base: I) -> Self {
70        Self { base }
71    }
72}
73
74impl<I> ParallelIterator for Flatten<I>
75where
76    I: ParallelIterator,
77    I::Item: IntoIterator,
78    <I::Item as IntoIterator>::Item: Send + Sync + 'static,
79{
80    type Item = <I::Item as IntoIterator>::Item;
81
82    fn seq_items(self) -> Vec<Self::Item> {
83        self.base.seq_items().into_iter().flatten().collect()
84    }
85
86    fn seq_try_fold<Acc, B, FoldFn>(self, init: Acc, mut fold_fn: FoldFn) -> ControlFlow<B, Acc>
87    where
88        FoldFn: FnMut(Acc, Self::Item) -> ControlFlow<B, Acc>,
89    {
90        self.base.seq_try_fold(init, move |accumulator, item| {
91            item.into_iter().try_fold(accumulator, &mut fold_fn)
92        })
93    }
94
95    fn drive<C, R>(self, consumer: C) -> R
96    where
97        C: Consumer<Self::Item, Result = R> + Send + Sync,
98        R: Send,
99    {
100        // Flattening is `flat_map` with the identity expansion, so it reuses
101        // that consumer rather than duplicating the split and combine
102        // forwarding.
103        self.base
104            .drive(FlatMapConsumer::new(consumer, |item: I::Item| item))
105    }
106}