use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use crate::ctx::Ctx;
use crate::error::{join_panic_error, panic_error, CordisError};
use crate::event::{
CatchUnwind, DynEvent, ErasedValue, Event, EventOptions, Listener, ListenerAdapter, Terminal,
TerminalAdapter, WaterfallAdapter, WaterfallListener,
};
use crate::key::TypeKey;
use crate::{BoxFuture, Disposer, Effect};
pub(crate) struct ErasedNext<'a> {
chain: &'a [Arc<dyn ErasedWaterfallCall>],
index: usize,
ctx: &'a Ctx,
event: &'a DynEvent,
terminal: &'a mut (dyn ErasedTerminal + 'a),
}
impl<'a> ErasedNext<'a> {
pub(crate) fn invoke(self) -> BoxFuture<'a, Result<ErasedValue, CordisError>> {
if self.index < self.chain.len() {
let ErasedNext {
chain,
index,
ctx,
event,
terminal,
} = self;
chain[index].call(
ctx,
event,
ErasedNext {
chain,
index: index + 1,
ctx,
event,
terminal,
},
)
} else {
let ErasedNext {
ctx,
event,
terminal,
..
} = self;
terminal.call(ctx, event)
}
}
}
pub(crate) trait ErasedCall: Send + Sync + 'static {
fn call<'a>(
&'a self,
ctx: &'a Ctx,
e: &'a DynEvent,
) -> BoxFuture<'a, Result<Option<ErasedValue>, CordisError>>;
}
pub(crate) trait ErasedWaterfallCall: Send + Sync + 'static {
fn call<'a>(
&'a self,
ctx: &'a Ctx,
e: &'a DynEvent,
next: ErasedNext<'a>,
) -> BoxFuture<'a, Result<ErasedValue, CordisError>>;
}
pub(crate) trait ErasedTerminal: Send {
fn call<'a>(
&'a mut self,
ctx: &'a Ctx,
e: &'a DynEvent,
) -> BoxFuture<'a, Result<ErasedValue, CordisError>>;
}
struct Hook<C> {
call: C,
once: bool,
}
fn insert_hook<C>(list: &mut Vec<Arc<Hook<C>>>, hook: Arc<Hook<C>>, prepend: bool) {
if prepend {
list.insert(0, hook);
} else {
list.push(hook);
}
}
fn retain_hook<C>(list: &mut Vec<Arc<Hook<C>>>, hook: &Arc<Hook<C>>) {
list.retain(|h| !Arc::ptr_eq(h, hook));
}
fn claim_once<C>(list: &mut Vec<Arc<Hook<C>>>) -> Vec<Arc<Hook<C>>> {
let snapshot = list.clone();
list.retain(|h| !h.once);
snapshot
}
#[derive(Default)]
struct BusInner {
hooks: HashMap<TypeKey, Vec<Arc<Hook<Arc<dyn ErasedCall>>>>>,
wf_hooks: HashMap<TypeKey, Vec<Arc<Hook<Arc<dyn ErasedWaterfallCall>>>>>,
dispatch_tail: HashMap<TypeKey, tokio::task::JoinHandle<()>>,
}
#[derive(Clone)]
pub struct EventBus {
inner: Arc<Mutex<BusInner>>,
}
impl EventBus {
pub(crate) fn new() -> Self {
Self {
inner: Arc::new(Mutex::new(BusInner::default())),
}
}
pub fn on<E: Event>(&self, ctx: &Ctx, l: impl Listener<E>) -> Result<Disposer, CordisError> {
self.add_hook(TypeKey::of::<E>(), ctx, l, EventOptions::default(), false)
}
pub fn on_opt<E: Event>(
&self,
ctx: &Ctx,
l: impl Listener<E>,
opts: EventOptions,
) -> Result<Disposer, CordisError> {
self.add_hook(TypeKey::of::<E>(), ctx, l, opts, false)
}
pub fn once<E: Event>(&self, ctx: &Ctx, l: impl Listener<E>) -> Result<Disposer, CordisError> {
self.add_hook(TypeKey::of::<E>(), ctx, l, EventOptions::default(), true)
}
pub fn on_keyed<E: Event>(
&self,
ctx: &Ctx,
name: impl Into<std::sync::Arc<str>>,
l: impl Listener<E>,
) -> Result<Disposer, CordisError> {
self.add_hook(
TypeKey::keyed_dynamic::<E>(name),
ctx,
l,
EventOptions::default(),
false,
)
}
pub fn on_keyed_opt<E: Event>(
&self,
ctx: &Ctx,
name: impl Into<std::sync::Arc<str>>,
l: impl Listener<E>,
opts: EventOptions,
) -> Result<Disposer, CordisError> {
self.add_hook(TypeKey::keyed_dynamic::<E>(name), ctx, l, opts, false)
}
pub fn once_keyed<E: Event>(
&self,
ctx: &Ctx,
name: impl Into<std::sync::Arc<str>>,
l: impl Listener<E>,
) -> Result<Disposer, CordisError> {
self.add_hook(
TypeKey::keyed_dynamic::<E>(name),
ctx,
l,
EventOptions::default(),
true,
)
}
pub fn on_waterfall<E: Event>(
&self,
ctx: &Ctx,
l: impl WaterfallListener<E>,
) -> Result<Disposer, CordisError> {
self.add_wf_hook(TypeKey::of::<E>(), ctx, l, EventOptions::default(), false)
}
pub fn on_waterfall_opt<E: Event>(
&self,
ctx: &Ctx,
l: impl WaterfallListener<E>,
opts: EventOptions,
) -> Result<Disposer, CordisError> {
self.add_wf_hook(TypeKey::of::<E>(), ctx, l, opts, false)
}
pub fn on_waterfall_keyed<E: Event>(
&self,
ctx: &Ctx,
name: impl Into<std::sync::Arc<str>>,
l: impl WaterfallListener<E>,
) -> Result<Disposer, CordisError> {
self.add_wf_hook(
TypeKey::keyed_dynamic::<E>(name),
ctx,
l,
EventOptions::default(),
false,
)
}
fn add_hook<E: Event>(
&self,
key: TypeKey,
ctx: &Ctx,
l: impl Listener<E>,
opts: EventOptions,
once: bool,
) -> Result<Disposer, CordisError> {
let hook: Arc<Hook<Arc<dyn ErasedCall>>> = Arc::new(Hook {
call: Arc::new(ListenerAdapter(l, std::marker::PhantomData)),
once,
});
let bus = self.clone();
ctx.effect(move || {
{
let mut inner = bus.inner.lock().unwrap();
let list = inner.hooks.entry(key.clone()).or_default();
insert_hook(list, hook.clone(), opts.prepend);
}
Effect::Disposer(Box::new(move || {
let mut inner = bus.inner.lock().unwrap();
if let Some(list) = inner.hooks.get_mut(&key) {
retain_hook(list, &hook);
}
Ok(())
}))
})
}
fn add_wf_hook<E: Event>(
&self,
key: TypeKey,
ctx: &Ctx,
l: impl WaterfallListener<E>,
opts: EventOptions,
once: bool,
) -> Result<Disposer, CordisError> {
let hook: Arc<Hook<Arc<dyn ErasedWaterfallCall>>> = Arc::new(Hook {
call: Arc::new(WaterfallAdapter(l, std::marker::PhantomData)),
once,
});
let bus = self.clone();
ctx.effect(move || {
{
let mut inner = bus.inner.lock().unwrap();
let list = inner.wf_hooks.entry(key.clone()).or_default();
insert_hook(list, hook.clone(), opts.prepend);
}
Effect::Disposer(Box::new(move || {
let mut inner = bus.inner.lock().unwrap();
if let Some(list) = inner.wf_hooks.get_mut(&key) {
retain_hook(list, &hook);
}
Ok(())
}))
})
}
fn take_hooks(&self, key: &TypeKey) -> Vec<Arc<Hook<Arc<dyn ErasedCall>>>> {
let mut inner = self.inner.lock().unwrap();
let Some(list) = inner.hooks.get_mut(key) else {
return Vec::new();
};
claim_once(list)
}
fn take_wf_hooks(&self, key: &TypeKey) -> Vec<Arc<dyn ErasedWaterfallCall>> {
let mut inner = self.inner.lock().unwrap();
let Some(list) = inner.wf_hooks.get_mut(key) else {
return Vec::new();
};
claim_once(list)
.into_iter()
.map(|h| h.call.clone())
.collect()
}
pub fn emit<E: Event>(&self, ctx: &Ctx, e: Arc<E>) {
self.emit_keyed_inner(TypeKey::of::<E>(), ctx, e)
}
pub fn emit_keyed<E: Event>(&self, ctx: &Ctx, name: impl Into<std::sync::Arc<str>>, e: Arc<E>) {
self.emit_keyed_inner(TypeKey::keyed_dynamic::<E>(name), ctx, e)
}
fn emit_keyed_inner<E: Event>(&self, key: TypeKey, ctx: &Ctx, e: Arc<E>) {
let hooks = self.take_hooks(&key);
if hooks.is_empty() {
return; }
let ctx2 = ctx.clone();
let sink = ctx.error_sink();
let handle = ctx.handle().clone();
let mut inner = self.inner.lock().unwrap();
let prev = inner.dispatch_tail.remove(&key);
let tail = handle.spawn(async move {
if let Some(prev) = prev {
let _ = prev.await;
}
for hook in hooks {
let out = CatchUnwind::new(hook.call.call(&ctx2, &*e as &DynEvent)).await;
match out {
Ok(Ok(_)) => {}
Ok(Err(err)) => sink(Arc::new(err)),
Err(p) => sink(Arc::new(panic_error(p))),
}
}
});
inner.dispatch_tail.insert(key, tail);
}
pub async fn parallel<E: Event>(&self, ctx: &Ctx, e: Arc<E>) -> Result<(), CordisError> {
self.parallel_keyed_inner(TypeKey::of::<E>(), ctx, e).await
}
pub async fn parallel_keyed<E: Event>(
&self,
ctx: &Ctx,
name: impl Into<std::sync::Arc<str>>,
e: Arc<E>,
) -> Result<(), CordisError> {
self.parallel_keyed_inner(TypeKey::keyed_dynamic::<E>(name), ctx, e)
.await
}
async fn parallel_keyed_inner<E: Event>(
&self,
key: TypeKey,
ctx: &Ctx,
e: Arc<E>,
) -> Result<(), CordisError> {
let hooks = self.take_hooks(&key);
if hooks.is_empty() {
return Ok(());
}
let mut set = tokio::task::JoinSet::new();
for hook in hooks {
let ctx2 = ctx.clone();
let e2 = e.clone();
set.spawn_on(
async move { hook.call.call(&ctx2, &*e2 as &DynEvent).await },
ctx.handle(),
);
}
let mut errors: Vec<CordisError> = Vec::new();
while let Some(joined) = set.join_next().await {
match joined {
Ok(Ok(_)) => {}
Ok(Err(err)) => errors.push(err),
Err(join_err) => errors.push(join_panic_error(join_err)),
}
}
match crate::error::aggregate_errors(errors) {
Some(e) => Err(e),
None => Ok(()),
}
}
pub async fn serial<E: Event>(
&self,
ctx: &Ctx,
e: &E,
) -> Result<Option<E::Value>, CordisError> {
self.serial_keyed_inner(TypeKey::of::<E>(), ctx, e).await
}
pub async fn serial_keyed<E: Event>(
&self,
ctx: &Ctx,
name: impl Into<std::sync::Arc<str>>,
e: &E,
) -> Result<Option<E::Value>, CordisError> {
self.serial_keyed_inner(TypeKey::keyed_dynamic::<E>(name), ctx, e)
.await
}
async fn serial_keyed_inner<E: Event>(
&self,
key: TypeKey,
ctx: &Ctx,
e: &E,
) -> Result<Option<E::Value>, CordisError> {
for hook in self.take_hooks(&key) {
let outcome = CatchUnwind::new(hook.call.call(ctx, e as &DynEvent)).await;
match outcome {
Ok(Ok(Some(boxed))) => {
return match boxed.downcast::<E::Value>() {
Ok(v) => Ok(Some(*v)),
Err(_) => Err(CordisError::PluginFailed(
"serial value type mismatch".into(),
)),
};
}
Ok(Ok(None)) => continue,
Ok(Err(err)) => return Err(err),
Err(p) => return Err(panic_error(p)),
}
}
Ok(None)
}
pub fn waterfall<'a, E: Event, T: Terminal<E> + 'a>(
&self,
ctx: &'a Ctx,
e: &'a E,
terminal: T,
) -> BoxFuture<'a, Result<E::Value, CordisError>> {
self.waterfall_keyed_inner(TypeKey::of::<E>(), ctx, e, terminal)
}
pub fn waterfall_keyed<'a, E: Event, T: Terminal<E> + 'a>(
&self,
ctx: &'a Ctx,
name: impl Into<std::sync::Arc<str>>,
e: &'a E,
terminal: T,
) -> BoxFuture<'a, Result<E::Value, CordisError>> {
self.waterfall_keyed_inner(TypeKey::keyed_dynamic::<E>(name), ctx, e, terminal)
}
fn waterfall_keyed_inner<'a, E: Event, T: Terminal<E> + 'a>(
&self,
key: TypeKey,
ctx: &'a Ctx,
e: &'a E,
terminal: T,
) -> BoxFuture<'a, Result<E::Value, CordisError>> {
let bus = self.clone();
Box::pin(async move {
let chain = bus.take_wf_hooks(&key);
let mut terminal: Box<dyn ErasedTerminal + 'a> =
Box::new(TerminalAdapter(terminal, std::marker::PhantomData));
let next = ErasedNext {
chain: &chain,
index: 0,
ctx,
event: e as &DynEvent,
terminal: terminal.as_mut(),
};
let boxed = next.invoke().await?;
match boxed.downcast::<E::Value>() {
Ok(v) => Ok(*v),
Err(_) => Err(CordisError::PluginFailed(
"waterfall value type mismatch".into(),
)),
}
})
}
}