Skip to main content

moirai_iter/async_iter/
traits.rs

1//! Core async iterator traits.
2
3use 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
11/// Core async iterator trait for async/await compatible iteration
12pub trait AsyncIterator: Send {
13    /// The type of items yielded by this async iterator
14    type Item: Send;
15
16    /// Materialize the iterator into its logical item sequence.
17    fn into_vec(self) -> Vec<Self::Item>
18    where
19        Self: Sized;
20
21    /// Async map operation that transforms each element
22    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    /// Async filter operation
33    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    /// Retain at most `count` items from the logical async stream prefix.
43    fn take(self, count: usize) -> AsyncTake<Self>
44    where
45        Self: Sized,
46    {
47        AsyncTake::new(self, count)
48    }
49
50    /// Discard `count` items from the logical async stream prefix.
51    fn skip(self, count: usize) -> AsyncSkip<Self>
52    where
53        Self: Sized,
54    {
55        AsyncSkip::new(self, count)
56    }
57
58    /// Pair each item with its zero-based logical stream position.
59    fn enumerate(self) -> AsyncEnumerate<Self>
60    where
61        Self: Sized,
62    {
63        AsyncEnumerate::new(self)
64    }
65
66    /// Pair items with another async iterator, stopping at the shorter input.
67    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    /// Async for_each operation with side effects
76    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    /// Collect into a vector asynchronously
86    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    /// Fold operation with async function
95    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    /// Reduce operation for async iterators
106    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    /// Convert to parallel iterator for hybrid processing
116    fn into_parallel(self) -> AsyncParallelAdapter<Self>
117    where
118        Self: Sized,
119    {
120        AsyncParallelAdapter::new(self)
121    }
122}
123
124/// Parallel async iterator for CPU+async hybrid workloads
125pub trait AsyncParallelIterator: AsyncIterator {
126    /// Execute async operations in parallel with controlled concurrency
127    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    /// Parallel async filter with concurrency control
138    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    /// Execute side effects in parallel with async operations.
148    ///
149    /// Returns a real `Future` driven by the caller's runtime; unlike a blocking
150    /// consumer it never blocks the async executor.
151    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
161/// Trait for converting types into async iterators
162pub trait IntoAsyncIterator {
163    /// The type of items yielded by this async iterator
164    type Item: Send;
165    /// The resulting async iterator type
166    type IntoAsyncIter: AsyncIterator<Item = Self::Item>;
167
168    /// Convert the type into an async iterator
169    fn into_async_iter(self) -> Self::IntoAsyncIter;
170}