use tokio::sync::{broadcast, mpsc::channel};
use tokio_util::sync::CancellationToken;
use crate::{handle, Error, Task, BUFFER_CAPACITY};
mod modes {
use tokio::sync::{broadcast, mpsc::Sender};
pub struct Isolated {}
pub struct OneWay<Message> {
pub(super) sender_to_tsk: Sender<Message>,
}
pub struct TwoWay<Message, TaskMessage> {
pub(super) sender_to_tsk: Sender<Message>,
pub(super) broadcast_from_tsk: broadcast::Sender<TaskMessage>,
}
}
pub use modes::*;
pub struct Worker<Mode> {
termination_token: CancellationToken,
mode: Mode,
}
impl<Mode> Worker<Mode> {
pub fn terminate(self) {
self.termination_token.cancel();
}
}
impl Worker<Isolated> {
pub fn spawn<T>(task: T) -> Worker<Isolated>
where
T: Task<Handle = handle::Worker<handle::Isolated>>,
<T as Task>::Output: Send + 'static,
{
let token = CancellationToken::new();
let wkh = handle::Worker::isolated(token.clone());
tokio::spawn(task.spawn(wkh));
Worker {
termination_token: token,
mode: Isolated {},
}
}
}
impl<Message> Worker<OneWay<Message>> {
pub fn spawn<T>(task: T) -> Worker<OneWay<Message>>
where
T: Task<Handle = handle::Worker<handle::OneWay<Message>>>,
<T as Task>::Output: Send + 'static,
{
let token = CancellationToken::new();
let (send_to_task, recv_from_wk) = channel::<Message>(BUFFER_CAPACITY);
let wkh = handle::Worker::one_way(token.clone(), recv_from_wk);
tokio::spawn(task.spawn(wkh));
Worker {
termination_token: token,
mode: OneWay {
sender_to_tsk: send_to_task,
},
}
}
pub async fn post_message(&self, msg: Message) -> Result<(), Error> {
self.mode
.sender_to_tsk
.send(msg)
.await
.map_err(|e| Error::from(&e))
}
}
impl<Message, TaskMessage: Clone> Worker<TwoWay<Message, TaskMessage>> {
pub fn spawn<T>(task: T) -> Worker<TwoWay<Message, TaskMessage>>
where
T: Task<Handle = handle::Worker<handle::TwoWay<Message, TaskMessage>>>,
<T as Task>::Output: Send + 'static,
{
let token = CancellationToken::new();
let (send_to_task, recv_from_wk) = channel::<Message>(BUFFER_CAPACITY);
let (broadcast_to_wk, _) = broadcast::channel::<TaskMessage>(BUFFER_CAPACITY);
let wkh = handle::Worker::two_way(token.clone(), recv_from_wk, broadcast_to_wk.to_owned());
tokio::spawn(task.spawn(wkh));
Worker {
termination_token: token,
mode: TwoWay {
sender_to_tsk: send_to_task,
broadcast_from_tsk: broadcast_to_wk,
},
}
}
pub async fn post_message(&self, msg: Message) -> Result<(), Error> {
self.mode
.sender_to_tsk
.send(msg)
.await
.map_err(|e| Error::from(&e))
}
pub fn on_message<T>(&self, task: T) -> Worker<Isolated>
where
T: Task<Handle = handle::Worker<handle::OnEvent<TaskMessage>>>,
<T as Task>::Output: Send + 'static,
{
let token = CancellationToken::new();
let wkh = handle::Worker::on_event(token.clone(), self.mode.broadcast_from_tsk.subscribe());
tokio::spawn(task.spawn(wkh));
Worker {
termination_token: token,
mode: Isolated {},
}
}
}