1use std::future::Future;
15use std::marker::PhantomData;
16use std::pin::Pin;
17use std::task::{Context, Poll};
18
19use super::traits::AsyncIterator;
20
21type ItemFuture<O> = Pin<Box<dyn Future<Output = O> + Send>>;
23
24pub struct AsyncForEach<I, F>
26where
27 I: AsyncIterator,
28{
29 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 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
89pub 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
127pub 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
206pub 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 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 return Poll::Ready(this.accumulator.take());
277 }
278 }
279 }
280 }
281}