moirai_iter/async_iter/
adapters.rs1use super::traits::{AsyncIterator, AsyncParallelIterator};
4use std::future::Future;
5
6pub 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
36pub 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
65pub 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
90pub 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
119pub 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
141pub 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
169pub 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> {}