summaryrefslogtreecommitdiffstats
path: root/third_party/rust/tokio-stream/tests/support/mpsc.rs
blob: 09dbe04215eb95081f2c163c0e4c266cdc0056c6 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
use async_stream::stream;
use tokio::sync::mpsc::{self, UnboundedSender};
use tokio_stream::Stream;

pub fn unbounded_channel_stream<T: Unpin>() -> (UnboundedSender<T>, impl Stream<Item = T>) {
    let (tx, mut rx) = mpsc::unbounded_channel();

    let stream = stream! {
        while let Some(item) = rx.recv().await {
            yield item;
        }
    };

    (tx, stream)
}