1use std::future::Future;
4
5use super::adapters::{
6 AsyncEnumerate, AsyncFilter, AsyncMap, AsyncParallelAdapter, AsyncSkip, AsyncTake, AsyncZip,
7};
8use super::consumers::{AsyncCollect, AsyncFold, AsyncForEach, AsyncReduce};
9use super::parallel::{self, ParAsyncFilter, ParAsyncMap};
10
11pub trait AsyncIterator: Send {
13 type Item: Send;
15
16 fn into_vec(self) -> Vec<Self::Item>
18 where
19 Self: Sized;
20
21 fn map<F, Fut, R>(self, map_fn: F) -> AsyncMap<Self, F>
23 where
24 Self: Sized,
25 F: Fn(Self::Item) -> Fut + Send + Sync,
26 Fut: Future<Output = R> + Send,
27 R: Send,
28 {
29 AsyncMap::new(self, map_fn)
30 }
31
32 fn filter<F, Fut>(self, filter_fn: F) -> AsyncFilter<Self, F>
34 where
35 Self: Sized,
36 F: Fn(&Self::Item) -> Fut + Send + Sync,
37 Fut: Future<Output = bool> + Send,
38 {
39 AsyncFilter::new(self, filter_fn)
40 }
41
42 fn take(self, count: usize) -> AsyncTake<Self>
44 where
45 Self: Sized,
46 {
47 AsyncTake::new(self, count)
48 }
49
50 fn skip(self, count: usize) -> AsyncSkip<Self>
52 where
53 Self: Sized,
54 {
55 AsyncSkip::new(self, count)
56 }
57
58 fn enumerate(self) -> AsyncEnumerate<Self>
60 where
61 Self: Sized,
62 {
63 AsyncEnumerate::new(self)
64 }
65
66 fn zip<J>(self, other: J) -> AsyncZip<Self, J>
68 where
69 Self: Sized,
70 J: AsyncIterator,
71 {
72 AsyncZip::new(self, other)
73 }
74
75 fn for_each<F, Fut>(self, func: F) -> AsyncForEach<Self, F>
77 where
78 Self: Sized,
79 F: Fn(Self::Item) -> Fut + Send + Sync,
80 Fut: Future<Output = ()> + Send,
81 {
82 AsyncForEach::new(self, func)
83 }
84
85 fn collect<C>(self) -> AsyncCollect<Self, C>
87 where
88 Self: Sized,
89 C: Default + Extend<Self::Item> + Send,
90 {
91 AsyncCollect::new(self)
92 }
93
94 fn fold<T, F, Fut>(self, init: T, fold_fn: F) -> AsyncFold<Self, T, F>
96 where
97 Self: Sized,
98 F: Fn(T, Self::Item) -> Fut + Send + Sync,
99 Fut: Future<Output = T> + Send,
100 T: Send,
101 {
102 AsyncFold::new(self, init, fold_fn)
103 }
104
105 fn reduce<F, Fut>(self, reduce_fn: F) -> AsyncReduce<Self, F>
107 where
108 Self: Sized,
109 F: Fn(Self::Item, Self::Item) -> Fut + Send + Sync,
110 Fut: Future<Output = Self::Item> + Send,
111 {
112 AsyncReduce::new(self, reduce_fn)
113 }
114
115 fn into_parallel(self) -> AsyncParallelAdapter<Self>
117 where
118 Self: Sized,
119 {
120 AsyncParallelAdapter::new(self)
121 }
122}
123
124pub trait AsyncParallelIterator: AsyncIterator {
126 fn par_map<F, Fut, R>(self, concurrency: usize, map_fn: F) -> ParAsyncMap<Self, F>
128 where
129 Self: Sized,
130 F: Fn(Self::Item) -> Fut + Send + Sync,
131 Fut: Future<Output = R> + Send,
132 R: Send,
133 {
134 ParAsyncMap::new(self, concurrency, map_fn)
135 }
136
137 fn par_filter<F, Fut>(self, concurrency: usize, filter_fn: F) -> ParAsyncFilter<Self, F>
139 where
140 Self: Sized,
141 F: Fn(&Self::Item) -> Fut + Send + Sync,
142 Fut: Future<Output = bool> + Send,
143 {
144 ParAsyncFilter::new(self, concurrency, filter_fn)
145 }
146
147 fn par_for_each<F, Fut>(self, concurrency: usize, func: F) -> impl Future<Output = ()> + Send
152 where
153 Self: Sized,
154 F: Fn(Self::Item) -> Fut + Send + Sync,
155 Fut: Future<Output = ()> + Send,
156 {
157 parallel::for_each(self, concurrency, func)
158 }
159}
160
161pub trait IntoAsyncIterator {
163 type Item: Send;
165 type IntoAsyncIter: AsyncIterator<Item = Self::Item>;
167
168 fn into_async_iter(self) -> Self::IntoAsyncIter;
170}