Skip to main content

tuwunel_core/utils/stream/
try_broadband.rs

1//! Synchronous combinator extensions to futures::TryStream
2
3use futures::{TryFuture, TryStream, TryStreamExt, future::ready};
4
5use super::automatic_width;
6use crate::Result;
7
8/// Adds bounded concurrent transformations with completion-ordered outputs.
9///
10/// Successful item futures may run ahead of downstream demand, and their
11/// results are yielded as they complete. Source errors bypass the transform and
12/// may overtake queued item futures.
13pub trait TryBroadbandExt<T, E>
14where
15	Self: TryStream<Ok = T, Error = E, Item = Result<T, E>> + Send + Sized,
16{
17	/// Tests all successful items concurrently with an explicit width.
18	///
19	/// `n` limits in-flight predicates, `None` selects the automatic width,
20	/// and an explicit zero cannot make progress. The first observed false or
21	/// source or predicate error short-circuits; source errors may overtake
22	/// queued predicates. An empty stream resolves `Ok(true)`.
23	fn broadn_try_all<F, Fut, N>(
24		self,
25		n: N,
26		f: F,
27	) -> impl Future<Output = Result<bool, E>> + Send
28	where
29		N: Into<Option<usize>>,
30		F: Fn(Self::Ok) -> Fut + Send,
31		Fut: TryFuture<Ok = bool, Error = E, Output = Result<bool, E>> + Send;
32
33	/// Tests successful items concurrently for any match with an explicit width.
34	///
35	/// `n` limits in-flight predicates, `None` selects the automatic width,
36	/// and an explicit zero cannot make progress. The first observed true or
37	/// source or predicate error short-circuits; source errors may overtake
38	/// queued predicates. An empty stream resolves `Ok(false)`.
39	fn broadn_try_any<F, Fut, N>(
40		self,
41		n: N,
42		f: F,
43	) -> impl Future<Output = Result<bool, E>> + Send
44	where
45		N: Into<Option<usize>>,
46		F: Fn(Self::Ok) -> Fut + Send,
47		Fut: TryFuture<Ok = bool, Error = E, Output = Result<bool, E>> + Send;
48
49	/// Transforms successful items concurrently with an explicit width.
50	///
51	/// `n` limits in-flight item futures, while `None` selects the automatic
52	/// width; an explicit zero cannot make progress. Transformation results are
53	/// completion ordered, while source errors bypass `f` and may overtake
54	/// them.
55	fn broadn_and_then<U, F, Fut, N>(
56		self,
57		n: N,
58		f: F,
59	) -> impl TryStream<Ok = U, Error = E, Item = Result<U, E>> + Send
60	where
61		N: Into<Option<usize>>,
62		F: Fn(Self::Ok) -> Fut + Send,
63		Fut: TryFuture<Ok = U, Error = E, Output = Result<U, E>> + Send;
64
65	/// Tests all successful items concurrently with the automatic width.
66	///
67	/// The first observed false or source or predicate error short-circuits,
68	/// so a false result can precede a pending error. An empty stream resolves
69	/// `Ok(true)`.
70	fn broad_try_all<F, Fut>(self, f: F) -> impl Future<Output = Result<bool, E>> + Send
71	where
72		F: Fn(Self::Ok) -> Fut + Send,
73		Fut: TryFuture<Ok = bool, Error = E, Output = Result<bool, E>> + Send,
74	{
75		self.broadn_try_all(None, f)
76	}
77
78	/// Tests successful items concurrently for any match with the automatic width.
79	///
80	/// The first observed true or source or predicate error short-circuits,
81	/// so a true result can precede a pending error. An empty stream resolves
82	/// `Ok(false)`.
83	fn broad_try_any<F, Fut>(self, f: F) -> impl Future<Output = Result<bool, E>> + Send
84	where
85		F: Fn(Self::Ok) -> Fut + Send,
86		Fut: TryFuture<Ok = bool, Error = E, Output = Result<bool, E>> + Send,
87	{
88		self.broadn_try_any(None, f)
89	}
90
91	/// Transforms successful items concurrently with the automatic width.
92	///
93	/// Transformation results are yielded in completion order rather than
94	/// source order. Existing source errors bypass `f` and may overtake queued
95	/// futures.
96	fn broad_and_then<U, F, Fut>(
97		self,
98		f: F,
99	) -> impl TryStream<Ok = U, Error = E, Item = Result<U, E>> + Send
100	where
101		F: Fn(Self::Ok) -> Fut + Send,
102		Fut: TryFuture<Ok = U, Error = E, Output = Result<U, E>> + Send,
103	{
104		self.broadn_and_then(None, f)
105	}
106}
107
108impl<T, E, S> TryBroadbandExt<T, E> for S
109where
110	S: TryStream<Ok = T, Error = E, Item = Result<T, E>> + Send + Sized,
111{
112	fn broadn_try_all<F, Fut, N>(self, n: N, f: F) -> impl Future<Output = Result<bool, E>> + Send
113	where
114		N: Into<Option<usize>>,
115		F: Fn(Self::Ok) -> Fut + Send,
116		Fut: TryFuture<Ok = bool, Error = E, Output = Result<bool, E>> + Send,
117	{
118		self.broadn_and_then(n, f).try_all(ready)
119	}
120
121	fn broadn_try_any<F, Fut, N>(self, n: N, f: F) -> impl Future<Output = Result<bool, E>> + Send
122	where
123		N: Into<Option<usize>>,
124		F: Fn(Self::Ok) -> Fut + Send,
125		Fut: TryFuture<Ok = bool, Error = E, Output = Result<bool, E>> + Send,
126	{
127		self.broadn_and_then(n, f).try_any(ready)
128	}
129
130	fn broadn_and_then<U, F, Fut, N>(
131		self,
132		n: N,
133		f: F,
134	) -> impl TryStream<Ok = U, Error = E, Item = Result<U, E>> + Send
135	where
136		N: Into<Option<usize>>,
137		F: Fn(Self::Ok) -> Fut + Send,
138		Fut: TryFuture<Ok = U, Error = E, Output = Result<U, E>> + Send,
139	{
140		self.map_ok(f)
141			.try_buffer_unordered(n.into().unwrap_or_else(automatic_width))
142	}
143}