use tokio::sync::{
broadcast::{self, error::SendError},
mpsc::Receiver,
};
use tokio_util::sync::{CancellationToken, WaitForCancellationFuture};
mod private {
pub trait Sealed {}
}
pub trait Handle: private::Sealed {}
mod modes {
use tokio::sync::{broadcast, mpsc::Receiver};
pub struct Isolated {}
pub struct OneWay<Message> {
pub(super) receiver_from_wk: Receiver<Message>,
}
pub struct TwoWay<InMessage, OutMessage> {
pub(super) receiver_from_wk: Receiver<InMessage>,
pub(crate) broadcast_from_task: broadcast::Sender<OutMessage>,
}
pub struct OneWayBack<OutMessage> {
pub(crate) broadcast_from_task: broadcast::Sender<OutMessage>,
}
pub struct OnEvent<Event> {
pub(crate) receiver_from_task: broadcast::Receiver<Event>,
}
}
pub use modes::*;
pub struct Worker<Mode> {
termination_token: CancellationToken,
mode: Mode,
}
impl<Mode> Worker<Mode> {
pub fn terminated(&self) -> WaitForCancellationFuture<'_> {
self.termination_token.cancelled()
}
pub fn terminate(self) {
self.termination_token.cancel();
}
}
impl private::Sealed for Worker<Isolated> {}
impl Handle for Worker<Isolated> {}
impl Worker<Isolated> {
pub(crate) fn isolated(token: CancellationToken) -> Worker<Isolated> {
Self {
termination_token: token,
mode: Isolated {},
}
}
}
impl<Message> private::Sealed for Worker<OneWay<Message>> {}
impl<Message> Handle for Worker<OneWay<Message>> {}
impl<Message> Worker<OneWay<Message>> {
pub(crate) fn one_way(
token: CancellationToken,
from_wk: Receiver<Message>,
) -> Worker<OneWay<Message>> {
Self {
termination_token: token,
mode: OneWay {
receiver_from_wk: from_wk,
},
}
}
pub fn receiver(self) -> (Receiver<Message>, Worker<Isolated>) {
let Worker {
termination_token,
mode,
} = self;
let OneWay { receiver_from_wk } = mode;
(receiver_from_wk, Worker::isolated(termination_token))
}
}
impl<InMessage, OutMessage> private::Sealed for Worker<TwoWay<InMessage, OutMessage>> {}
impl<InMessage, OutMessage> Handle for Worker<TwoWay<InMessage, OutMessage>> {}
impl<InMessage, OutMessage> Worker<TwoWay<InMessage, OutMessage>> {
pub(crate) fn two_way(
token: CancellationToken,
from_wk: Receiver<InMessage>,
to_task: broadcast::Sender<OutMessage>,
) -> Worker<TwoWay<InMessage, OutMessage>> {
Self {
termination_token: token,
mode: TwoWay {
receiver_from_wk: from_wk,
broadcast_from_task: to_task,
},
}
}
pub fn receiver(self) -> (Receiver<InMessage>, Worker<OneWayBack<OutMessage>>) {
let Worker {
termination_token,
mode,
} = self;
let TwoWay {
receiver_from_wk,
broadcast_from_task,
} = mode;
(
receiver_from_wk,
Worker::one_way_back(termination_token, broadcast_from_task),
)
}
}
impl<OutMessage> private::Sealed for Worker<OneWayBack<OutMessage>> {}
impl<OutMessage> Handle for Worker<OneWayBack<OutMessage>> {}
impl<OutMessage> Worker<OneWayBack<OutMessage>> {
pub(crate) fn one_way_back(
token: CancellationToken,
to_task: broadcast::Sender<OutMessage>,
) -> Worker<OneWayBack<OutMessage>> {
Self {
termination_token: token,
mode: OneWayBack {
broadcast_from_task: to_task,
},
}
}
pub async fn post_message(&self, msg: OutMessage) -> Result<usize, SendError<OutMessage>> {
self.mode.broadcast_from_task.send(msg)
}
}
impl<Event> private::Sealed for Worker<OnEvent<Event>> {}
impl<Event> Handle for Worker<OnEvent<Event>> {}
impl<Event> Worker<OnEvent<Event>> {
pub(crate) fn on_event(
token: CancellationToken,
from_task: broadcast::Receiver<Event>,
) -> Worker<OnEvent<Event>> {
Self {
termination_token: token,
mode: OnEvent {
receiver_from_task: from_task,
},
}
}
pub fn receiver(self) -> (broadcast::Receiver<Event>, Worker<Isolated>) {
let Worker {
termination_token,
mode,
} = self;
let OnEvent { receiver_from_task } = mode;
(receiver_from_task, Worker::isolated(termination_token))
}
}