use super::{Dispatcher, DispatcherError, MessageType, ObserverRef};
use tokio::sync::broadcast::{channel, error, Receiver, Sender};
pub struct Broadcaster<M> {
local: Dispatcher<M>,
broadcast_sender: Sender<M>,
_broadcast_receiver: Option<Receiver<M>>,
}
impl<M> Default for Broadcaster<M>
where
M: Clone + MessageType + std::default::Default,
{
fn default() -> Self {
Self::new(100)
}
}
impl<M> Clone for Broadcaster<M>
where
M: Clone + MessageType + std::default::Default,
{
fn clone(&self) -> Self {
Self {
local: Dispatcher::default(),
broadcast_sender: self.broadcast_sender.clone(),
_broadcast_receiver: None,
}
}
}
impl<M> Broadcaster<M>
where
M: Clone + MessageType + std::default::Default,
{
pub fn new(capacity: usize) -> Self {
let (broadcast_sender, broadcast_receiver) = channel(capacity);
Self {
local: Dispatcher::default(),
broadcast_sender,
_broadcast_receiver: Some(broadcast_receiver),
}
}
pub fn register_handler(&mut self, message_type: &str, observer: ObserverRef<M>, tag: &str) {
self.local.register_handler(message_type, observer, tag);
}
pub fn unregister_handler(&mut self, message_type: &str, tag: &str) {
self.local.unregister_handler(message_type, tag);
}
pub fn dispatch(&self, message: M) -> Result<usize, DispatcherError> {
let n1 = self.local.dispatch(&message)?;
self.broadcast_sender
.send(message)
.map(|n| n1 + n)
.map_err(|err| DispatcherError::SendError(err.to_string()))
}
pub fn dispatch_local(&self, message: &M) -> Result<usize, DispatcherError> {
self.local.dispatch(message)
}
pub fn send(&self, message: M) -> Result<usize, error::SendError<M>> {
self.broadcast_sender.send(message)
}
pub fn receiver(&self) -> Receiver<M> {
self.broadcast_sender.subscribe()
}
pub fn sender(&self) -> Sender<M> {
self.broadcast_sender.clone()
}
}