use std::any::TypeId;
use std::sync::Arc;
use arrayvec::ArrayVec;
use crate::ctx::Ctx;
use crate::error::Result;
use crate::monitor::async_handler::BoxFuture;
pub const MAX_EVENT_TYPES: usize = 16;
pub(crate) type BoxedHandler = Arc<dyn Fn(*const (), &mut Ctx<'_>) -> Result<()> + Send + Sync>;
pub(crate) trait DynAsyncHandler: Send + Sync {
fn call(&self, ptr: *const ()) -> BoxFuture<Result<()>>;
}
pub(crate) type BoxedAsyncHandler = Arc<dyn DynAsyncHandler>;
pub(crate) struct HandlerSlot {
pub(crate) handler: BoxedHandler,
}
pub(crate) struct AsyncHandlerSlot {
pub(crate) handler: BoxedAsyncHandler,
}
pub struct Dispatcher {
slot_by_type: ArrayVec<(TypeId, u8), MAX_EVENT_TYPES>,
slots: Box<[Vec<HandlerSlot>]>,
async_slots: Box<[Vec<AsyncHandlerSlot>]>,
}
impl Dispatcher {
pub(crate) fn new(
slot_by_type: ArrayVec<(TypeId, u8), MAX_EVENT_TYPES>,
slots: Box<[Vec<HandlerSlot>]>,
async_slots: Box<[Vec<AsyncHandlerSlot>]>,
) -> Self {
Self {
slot_by_type,
slots,
async_slots,
}
}
#[inline]
pub fn dispatch<P: 'static>(&mut self, payload: &P, ctx: &mut Ctx<'_>) -> Result<()> {
let target = TypeId::of::<P>();
let Some((_, slot_idx)) = self
.slot_by_type
.iter()
.copied()
.find(|(t, _)| *t == target)
else {
return Ok(());
};
let ptr = payload as *const P as *const ();
for slot in &mut self.slots[slot_idx as usize] {
(slot.handler)(ptr, ctx)?;
}
Ok(())
}
pub async fn dispatch_async<P: 'static>(&mut self, payload: &P) -> Result<()> {
let target = TypeId::of::<P>();
let Some((_, slot_idx)) = self
.slot_by_type
.iter()
.copied()
.find(|(t, _)| *t == target)
else {
return Ok(());
};
let slots = &self.async_slots[slot_idx as usize];
match slots.len() {
0 => Ok(()),
1 => {
let fut = {
let ptr = payload as *const P as *const ();
slots[0].handler.call(ptr)
};
fut.await
}
_ => {
let mut futures: Vec<BoxFuture<Result<()>>> = Vec::with_capacity(slots.len());
{
let ptr = payload as *const P as *const ();
for slot in slots.iter() {
futures.push(slot.handler.call(ptr));
}
}
for fut in futures {
fut.await?;
}
Ok(())
}
}
}
#[allow(dead_code)]
pub(crate) fn clone_for_shard(&self) -> Self {
Self {
slot_by_type: self.slot_by_type.clone(),
slots: self
.slots
.iter()
.map(|v| {
v.iter()
.map(|s| HandlerSlot {
handler: Arc::clone(&s.handler),
})
.collect()
})
.collect::<Vec<_>>()
.into_boxed_slice(),
async_slots: self
.async_slots
.iter()
.map(|v| {
v.iter()
.map(|s| AsyncHandlerSlot {
handler: Arc::clone(&s.handler),
})
.collect()
})
.collect::<Vec<_>>()
.into_boxed_slice(),
}
}
pub fn type_count(&self) -> usize {
self.slot_by_type.len()
}
pub fn handler_count(&self) -> usize {
self.slots.iter().map(|s| s.len()).sum()
}
pub fn async_handler_count(&self) -> usize {
self.async_slots.iter().map(|s| s.len()).sum()
}
}
impl std::fmt::Debug for Dispatcher {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Dispatcher")
.field("type_count", &self.type_count())
.field("handler_count", &self.handler_count())
.finish()
}
}
#[cfg(test)]
mod tests {
use flowscope::Timestamp;
use super::*;
use crate::anomaly::sink::NoopSink;
use crate::ctx::{CounterRegistry, SourceIdx, StateMap};
fn fresh_ctx<'a>(
state: &'a mut StateMap,
sink: &'a mut NoopSink,
counters: &'a mut CounterRegistry,
flow_states: &'a mut crate::ctx::FlowStateRegistry,
) -> Ctx<'a> {
Ctx {
flow: None,
ts: Timestamp::new(0, 0),
source: SourceIdx(0),
monitor_name: None,
state_map: state,
sink,
counters,
flow_states,
label_table: crate::ctx::default_label_table(),
tracker: None,
}
}
#[test]
fn clone_for_shard_produces_independent_dispatcher_sharing_handlers() {
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicU32, Ordering};
let count = StdArc::new(AtomicU32::new(0));
let count_h = StdArc::clone(&count);
let handler: BoxedHandler = Arc::new(move |_ptr, _ctx| {
count_h.fetch_add(1, Ordering::Relaxed);
Ok(())
});
let mut slot_by_type = ArrayVec::new();
slot_by_type.push((TypeId::of::<u32>(), 0));
let slots = vec![vec![HandlerSlot { handler }]].into_boxed_slice();
let async_slots = vec![vec![]].into_boxed_slice();
let primary = Dispatcher::new(slot_by_type, slots, async_slots);
let mut shard = primary.clone_for_shard();
let mut primary = primary; assert_eq!(primary.type_count(), 1);
assert_eq!(shard.type_count(), 1);
assert_eq!(primary.handler_count(), 1);
assert_eq!(shard.handler_count(), 1);
let payload: u32 = 42;
let mut s = StateMap::default();
let mut sink = NoopSink;
let mut c = CounterRegistry::default();
let mut fs = crate::ctx::FlowStateRegistry::default();
let mut ctx = fresh_ctx(&mut s, &mut sink, &mut c, &mut fs);
primary.dispatch::<u32>(&payload, &mut ctx).unwrap();
shard.dispatch::<u32>(&payload, &mut ctx).unwrap();
shard.dispatch::<u32>(&payload, &mut ctx).unwrap();
assert_eq!(count.load(Ordering::Relaxed), 3);
}
#[test]
fn empty_dispatch_is_noop() {
let mut d = Dispatcher::new(
ArrayVec::new(),
Vec::new().into_boxed_slice(),
Vec::new().into_boxed_slice(),
);
let mut s = StateMap::default();
let mut k = NoopSink;
let mut c = CounterRegistry::default();
let mut fs = crate::ctx::FlowStateRegistry::default();
let mut ctx = fresh_ctx(&mut s, &mut k, &mut c, &mut fs);
let payload: u32 = 7;
assert!(d.dispatch::<u32>(&payload, &mut ctx).is_ok());
assert_eq!(d.type_count(), 0);
assert_eq!(d.handler_count(), 0);
}
#[test]
fn dispatch_routes_to_matching_slot_only() {
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
let u32_count = Arc::new(AtomicU32::new(0));
let u64_count = Arc::new(AtomicU32::new(0));
let u32_count_h = Arc::clone(&u32_count);
let u64_count_h = Arc::clone(&u64_count);
let u32_handler: BoxedHandler = Arc::new(move |ptr, _ctx| {
let _val: u32 = unsafe { *(ptr as *const u32) };
u32_count_h.fetch_add(1, Ordering::Relaxed);
Ok(())
});
let u64_handler: BoxedHandler = Arc::new(move |ptr, _ctx| {
let _val: u64 = unsafe { *(ptr as *const u64) };
u64_count_h.fetch_add(1, Ordering::Relaxed);
Ok(())
});
let mut slot_by_type = ArrayVec::new();
slot_by_type.push((TypeId::of::<u32>(), 0));
slot_by_type.push((TypeId::of::<u64>(), 1));
let slots: Box<[Vec<HandlerSlot>]> = vec![
vec![HandlerSlot {
handler: u32_handler,
}],
vec![HandlerSlot {
handler: u64_handler,
}],
]
.into_boxed_slice();
let async_slots: Box<[Vec<AsyncHandlerSlot>]> =
vec![Vec::new(), Vec::new()].into_boxed_slice();
let mut d = Dispatcher::new(slot_by_type, slots, async_slots);
let mut s = StateMap::default();
let mut k = NoopSink;
let mut c = CounterRegistry::default();
let mut fs = crate::ctx::FlowStateRegistry::default();
let mut ctx = fresh_ctx(&mut s, &mut k, &mut c, &mut fs);
let u32_payload: u32 = 7;
d.dispatch::<u32>(&u32_payload, &mut ctx).unwrap();
assert_eq!(u32_count.load(Ordering::Relaxed), 1);
assert_eq!(u64_count.load(Ordering::Relaxed), 0);
let u64_payload: u64 = 13;
d.dispatch::<u64>(&u64_payload, &mut ctx).unwrap();
assert_eq!(u32_count.load(Ordering::Relaxed), 1);
assert_eq!(u64_count.load(Ordering::Relaxed), 1);
assert_eq!(d.type_count(), 2);
assert_eq!(d.handler_count(), 2);
}
}