Skip to main content

moirai_iter/async_iter/
consumers.rs

1//! Consumer futures driving async iterators to completion.
2//!
3//! Each terminal drives the per-item user futures *cooperatively*: it holds the
4//! materialized item list plus the currently in-flight user future as state, and
5//! on every [`Future::poll`] it polls that in-flight future. On `Ready` it folds
6//! the result and advances to the next item; on `Pending` it returns `Pending`
7//! so the outer waker propagates. No terminal blocks the executor.
8//!
9//! The user closures return unnameable `impl Future` types, so the in-flight
10//! future is type-erased into a `Pin<Box<dyn Future>>` — one allocation per item.
11//! The alternative (threading a generic future type through the trait surface)
12//! would churn every caller; boxing keeps the public shape stable.
13
14use std::future::Future;
15use std::marker::PhantomData;
16use std::pin::Pin;
17use std::task::{Context, Poll};
18
19use super::traits::AsyncIterator;
20
21/// Type-erased in-flight user future.
22type ItemFuture<O> = Pin<Box<dyn Future<Output = O> + Send>>;
23
24/// Async for_each operation.
25pub struct AsyncForEach<I, F>
26where
27    I: AsyncIterator,
28{
29    /// Remaining items in forward order, consumed from the back for O(1) pops.
30    items: Option<Vec<I::Item>>,
31    func: F,
32    in_flight: Option<ItemFuture<()>>,
33}
34
35impl<I, F> AsyncForEach<I, F>
36where
37    I: AsyncIterator,
38{
39    pub(super) fn new(iter: I, func: F) -> Self
40    where
41        I: Sized,
42    {
43        let mut items = iter.into_vec();
44        // Reverse once so `pop` yields items in original order.
45        items.reverse();
46        Self {
47            items: Some(items),
48            func,
49            in_flight: None,
50        }
51    }
52}
53
54impl<I, F, Fut> Future for AsyncForEach<I, F>
55where
56    I: AsyncIterator,
57    F: Fn(I::Item) -> Fut + Send + Sync,
58    Fut: Future<Output = ()> + Send + 'static,
59    I: Unpin,
60    F: Unpin,
61    I::Item: Unpin,
62{
63    type Output = ();
64
65    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
66        let this = self.as_mut().get_mut();
67        loop {
68            if let Some(fut) = this.in_flight.as_mut() {
69                match fut.as_mut().poll(cx) {
70                    Poll::Ready(()) => this.in_flight = None,
71                    Poll::Pending => return Poll::Pending,
72                }
73            }
74            let items = this
75                .items
76                .as_mut()
77                .expect("async for_each polled after completion");
78            match items.pop() {
79                Some(item) => this.in_flight = Some(Box::pin((this.func)(item))),
80                None => {
81                    this.items = None;
82                    return Poll::Ready(());
83                }
84            }
85        }
86    }
87}
88
89/// Async collect operation.
90///
91/// `collect` performs no per-item async work — items are materialized by the
92/// source and extended into the target collection — so it completes in a single
93/// poll without blocking.
94pub struct AsyncCollect<I, C> {
95    iter: Option<I>,
96    _phantom: PhantomData<C>,
97}
98
99impl<I, C> AsyncCollect<I, C> {
100    pub(super) fn new(iter: I) -> Self {
101        Self {
102            iter: Some(iter),
103            _phantom: PhantomData,
104        }
105    }
106}
107
108impl<I, C> Future for AsyncCollect<I, C>
109where
110    I: AsyncIterator,
111    C: Default + Extend<I::Item> + Send,
112    I: Unpin,
113    C: Unpin,
114{
115    type Output = C;
116
117    fn poll(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Self::Output> {
118        let this = self.as_mut().get_mut();
119        let mut collection = C::default();
120        if let Some(iter) = this.iter.take() {
121            collection.extend(iter.into_vec());
122        }
123        Poll::Ready(collection)
124    }
125}
126
127/// Async fold operation.
128pub struct AsyncFold<I, T, F>
129where
130    I: AsyncIterator,
131{
132    items: Option<Vec<I::Item>>,
133    accumulator: Option<T>,
134    fold_fn: F,
135    in_flight: Option<ItemFuture<T>>,
136}
137
138impl<I, T, F> AsyncFold<I, T, F>
139where
140    I: AsyncIterator,
141{
142    pub(super) fn new(iter: I, init: T, fold_fn: F) -> Self
143    where
144        I: Sized,
145    {
146        let mut items = iter.into_vec();
147        items.reverse();
148        Self {
149            items: Some(items),
150            accumulator: Some(init),
151            fold_fn,
152            in_flight: None,
153        }
154    }
155}
156
157impl<I, T, F, Fut> Future for AsyncFold<I, T, F>
158where
159    I: AsyncIterator,
160    F: Fn(T, I::Item) -> Fut + Send + Sync,
161    Fut: Future<Output = T> + Send + 'static,
162    T: Send + Unpin,
163    I: Unpin,
164    F: Unpin,
165    I::Item: Unpin,
166{
167    type Output = T;
168
169    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
170        let this = self.as_mut().get_mut();
171        loop {
172            if let Some(fut) = this.in_flight.as_mut() {
173                match fut.as_mut().poll(cx) {
174                    Poll::Ready(acc) => {
175                        this.accumulator = Some(acc);
176                        this.in_flight = None;
177                    }
178                    Poll::Pending => return Poll::Pending,
179                }
180            }
181            let items = this
182                .items
183                .as_mut()
184                .expect("async fold polled after completion");
185            match items.pop() {
186                Some(item) => {
187                    let acc = this
188                        .accumulator
189                        .take()
190                        .expect("fold accumulator present between items");
191                    this.in_flight = Some(Box::pin((this.fold_fn)(acc, item)));
192                }
193                None => {
194                    this.items = None;
195                    return Poll::Ready(
196                        this.accumulator
197                            .take()
198                            .expect("fold accumulator present at completion"),
199                    );
200                }
201            }
202        }
203    }
204}
205
206/// Async reduce operation.
207pub struct AsyncReduce<I, F>
208where
209    I: AsyncIterator,
210{
211    items: Option<Vec<I::Item>>,
212    accumulator: Option<I::Item>,
213    reduce_fn: F,
214    in_flight: Option<ItemFuture<I::Item>>,
215}
216
217impl<I, F> AsyncReduce<I, F>
218where
219    I: AsyncIterator,
220{
221    pub(super) fn new(iter: I, reduce_fn: F) -> Self
222    where
223        I: Sized,
224    {
225        let mut items = iter.into_vec();
226        items.reverse();
227        // Seed the accumulator with the first logical item (last after reverse).
228        let accumulator = items.pop();
229        Self {
230            items: Some(items),
231            accumulator,
232            reduce_fn,
233            in_flight: None,
234        }
235    }
236}
237
238impl<I, F, Fut> Future for AsyncReduce<I, F>
239where
240    I: AsyncIterator,
241    F: Fn(I::Item, I::Item) -> Fut + Send + Sync,
242    Fut: Future<Output = I::Item> + Send + 'static,
243    I: Unpin,
244    F: Unpin,
245    I::Item: Unpin,
246{
247    type Output = Option<I::Item>;
248
249    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
250        let this = self.as_mut().get_mut();
251        loop {
252            if let Some(fut) = this.in_flight.as_mut() {
253                match fut.as_mut().poll(cx) {
254                    Poll::Ready(acc) => {
255                        this.accumulator = Some(acc);
256                        this.in_flight = None;
257                    }
258                    Poll::Pending => return Poll::Pending,
259                }
260            }
261            let items = this
262                .items
263                .as_mut()
264                .expect("async reduce polled after completion");
265            match items.pop() {
266                Some(item) => {
267                    let acc = this
268                        .accumulator
269                        .take()
270                        .expect("reduce accumulator present once seeded");
271                    this.in_flight = Some(Box::pin((this.reduce_fn)(acc, item)));
272                }
273                None => {
274                    this.items = None;
275                    // Empty input yields `None`; otherwise the folded accumulator.
276                    return Poll::Ready(this.accumulator.take());
277                }
278            }
279        }
280    }
281}