Skip to main content

moirai_iter/parallel/adapters/
flat.rs

1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3/// Flat-map adapter with standard left-to-right flattening semantics.
4pub struct FlatMap<I, F> {
5    pub(super) base: I,
6    pub(super) flat_map_fn: F,
7}
8
9impl<I, F> FlatMap<I, F> {
10    pub(crate) fn new(base: I, flat_map_fn: F) -> Self {
11        Self { base, flat_map_fn }
12    }
13}
14
15impl<I, F, U> ParallelIterator for FlatMap<I, F>
16where
17    I: ParallelIterator,
18    F: Fn(I::Item) -> U + Send + Sync + Clone,
19    U: IntoIterator,
20    U::Item: Send + Sync + 'static,
21{
22    type Item = U::Item;
23
24    fn seq_items(self) -> Vec<Self::Item> {
25        self.base
26            .seq_items()
27            .into_iter()
28            .flat_map(self.flat_map_fn)
29            .collect()
30    }
31
32    fn drive<C, R>(self, consumer: C) -> R
33    where
34        C: Consumer<Self::Item, Result = R> + Send + Sync,
35        R: Send,
36    {
37        consumer.consume(VecParIter::new(self.seq_items()))
38    }
39}
40
41/// Flatten adapter with standard left-to-right nested stream semantics.
42pub struct Flatten<I> {
43    pub(super) base: I,
44}
45
46impl<I> Flatten<I> {
47    pub(crate) fn new(base: I) -> Self {
48        Self { base }
49    }
50}
51
52impl<I> ParallelIterator for Flatten<I>
53where
54    I: ParallelIterator,
55    I::Item: IntoIterator,
56    <I::Item as IntoIterator>::Item: Send + Sync + 'static,
57{
58    type Item = <I::Item as IntoIterator>::Item;
59
60    fn seq_items(self) -> Vec<Self::Item> {
61        self.base.seq_items().into_iter().flatten().collect()
62    }
63
64    fn drive<C, R>(self, consumer: C) -> R
65    where
66        C: Consumer<Self::Item, Result = R> + Send + Sync,
67        R: Send,
68    {
69        consumer.consume(VecParIter::new(self.seq_items()))
70    }
71}