Skip to main content

moirai_iter/async_iter/
adapters.rs

1//! Adapters transforming async iterators.
2
3use super::traits::{AsyncIterator, AsyncParallelIterator};
4use std::future::Future;
5
6/// Async map operation
7pub struct AsyncMap<I, F> {
8    iter: I,
9    map_fn: F,
10}
11
12impl<I, F> AsyncMap<I, F> {
13    pub(super) fn new(iter: I, map_fn: F) -> Self {
14        Self { iter, map_fn }
15    }
16}
17
18impl<I, F, Fut, R> AsyncIterator for AsyncMap<I, F>
19where
20    I: AsyncIterator,
21    F: Fn(I::Item) -> Fut + Send + Sync,
22    Fut: Future<Output = R> + Send,
23    R: Send,
24{
25    type Item = R;
26
27    fn into_vec(self) -> Vec<Self::Item> {
28        self.iter
29            .into_vec()
30            .into_iter()
31            .map(|item| futures::executor::block_on((self.map_fn)(item)))
32            .collect()
33    }
34}
35
36/// Async filter operation
37pub struct AsyncFilter<I, F> {
38    iter: I,
39    filter_fn: F,
40}
41
42impl<I, F> AsyncFilter<I, F> {
43    pub(super) fn new(iter: I, filter_fn: F) -> Self {
44        Self { iter, filter_fn }
45    }
46}
47
48impl<I, F, Fut> AsyncIterator for AsyncFilter<I, F>
49where
50    I: AsyncIterator,
51    F: Fn(&I::Item) -> Fut + Send + Sync,
52    Fut: Future<Output = bool> + Send,
53{
54    type Item = I::Item;
55
56    fn into_vec(self) -> Vec<Self::Item> {
57        self.iter
58            .into_vec()
59            .into_iter()
60            .filter(|item| futures::executor::block_on((self.filter_fn)(item)))
61            .collect()
62    }
63}
64
65/// Async take operation with prefix-bounded value semantics.
66pub struct AsyncTake<I> {
67    iter: I,
68    count: usize,
69}
70
71impl<I> AsyncTake<I> {
72    pub(super) fn new(iter: I, count: usize) -> Self {
73        Self { iter, count }
74    }
75}
76
77impl<I> AsyncIterator for AsyncTake<I>
78where
79    I: AsyncIterator,
80{
81    type Item = I::Item;
82
83    fn into_vec(self) -> Vec<Self::Item> {
84        let mut items = self.iter.into_vec();
85        items.truncate(self.count);
86        items
87    }
88}
89
90/// Async skip operation with prefix-discarding value semantics.
91pub struct AsyncSkip<I> {
92    iter: I,
93    count: usize,
94}
95
96impl<I> AsyncSkip<I> {
97    pub(super) fn new(iter: I, count: usize) -> Self {
98        Self { iter, count }
99    }
100}
101
102impl<I> AsyncIterator for AsyncSkip<I>
103where
104    I: AsyncIterator,
105{
106    type Item = I::Item;
107
108    fn into_vec(self) -> Vec<Self::Item> {
109        let mut items = self.iter.into_vec();
110        if self.count >= items.len() {
111            Vec::new()
112        } else {
113            items.drain(..self.count);
114            items
115        }
116    }
117}
118
119/// Async enumerate operation with zero-based logical positions.
120pub struct AsyncEnumerate<I> {
121    iter: I,
122}
123
124impl<I> AsyncEnumerate<I> {
125    pub(super) fn new(iter: I) -> Self {
126        Self { iter }
127    }
128}
129
130impl<I> AsyncIterator for AsyncEnumerate<I>
131where
132    I: AsyncIterator,
133{
134    type Item = (usize, I::Item);
135
136    fn into_vec(self) -> Vec<Self::Item> {
137        self.iter.into_vec().into_iter().enumerate().collect()
138    }
139}
140
141/// Async zip operation with shortest-input semantics.
142pub struct AsyncZip<I, J> {
143    left: I,
144    right: J,
145}
146
147impl<I, J> AsyncZip<I, J> {
148    pub(super) fn new(left: I, right: J) -> Self {
149        Self { left, right }
150    }
151}
152
153impl<I, J> AsyncIterator for AsyncZip<I, J>
154where
155    I: AsyncIterator,
156    J: AsyncIterator,
157{
158    type Item = (I::Item, J::Item);
159
160    fn into_vec(self) -> Vec<Self::Item> {
161        self.left
162            .into_vec()
163            .into_iter()
164            .zip(self.right.into_vec())
165            .collect()
166    }
167}
168
169/// Adapter to make async iterators work with parallel processing
170pub struct AsyncParallelAdapter<I> {
171    iter: I,
172}
173
174impl<I> AsyncParallelAdapter<I> {
175    pub(super) fn new(iter: I) -> Self {
176        Self { iter }
177    }
178}
179
180impl<I: AsyncIterator> AsyncIterator for AsyncParallelAdapter<I> {
181    type Item = I::Item;
182
183    fn into_vec(self) -> Vec<Self::Item> {
184        self.iter.into_vec()
185    }
186}
187
188impl<I: AsyncIterator> AsyncParallelIterator for AsyncParallelAdapter<I> {}