Skip to main content

tuwunel_core/utils/stream/
try_tools.rs

1//! TryStreamTools for futures::TryStream
2#![expect(clippy::type_complexity)]
3
4use futures::{
5	TryStream, TryStreamExt, future,
6	future::{Ready, ready},
7	stream::TryTakeWhile,
8};
9
10use crate::Result;
11
12/// Adds general-purpose operations to fallible streams.
13///
14/// Operations preserve the source error type and process successful items in
15/// source order. They support limiting the stream and collecting successful pairs.
16pub trait TryTools<T, E, S>
17where
18	S: TryStream<Ok = T, Error = E, Item = Result<T, E>> + ?Sized,
19	Self: TryStream + Sized,
20{
21	/// Limits the stream to at most `n` successful items.
22	///
23	/// After yielding `n` successes, the adapter consumes one additional
24	/// success to detect the limit. Earlier source errors are still forwarded;
25	/// with zero, they precede consumption of the first unyielded success.
26	fn try_take(
27		self,
28		n: usize,
29	) -> TryTakeWhile<
30		Self,
31		Ready<Result<bool, S::Error>>,
32		impl FnMut(&S::Ok) -> Ready<Result<bool, S::Error>>,
33	>;
34
35	/// Collects successful pairs into two collections.
36	///
37	/// Collections are extended in source order, starting from their defaults.
38	/// The first source error ends the fold without returning partial collections.
39	fn try_unzip<FromA, FromB>(self) -> impl Future<Output = Result<(FromA, FromB), S::Error>>
40	where
41		(FromA, FromB): Default + Extend<T>;
42}
43
44impl<T, E, S> TryTools<T, E, S> for S
45where
46	S: TryStream<Ok = T, Error = E, Item = Result<T, E>> + ?Sized,
47	Self: TryStream + Sized,
48{
49	#[inline]
50	fn try_take(
51		self,
52		mut n: usize,
53	) -> TryTakeWhile<
54		Self,
55		Ready<Result<bool, S::Error>>,
56		impl FnMut(&S::Ok) -> Ready<Result<bool, S::Error>>,
57	> {
58		self.try_take_while(move |_| {
59			let res = future::ok(n > 0);
60			n = n.saturating_sub(1);
61			res
62		})
63	}
64
65	#[inline]
66	fn try_unzip<FromA, FromB>(self) -> impl Future<Output = Result<(FromA, FromB), S::Error>>
67	where
68		(FromA, FromB): Default + Extend<T>,
69	{
70		self.try_fold(<(FromA, FromB)>::default(), |mut collections, item| {
71			collections.extend([item]);
72			ready(Ok(collections))
73		})
74	}
75}