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