use crate::actor::addr::Addr;
use crate::actor::event_bus::rpc::{AsyncBus, RpcCall, RpcRequest};
use crate::actor::short_type_name;
use crate::actor::traits::Handler;
use crate::trace::Bus;
use std::any::{Any, TypeId};
use std::marker::PhantomData;
use std::rc::Weak;
use super::EventBus;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct SubscriptionId {
pub(super) seq: u64,
pub(super) event: TypeId,
}
pub struct BusSubscription {
pub(super) bus: Weak<EventBus>,
pub(super) id: SubscriptionId,
}
impl BusSubscription {
pub fn leak(self) {
std::mem::forget(self);
}
}
impl Drop for BusSubscription {
fn drop(&mut self) {
if let Some(bus) = self.bus.upgrade() {
bus.remove(self.id);
}
}
}
impl crate::scope::Teardown for BusSubscription {
fn teardown(self) {
drop(self);
}
}
pub trait Event: Clone + Send + 'static {
#[doc(hidden)]
fn overheard(self) -> Self {
self
}
}
pub trait UntypedSubscriber: 'static {
fn deliver(&self, msg: Box<dyn Any>, bus: Bus);
fn seq(&self) -> u64;
fn event(&self) -> &'static str;
fn answerer(&self) -> Option<&'static str> {
None
}
fn is_asleep(&self) -> bool {
false
}
}
pub struct Subscriber<A: Handler<M>, M: Event> {
pub(super) seq: u64,
pub(super) addr: Addr<A>,
pub(super) _marker: PhantomData<M>,
}
impl<A, M> UntypedSubscriber for Subscriber<A, M>
where
A: Handler<M> + 'static,
M: Event,
{
fn deliver(&self, msg: Box<dyn Any>, _bus: Bus) {
if self.addr.is_asleep() {
return;
}
if let Ok(concrete_msg) = msg.downcast::<M>() {
let message = match <A as Handler<M>>::ANSWERS {
true => *concrete_msg,
false => concrete_msg.overheard(),
};
self.addr.send(message);
}
}
fn seq(&self) -> u64 {
self.seq
}
fn event(&self) -> &'static str {
short_type_name::<M>()
}
fn answerer(&self) -> Option<&'static str> {
<A as Handler<M>>::ANSWERS.then(short_type_name::<A>)
}
fn is_asleep(&self) -> bool {
self.addr.is_asleep()
}
}
pub struct AnswerFn<Req: RpcCall> {
pub(super) seq: u64,
pub(super) answer: Box<dyn Fn(Req) -> Req::Response>,
}
impl<Req: RpcCall> UntypedSubscriber for AnswerFn<Req> {
fn deliver(&self, msg: Box<dyn Any>, _bus: Bus) {
if let Ok(request) = msg.downcast::<RpcRequest<Req>>() {
let RpcRequest {
correlation_id,
payload,
..
} = *request;
AsyncBus::reply(correlation_id, (self.answer)(payload));
}
}
fn seq(&self) -> u64 {
self.seq
}
fn event(&self) -> &'static str {
short_type_name::<RpcRequest<Req>>()
}
fn answerer(&self) -> Option<&'static str> {
Some("a callback")
}
}
pub struct FnSubscriber<M: Event> {
pub(super) seq: u64,
pub(super) callback: std::sync::Arc<dyn Fn(M) + 'static>,
}
impl<M: Event> UntypedSubscriber for FnSubscriber<M> {
fn deliver(&self, msg: Box<dyn Any>, _bus: Bus) {
if let Ok(concrete_msg) = msg.downcast::<M>() {
(self.callback)(concrete_msg.overheard());
}
}
fn seq(&self) -> u64 {
self.seq
}
fn event(&self) -> &'static str {
short_type_name::<M>()
}
}