use std::any::TypeId;
use std::sync::Arc;
use arrayvec::ArrayVec;
use rustc_hash::FxHashMap;
use crate::ctx::Ctx;
use crate::error::Result;
use crate::monitor::async_handler::BoxFuture;
pub const MAX_EVENT_TYPES: usize = u16::MAX as usize;
const INLINE_EVENT_TYPES: usize = 16;
#[allow(clippy::large_enum_variant)]
#[derive(Clone)]
pub(crate) enum TypeSlotTable {
Inline(ArrayVec<(TypeId, u16), INLINE_EVENT_TYPES>),
Spilled(FxHashMap<TypeId, u16>),
}
impl TypeSlotTable {
pub(crate) fn new() -> Self {
TypeSlotTable::Inline(ArrayVec::new())
}
#[inline]
pub(crate) fn get(&self, ty: TypeId) -> Option<u16> {
match self {
TypeSlotTable::Inline(v) => v.iter().copied().find(|(t, _)| *t == ty).map(|(_, i)| i),
TypeSlotTable::Spilled(m) => m.get(&ty).copied(),
}
}
pub(crate) fn insert(&mut self, ty: TypeId, idx: u16) {
match self {
TypeSlotTable::Inline(v) => {
if v.try_push((ty, idx)).is_err() {
let mut m: FxHashMap<TypeId, u16> = v.iter().copied().collect();
m.insert(ty, idx);
*self = TypeSlotTable::Spilled(m);
}
}
TypeSlotTable::Spilled(m) => {
m.insert(ty, idx);
}
}
}
pub(crate) fn len(&self) -> usize {
match self {
TypeSlotTable::Inline(v) => v.len(),
TypeSlotTable::Spilled(m) => m.len(),
}
}
}
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) trait DynEffectHandler: Send + Sync {
fn call(
&self,
ptr: *const (),
ctx: &Ctx<'_>,
) -> BoxFuture<Result<crate::monitor::effect::Effects>>;
}
pub(crate) type BoxedEffectHandler = Arc<dyn DynEffectHandler>;
fn call_handler_catching(handler: &BoxedHandler, ptr: *const (), ctx: &mut Ctx<'_>) -> Result<()> {
use std::panic::{AssertUnwindSafe, catch_unwind};
match catch_unwind(AssertUnwindSafe(|| handler(ptr, ctx))) {
Ok(res) => res,
Err(payload) => Err(crate::error::Error::HandlerPanic(panic_message(&*payload))),
}
}
#[cold]
fn panic_message(payload: &(dyn std::any::Any + Send)) -> String {
if let Some(s) = payload.downcast_ref::<&str>() {
(*s).to_string()
} else if let Some(s) = payload.downcast_ref::<String>() {
s.clone()
} else {
"handler panicked (non-string payload)".to_string()
}
}
pub(crate) struct HandlerSlot {
pub(crate) handler: BoxedHandler,
}
pub(crate) struct AsyncHandlerSlot {
pub(crate) handler: BoxedAsyncHandler,
}
pub(crate) struct EffectHandlerSlot {
pub(crate) handler: BoxedEffectHandler,
}
pub struct Dispatcher {
slot_by_type: TypeSlotTable,
slots: Box<[Vec<HandlerSlot>]>,
async_slots: Box<[Vec<AsyncHandlerSlot>]>,
effect_slots: Box<[Vec<EffectHandlerSlot>]>,
#[allow(dead_code)] slot_types: Box<[TypeId]>,
catch_panics: bool,
}
impl Dispatcher {
pub(crate) fn new(
slot_by_type: TypeSlotTable,
slots: Box<[Vec<HandlerSlot>]>,
async_slots: Box<[Vec<AsyncHandlerSlot>]>,
effect_slots: Box<[Vec<EffectHandlerSlot>]>,
slot_types: Box<[TypeId]>,
) -> Self {
Self {
slot_by_type,
slots,
async_slots,
effect_slots,
slot_types,
catch_panics: false,
}
}
pub(crate) fn set_catch_panics(&mut self, on: bool) {
self.catch_panics = on;
}
#[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.get(target) else {
return Ok(());
};
#[cfg(debug_assertions)]
debug_assert_eq!(
self.slot_types[slot_idx as usize], target,
"dispatcher slot/type desync — type-erased cast would be UB"
);
let ptr = payload as *const P as *const ();
let catch = self.catch_panics;
for slot in &mut self.slots[slot_idx as usize] {
if catch {
call_handler_catching(&slot.handler, ptr, ctx)?;
} else {
(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.get(target) else {
return Ok(());
};
#[cfg(debug_assertions)]
debug_assert_eq!(
self.slot_types[slot_idx as usize], target,
"dispatcher slot/type desync — type-erased cast would be UB"
);
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(())
}
}
}
pub async fn dispatch_effects<P: 'static>(
&mut self,
payload: &P,
ctx: &mut Ctx<'_>,
) -> Result<()> {
let target = TypeId::of::<P>();
let Some(slot_idx) = self.slot_by_type.get(target) else {
return Ok(());
};
#[cfg(debug_assertions)]
debug_assert_eq!(
self.slot_types[slot_idx as usize], target,
"dispatcher slot/type desync — type-erased cast would be UB"
);
if self.effect_slots[slot_idx as usize].is_empty() {
return Ok(());
}
let mut futures: Vec<BoxFuture<Result<crate::monitor::effect::Effects>>> =
Vec::with_capacity(self.effect_slots[slot_idx as usize].len());
{
let ptr = payload as *const P as *const ();
for slot in &self.effect_slots[slot_idx as usize] {
futures.push(slot.handler.call(ptr, &*ctx));
}
}
for fut in futures {
let effects = fut.await?;
if !effects.is_empty() {
effects.apply(&mut *ctx.sink);
}
}
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(),
effect_slots: self
.effect_slots
.iter()
.map(|v| {
v.iter()
.map(|s| EffectHandlerSlot {
handler: Arc::clone(&s.handler),
})
.collect()
})
.collect::<Vec<_>>()
.into_boxed_slice(),
slot_types: self.slot_types.clone(),
catch_panics: self.catch_panics,
}
}
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()
}
pub fn effect_handler_count(&self) -> usize {
self.effect_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 = TypeSlotTable::new();
slot_by_type.insert(TypeId::of::<u32>(), 0);
let slots = vec![vec![HandlerSlot { handler }]].into_boxed_slice();
let async_slots = vec![vec![]].into_boxed_slice();
let effect_slots = vec![vec![]].into_boxed_slice();
let slot_types = vec![TypeId::of::<u32>()].into_boxed_slice();
let primary = Dispatcher::new(slot_by_type, slots, async_slots, effect_slots, slot_types);
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 type_slot_table_stays_inline_then_spills_preserving_lookups() {
let mut t = TypeSlotTable::new();
let tys: [TypeId; 4] = [
TypeId::of::<u8>(),
TypeId::of::<u16>(),
TypeId::of::<u32>(),
TypeId::of::<u64>(),
];
for (i, ty) in tys.iter().enumerate() {
t.insert(*ty, i as u16);
}
assert!(matches!(t, TypeSlotTable::Inline(_)), "≤16 stays inline");
for (i, ty) in tys.iter().enumerate() {
assert_eq!(t.get(*ty), Some(i as u16));
}
assert_eq!(t.get(TypeId::of::<i8>()), None);
assert_eq!(t.len(), 4);
let mut big = TypeSlotTable::new();
fn synth_ty(n: usize) -> TypeId {
const TYS: [fn() -> TypeId; 17] = [
|| TypeId::of::<[u8; 0]>(),
|| TypeId::of::<[u8; 1]>(),
|| TypeId::of::<[u8; 2]>(),
|| TypeId::of::<[u8; 3]>(),
|| TypeId::of::<[u8; 4]>(),
|| TypeId::of::<[u8; 5]>(),
|| TypeId::of::<[u8; 6]>(),
|| TypeId::of::<[u8; 7]>(),
|| TypeId::of::<[u8; 8]>(),
|| TypeId::of::<[u8; 9]>(),
|| TypeId::of::<[u8; 10]>(),
|| TypeId::of::<[u8; 11]>(),
|| TypeId::of::<[u8; 12]>(),
|| TypeId::of::<[u8; 13]>(),
|| TypeId::of::<[u8; 14]>(),
|| TypeId::of::<[u8; 15]>(),
|| TypeId::of::<[u8; 16]>(),
];
TYS[n]()
}
for i in 0..17 {
big.insert(synth_ty(i), i as u16);
}
assert!(matches!(big, TypeSlotTable::Spilled(_)), "17 spills to map");
assert_eq!(big.len(), 17);
assert_eq!(big.get(synth_ty(0)), Some(0));
assert_eq!(big.get(synth_ty(16)), Some(16));
}
#[test]
fn empty_dispatch_is_noop() {
let mut d = Dispatcher::new(
TypeSlotTable::new(),
Vec::new().into_boxed_slice(),
Vec::new().into_boxed_slice(),
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 = TypeSlotTable::new();
slot_by_type.insert(TypeId::of::<u32>(), 0);
slot_by_type.insert(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 effect_slots: Box<[Vec<EffectHandlerSlot>]> =
vec![Vec::new(), Vec::new()].into_boxed_slice();
let slot_types = vec![TypeId::of::<u32>(), TypeId::of::<u64>()].into_boxed_slice();
let mut d = Dispatcher::new(slot_by_type, slots, async_slots, effect_slots, slot_types);
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);
}
#[test]
fn catch_panics_converts_a_handler_panic_into_handler_panic_error() {
let handler: BoxedHandler = Arc::new(|_ptr, _ctx| panic!("boom in handler"));
let mut slot_by_type = TypeSlotTable::new();
slot_by_type.insert(TypeId::of::<u32>(), 0);
let slots = vec![vec![HandlerSlot { handler }]].into_boxed_slice();
let async_slots = vec![vec![]].into_boxed_slice();
let effect_slots = vec![vec![]].into_boxed_slice();
let slot_types = vec![TypeId::of::<u32>()].into_boxed_slice();
let mut d = Dispatcher::new(slot_by_type, slots, async_slots, effect_slots, slot_types);
d.set_catch_panics(true);
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);
let payload: u32 = 1;
let prev = std::panic::take_hook();
std::panic::set_hook(Box::new(|_| {}));
let res = d.dispatch::<u32>(&payload, &mut ctx);
std::panic::set_hook(prev);
match res {
Err(crate::error::Error::HandlerPanic(msg)) => {
assert!(msg.contains("boom in handler"), "msg = {msg}");
}
other => panic!("expected HandlerPanic, got {other:?}"),
}
}
}