File: mpsc.rs

package info (click to toggle)
rust-tokio-stream 0.1.17-1
  • links: PTS, VCS
  • area: main
  • in suites: forky, sid
  • size: 432 kB
  • sloc: makefile: 2
file content (15 lines) | stat: -rw-r--r-- 382 bytes parent folder | download | duplicates (28)
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)
}