1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
use futures::{
channel::mpsc::{self, Sender},
future::{Either, FutureExt},
stream::{self, StreamExt},
Future, Stream,
};
pub fn gen_z_sender<
'fut,
T: 'fut + Send,
Fut: Future<Output = ()> + 'fut + Send,
Gen: FnOnce(Sender<T>) -> Fut,
>(
gen: Gen,
) -> impl Stream<Item = T> + 'fut {
let (send, recv) = mpsc::channel::<T>(0);
let fut: Fut = gen(send);
let fut = fut.map(Either::Left).into_stream();
let joined = stream::select(fut, recv.map(Either::Right)).filter_map(|x| {
futures::future::ready(match x {
Either::Left(_) => None,
Either::Right(y) => Some(y),
})
});
joined.fuse()
}
async fn send_infallible<T: Send>(sender: &mut Sender<T>, item: T) -> () {
use futures::sink::SinkExt;
SinkExt::send(sender, item)
.await
.expect("Infallible sends should only be used where they cannot fail")
}
pub struct Yielder<T: Send> {
sender: Sender<T>,
}
impl<T: Send> Yielder<T> {
pub fn new(sender: Sender<T>) -> Yielder<T> {
Yielder { sender }
}
pub fn send(&mut self, item: T) -> impl Future<Output = ()> + '_ {
send_infallible(&mut self.sender, item)
}
}
pub fn gen_z<
'fut,
T: 'fut + Send,
Fut: Future<Output = ()> + 'fut + Send,
Gen: 'fut + FnOnce(Yielder<T>) -> Fut,
>(
gen: Gen,
) -> impl Stream<Item = T> + 'fut {
gen_z_sender(|sender| gen(Yielder::<T>::new(sender)))
}