use std::io::{stdin, BufRead, BufReader};
use tokio::sync::mpsc;
pub fn recv_from_stdin(buffer_size: usize) -> mpsc::Receiver<String> {
let (tx, rx) = mpsc::channel::<String>(buffer_size);
let stdin = BufReader::new(stdin());
std::thread::spawn(move || read_loop(stdin, tx));
rx
}
fn read_loop<R>(reader: R, tx: mpsc::Sender<String>)
where
R: BufRead,
{
let mut lines = reader.lines();
loop {
if let Some(Ok(line)) = lines.next() {
let _ = tx.blocking_send(line);
}
}
}
#[cfg(test)]
mod tests {
#[tokio::test]
async fn test_blocking_read() {
let (tx, mut rx) = tokio::sync::mpsc::channel::<String>(10);
let reader = std::io::BufReader::new("hello".as_bytes());
std::thread::spawn(move || super::read_loop(reader, tx));
let s = rx.recv().await.unwrap();
assert_eq!(s, "hello");
}
}