tuwunel_core/utils/stream/
try_tools.rs1#![expect(clippy::type_complexity)]
3
4use futures::{
5 TryStream, TryStreamExt, future,
6 future::{Ready, ready},
7 stream::TryTakeWhile,
8};
9
10use crate::Result;
11
12pub trait TryTools<T, E, S>
17where
18 S: TryStream<Ok = T, Error = E, Item = Result<T, E>> + ?Sized,
19 Self: TryStream + Sized,
20{
21 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 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}