pub(crate) mod handler;
mod receiver;
mod sender;
use crate::error::{StreamError, TaskError};
use crate::util::watch;
use tokio::sync::{mpsc, oneshot};
use tokio::task::JoinHandle;
pub use receiver::Receiver;
pub use sender::Sender;
pub(crate) enum SendBack<P> {
None,
Packet(P),
Close,
CloseWithPacket,
}
#[derive(Debug)]
pub(crate) struct TaskHandle {
pub close: oneshot::Sender<()>,
pub task: JoinHandle<Result<(), TaskError>>,
}
impl TaskHandle {
pub async fn closed(&mut self) {
self.close.closed().await;
}
pub async fn wait(self) -> Result<(), TaskError> {
self.task.await.map_err(TaskError::Join)?
}
pub async fn close(self) -> Result<(), TaskError> {
let _ = self.close.send(());
self.task.await.map_err(TaskError::Join)?
}
#[cfg(test)]
pub fn abort(self) {
self.task.abort();
}
}
#[derive(Debug, Clone)]
pub struct StreamSender<P> {
pub(crate) inner: mpsc::Sender<P>,
}
impl<P> StreamSender<P> {
pub(crate) fn new(inner: mpsc::Sender<P>) -> Self {
Self { inner }
}
pub async fn send(&self, packet: P) -> Result<(), StreamError> {
self.inner
.send(packet)
.await
.map_err(|_| StreamError::StreamAlreadyClosed)
}
}
#[derive(Debug)]
pub struct StreamReceiver<P> {
pub(crate) inner: mpsc::Receiver<P>,
}
impl<P> StreamReceiver<P> {
pub(crate) fn new(inner: mpsc::Receiver<P>) -> Self {
Self { inner }
}
pub async fn receive(&mut self) -> Option<P> {
self.inner.recv().await
}
pub fn close(&mut self) {
self.inner.close();
}
}
#[derive(Debug)]
pub struct Configurator<C> {
inner: watch::Sender<C>,
}
impl<C> Clone for Configurator<C> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
}
}
}
impl<C> Configurator<C> {
pub(crate) fn new(cfg: C) -> (Self, watch::Receiver<C>) {
let (tx, rx) = watch::channel(cfg);
(Self { inner: tx }, rx)
}
pub fn update(&self, cfg: C) {
self.inner.send(cfg);
}
pub fn read(&self) -> C
where
C: Clone,
{
self.inner.newest()
}
}