use crate::streams::error::{RecvError, SendError};
pub struct StreamReader<T> {
reader: tokio::sync::mpsc::Receiver<T>,
}
impl<T> StreamReader<T>
where
T: Send,
{
fn new(rx: tokio::sync::mpsc::Receiver<T>) -> Self {
Self { reader: rx }
}
pub async fn recv(&mut self) -> Result<T, RecvError> {
match self.reader.recv().await {
Some(s) => Ok(s),
None => Err(RecvError::Closed),
}
}
}
#[derive(Debug)]
pub struct StreamWriter<T> {
sender: tokio::sync::mpsc::Sender<T>,
}
impl<T> Clone for StreamWriter<T> {
fn clone(&self) -> Self {
Self {
sender: self.sender.clone(),
}
}
}
impl<T> StreamWriter<T>
where
T: Send,
{
fn new(tx: tokio::sync::mpsc::Sender<T>) -> Self {
Self { sender: tx }
}
pub async fn send(&self, data: T) -> Result<(), SendError> {
match self.sender.send(data).await {
Ok(_) => Ok(()),
Err(e) => Err(SendError::from(e)),
}
}
pub fn blocking_send(&self, data: T) -> Result<(), SendError> {
match self.sender.blocking_send(data) {
Ok(_) => Ok(()),
Err(e) => Err(SendError::from(e)),
}
}
}
pub fn stream<T>() -> (StreamWriter<T>, StreamReader<T>)
where
T: Send,
{
let (tx, rx) = tokio::sync::mpsc::channel(25);
(StreamWriter::new(tx), StreamReader::new(rx))
}