use super::*;
use crate::ctx::Ctx;
use crate::error::CordisError;
use crate::event::{Event, Listener, Next};
use crate::BoxFuture;
use std::time::Duration;
struct Ping;
impl Event for Ping {
const NAME: &'static str = "test::TransientPing";
type Value = ();
}
struct Nop;
impl Listener<Ping> for Nop {
fn call<'a>(
&'a self,
_ctx: &'a Ctx,
_e: &'a Ping,
) -> BoxFuture<'a, Result<Option<()>, CordisError>> {
Box::pin(async { Ok(None) })
}
}
struct NopWaterfall;
impl crate::event::WaterfallListener<Ping> for NopWaterfall {
fn call<'a>(
&'a self,
_ctx: &'a Ctx,
_e: &'a Ping,
next: Next<'a, Ping>,
) -> BoxFuture<'a, Result<(), CordisError>> {
Box::pin(async move { next.call().await })
}
}
#[tokio::test]
async fn keyed_channels_and_dispatch_tails_prune() {
let ctx = Ctx::root().expect("runtime in scope");
let bus = ctx.events().clone();
for i in 0..25 {
let name = format!("ch/{i}");
let listener = bus.on_keyed::<Ping>(&ctx, name.clone(), Nop).unwrap();
let wf = bus
.on_waterfall_keyed::<Ping>(&ctx, name.clone(), NopWaterfall)
.unwrap();
bus.emit_keyed::<Ping>(&ctx, name.clone(), Arc::new(Ping));
listener.dispose().await.unwrap();
wf.dispose().await.unwrap();
}
for _ in 0..500 {
let empty = {
let inner = bus.inner.lock().unwrap();
inner.hooks.is_empty() && inner.wf_hooks.is_empty() && inner.dispatch_tail.is_empty()
};
if empty {
break;
}
tokio::time::sleep(Duration::from_millis(2)).await;
}
{
let inner = bus.inner.lock().unwrap();
assert!(
inner.hooks.is_empty() && inner.wf_hooks.is_empty(),
"keyed channels must drop empty listener lists"
);
assert!(
inner.dispatch_tail.is_empty(),
"dispatch tails must self-remove on completion"
);
}
let listener = bus.on_keyed::<Ping>(&ctx, "ch/0", Nop).unwrap();
bus.emit_keyed::<Ping>(&ctx, "ch/0", Arc::new(Ping));
listener.dispose().await.unwrap();
}