moirai_iter/parallel/adapters/
flat.rs1use super::super::{Consumer, ParallelIterator, VecParIter};
2
3pub 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
41pub 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}