Skip to main content

moirai_iter/async_iter/
parallel.rs

1//! Parallel async operations with controlled concurrency.
2
3use futures::stream::{self, StreamExt};
4use std::future::Future;
5
6use super::traits::AsyncIterator;
7
8/// Parallel async map with concurrency control
9pub struct ParAsyncMap<I, F> {
10    iter: I,
11    concurrency: usize,
12    map_fn: F,
13}
14
15impl<I, F> ParAsyncMap<I, F> {
16    pub(super) fn new(iter: I, concurrency: usize, map_fn: F) -> Self {
17        Self {
18            iter,
19            concurrency,
20            map_fn,
21        }
22    }
23}
24
25impl<I, F, Fut, R> AsyncIterator for ParAsyncMap<I, F>
26where
27    I: AsyncIterator,
28    F: Fn(I::Item) -> Fut + Send + Sync,
29    Fut: Future<Output = R> + Send,
30    R: Send,
31{
32    type Item = R;
33
34    fn into_vec(self) -> Vec<Self::Item> {
35        let concurrency = self.concurrency.max(1);
36        let map_fn = self.map_fn;
37        let items = self.iter.into_vec();
38        futures::executor::block_on(async move {
39            stream::iter(items)
40                .map(|item| {
41                    let map_fn = &map_fn;
42                    async move { map_fn(item).await }
43                })
44                .buffered(concurrency)
45                .collect()
46                .await
47        })
48    }
49}
50
51/// Parallel async filter with concurrency control
52pub struct ParAsyncFilter<I, F> {
53    iter: I,
54    concurrency: usize,
55    filter_fn: F,
56}
57
58impl<I, F> ParAsyncFilter<I, F> {
59    pub(super) fn new(iter: I, concurrency: usize, filter_fn: F) -> Self {
60        Self {
61            iter,
62            concurrency,
63            filter_fn,
64        }
65    }
66}
67
68impl<I, F, Fut> AsyncIterator for ParAsyncFilter<I, F>
69where
70    I: AsyncIterator,
71    F: Fn(&I::Item) -> Fut + Send + Sync,
72    Fut: Future<Output = bool> + Send,
73{
74    type Item = I::Item;
75
76    fn into_vec(self) -> Vec<Self::Item> {
77        let concurrency = self.concurrency.max(1);
78        let filter_fn = self.filter_fn;
79        let items = self.iter.into_vec();
80        futures::executor::block_on(async move {
81            stream::iter(items)
82                .map(|item| {
83                    let filter_fn = &filter_fn;
84                    async move {
85                        let keep = filter_fn(&item).await;
86                        (item, keep)
87                    }
88                })
89                .buffered(concurrency)
90                .filter_map(|(item, keep)| async move { keep.then_some(item) })
91                .collect()
92                .await
93        })
94    }
95}
96
97/// Parallel async for_each with concurrency control.
98///
99/// Returns a real `Future` driven by the caller's runtime — it never blocks the
100/// executor. `for_each` is order-independent, so it uses `buffer_unordered`
101/// (no head-of-line blocking) while still keeping at most `concurrency` item
102/// futures in flight.
103pub(super) async fn for_each<I, F, Fut>(iter: I, concurrency: usize, func: F)
104where
105    I: AsyncIterator,
106    F: Fn(I::Item) -> Fut + Send + Sync,
107    Fut: Future<Output = ()> + Send,
108{
109    let items = iter.into_vec();
110    stream::iter(items)
111        .map(func)
112        .buffer_unordered(concurrency.max(1))
113        .for_each(|()| async {})
114        .await;
115}