Skip to main content

tuwunel_core/utils/stream/
ready.rs

1//! Synchronous combinator extensions to futures::Stream
2#![expect(clippy::type_complexity)]
3
4use futures::{
5	future::{FutureExt, Ready, ready},
6	stream::{All, Any, FilterMap, Fold, ForEach, Scan, Stream, StreamExt},
7};
8
9/// Adds synchronous counterparts to selected [`StreamExt`] combinators.
10///
11/// These methods accept ordinary closures for operations that do not need to
12/// await, avoiding an async block around iterator-like predicates and folds.
13pub trait ReadyExt<Item>
14where
15	Self: Stream<Item = Item> + Sized,
16{
17	/// Tests whether every item satisfies a synchronous predicate.
18	///
19	/// Evaluation stops at the first false result. An empty stream resolves to
20	/// true.
21	fn ready_all<F>(self, f: F) -> All<Self, Ready<bool>, impl FnMut(Item) -> Ready<bool>>
22	where
23		F: Fn(Item) -> bool;
24
25	/// Tests whether any item satisfies a synchronous predicate.
26	///
27	/// Evaluation stops at the first true result. An empty stream resolves to
28	/// false.
29	fn ready_any<F>(self, f: F) -> Any<Self, Ready<bool>, impl FnMut(Item) -> Ready<bool>>
30	where
31		F: Fn(Item) -> bool;
32
33	/// Finds the first item satisfying a synchronous predicate.
34	///
35	/// Items are tested in source order and the first match is returned. The
36	/// future resolves to `None` when the stream ends without a match.
37	fn ready_find<'a, F>(self, f: F) -> impl Future<Output = Option<Item>> + Send
38	where
39		Self: Send + Unpin + 'a,
40		F: Fn(&Item) -> bool + Send + 'a,
41		Item: Send;
42
43	/// Finds the first value produced by a synchronous mapping predicate.
44	///
45	/// Items are mapped in source order until `f` returns `Some`. The future
46	/// resolves to `None` when every item maps to absence.
47	fn ready_find_map<'a, F, U>(self, f: F) -> impl Future<Output = Option<U>> + Send
48	where
49		Self: Send + Unpin + 'a,
50		F: Fn(Item) -> Option<U> + Send + 'a,
51		U: Send;
52
53	/// Retains items accepted by a synchronous predicate.
54	///
55	/// The predicate borrows each item as it is polled. Accepted items preserve
56	/// their source order and value.
57	fn ready_filter<'a, F>(self, f: F) -> impl Stream<Item = Item> + Send + 'a
58	where
59		Self: Send + 'a,
60		Item: Send,
61		F: Fn(&Item) -> bool + Send + 'a;
62
63	/// Maps items synchronously while filtering absent results.
64	///
65	/// Values returned in `Some` are yielded in source order. Items mapped to
66	/// `None` are discarded.
67	fn ready_filter_map<F, U>(
68		self,
69		f: F,
70	) -> FilterMap<Self, Ready<Option<U>>, impl FnMut(Item) -> Ready<Option<U>>>
71	where
72		F: Fn(Item) -> Option<U>;
73
74	/// Folds items synchronously from an explicit initial accumulator.
75	///
76	/// `f` receives the accumulator and each item in source order. The future
77	/// resolves to the final accumulator after the stream ends.
78	fn ready_fold<T, F>(
79		self,
80		init: T,
81		f: F,
82	) -> Fold<Self, Ready<T>, T, impl FnMut(T, Item) -> Ready<T>>
83	where
84		F: Fn(T, Item) -> T;
85
86	/// Folds items synchronously from `T::default()`.
87	///
88	/// `f` receives the accumulator and each item in source order. An empty
89	/// stream resolves to the default accumulator unchanged.
90	fn ready_fold_default<T, F>(
91		self,
92		f: F,
93	) -> Fold<Self, Ready<T>, T, impl FnMut(T, Item) -> Ready<T>>
94	where
95		F: Fn(T, Item) -> T,
96		T: Default;
97
98	/// Applies a synchronous closure to every stream item.
99	///
100	/// Items are consumed in source order. The returned future resolves after
101	/// the stream ends and carries no output value.
102	fn ready_for_each<F>(self, f: F) -> ForEach<Self, Ready<()>, impl FnMut(Item) -> Ready<()>>
103	where
104		F: FnMut(Item);
105
106	/// Yields the leading items accepted by a synchronous predicate.
107	///
108	/// Evaluation stops at the first false result, and that boundary item is
109	/// not yielded. Accepted items retain source order.
110	fn ready_take_while<'a, F>(self, f: F) -> impl Stream<Item = Item> + Send + 'a
111	where
112		Self: Send + 'a,
113		Item: Send,
114		F: Fn(&Item) -> bool + Send + 'a;
115
116	/// Transforms items with mutable state and a synchronous scan closure.
117	///
118	/// Each `Some` result is yielded while `None` terminates the stream. State
119	/// is initialized once and retained between items.
120	fn ready_scan<B, T, F>(
121		self,
122		init: T,
123		f: F,
124	) -> Scan<Self, T, Ready<Option<B>>, impl FnMut(&mut T, Item) -> Ready<Option<B>>>
125	where
126		F: Fn(&mut T, Item) -> Option<B>;
127
128	/// Observes each item with mutable state and yields it unchanged.
129	///
130	/// The closure receives the persistent state and a shared reference to each
131	/// item. Every source item is then emitted in order.
132	fn ready_scan_each<T, F>(
133		self,
134		init: T,
135		f: F,
136	) -> Scan<Self, T, Ready<Option<Item>>, impl FnMut(&mut T, Item) -> Ready<Option<Item>>>
137	where
138		F: Fn(&mut T, &Item);
139
140	/// Skips the leading items accepted by a synchronous predicate.
141	///
142	/// The first false item and every later item are yielded in source order.
143	/// The predicate is no longer called after the first false result.
144	fn ready_skip_while<'a, F>(self, f: F) -> impl Stream<Item = Item> + Send + 'a
145	where
146		Self: Send + 'a,
147		Item: Send,
148		F: Fn(&Item) -> bool + Send + 'a;
149}
150
151impl<Item, S> ReadyExt<Item> for S
152where
153	S: Stream<Item = Item> + Sized,
154{
155	#[inline]
156	fn ready_all<F>(self, f: F) -> All<Self, Ready<bool>, impl FnMut(Item) -> Ready<bool>>
157	where
158		F: Fn(Item) -> bool,
159	{
160		self.all(move |t| ready(f(t)))
161	}
162
163	#[inline]
164	fn ready_any<F>(self, f: F) -> Any<Self, Ready<bool>, impl FnMut(Item) -> Ready<bool>>
165	where
166		F: Fn(Item) -> bool,
167	{
168		self.any(move |t| ready(f(t)))
169	}
170
171	#[inline]
172	fn ready_find<'a, F>(self, f: F) -> impl Future<Output = Option<Item>> + Send
173	where
174		Self: Send + Unpin + 'a,
175		F: Fn(&Item) -> bool + Send + 'a,
176		Item: Send,
177	{
178		self.ready_filter(f)
179			.take(1)
180			.into_future()
181			.map(|(curr, _next)| curr)
182	}
183
184	#[inline]
185	fn ready_find_map<'a, F, U>(self, f: F) -> impl Future<Output = Option<U>> + Send
186	where
187		Self: Send + Unpin + 'a,
188		F: Fn(Item) -> Option<U> + Send + 'a,
189		U: Send,
190	{
191		self.ready_filter_map(f)
192			.take(1)
193			.into_future()
194			.map(|(curr, _next)| curr)
195	}
196
197	#[inline]
198	fn ready_filter<'a, F>(self, f: F) -> impl Stream<Item = Item> + Send + 'a
199	where
200		Self: Send + 'a,
201		Item: Send,
202		F: Fn(&Item) -> bool + Send + 'a,
203	{
204		self.filter(move |t| ready(f(t)))
205	}
206
207	#[inline]
208	fn ready_filter_map<F, U>(
209		self,
210		f: F,
211	) -> FilterMap<Self, Ready<Option<U>>, impl FnMut(Item) -> Ready<Option<U>>>
212	where
213		F: Fn(Item) -> Option<U>,
214	{
215		self.filter_map(move |t| ready(f(t)))
216	}
217
218	#[inline]
219	fn ready_fold<T, F>(
220		self,
221		init: T,
222		f: F,
223	) -> Fold<Self, Ready<T>, T, impl FnMut(T, Item) -> Ready<T>>
224	where
225		F: Fn(T, Item) -> T,
226	{
227		self.fold(init, move |a, t| ready(f(a, t)))
228	}
229
230	#[inline]
231	fn ready_fold_default<T, F>(
232		self,
233		f: F,
234	) -> Fold<Self, Ready<T>, T, impl FnMut(T, Item) -> Ready<T>>
235	where
236		F: Fn(T, Item) -> T,
237		T: Default,
238	{
239		self.ready_fold(T::default(), f)
240	}
241
242	#[inline]
243	#[expect(clippy::unit_arg)]
244	fn ready_for_each<F>(
245		self,
246		mut f: F,
247	) -> ForEach<Self, Ready<()>, impl FnMut(Item) -> Ready<()>>
248	where
249		F: FnMut(Item),
250	{
251		self.for_each(move |t| ready(f(t)))
252	}
253
254	#[inline]
255	fn ready_take_while<'a, F>(self, f: F) -> impl Stream<Item = Item> + Send + 'a
256	where
257		Self: Send + 'a,
258		Item: Send,
259		F: Fn(&Item) -> bool + Send + 'a,
260	{
261		self.take_while(move |t| ready(f(t)))
262	}
263
264	#[inline]
265	fn ready_scan<B, T, F>(
266		self,
267		init: T,
268		f: F,
269	) -> Scan<Self, T, Ready<Option<B>>, impl FnMut(&mut T, Item) -> Ready<Option<B>>>
270	where
271		F: Fn(&mut T, Item) -> Option<B>,
272	{
273		self.scan(init, move |s, t| ready(f(s, t)))
274	}
275
276	#[inline]
277	fn ready_scan_each<T, F>(
278		self,
279		init: T,
280		f: F,
281	) -> Scan<Self, T, Ready<Option<Item>>, impl FnMut(&mut T, Item) -> Ready<Option<Item>>>
282	where
283		F: Fn(&mut T, &Item),
284	{
285		self.ready_scan(init, move |s, t| {
286			f(s, &t);
287			Some(t)
288		})
289	}
290
291	#[inline]
292	fn ready_skip_while<'a, F>(self, f: F) -> impl Stream<Item = Item> + Send + 'a
293	where
294		Self: Send + 'a,
295		Item: Send,
296		F: Fn(&Item) -> bool + Send + 'a,
297	{
298		self.skip_while(move |t| ready(f(t)))
299	}
300}