use std::collections::HashMap;
use std::sync::Arc;
use serde::Deserialize;
use tokio::sync::mpsc;
type EventHandlerResult = Result<(), Box<dyn std::error::Error + Send + Sync>>;
#[derive(Debug, Deserialize)]
struct RawEventEnvelope {
header: RawEventHeader,
}
#[derive(Debug, Deserialize)]
struct RawEventHeader {
#[serde(default)]
event_type: String,
}
pub trait EventHandler: Send + Sync + 'static {
fn handle(&self, payload: &[u8]) -> EventHandlerResult;
}
#[derive(Clone)]
pub struct EventDispatcherHandler {
payload_tx: Option<mpsc::UnboundedSender<Vec<u8>>>,
raw_handlers: HashMap<String, Arc<dyn EventHandler>>,
}
impl std::fmt::Debug for EventDispatcherHandler {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("EventDispatcherHandler")
.field(
"payload_tx",
&self.payload_tx.as_ref().map(|_| "configured"),
)
.field(
"raw_handler_keys",
&self.raw_handlers.keys().collect::<Vec<_>>(),
)
.finish()
}
}
impl EventDispatcherHandler {
pub const RAW_EVENT_KEY: &'static str = "raw";
pub fn builder() -> Self {
Self {
payload_tx: None,
raw_handlers: HashMap::new(),
}
}
pub fn build(self) -> Self {
self
}
pub fn payload_sender(mut self, payload_tx: mpsc::UnboundedSender<Vec<u8>>) -> Self {
self.payload_tx = Some(payload_tx);
self
}
pub fn register_raw<S, H>(mut self, key: S, handler: H) -> Result<Self, String>
where
S: Into<String>,
H: EventHandler,
{
let key = key.into();
if key.trim().is_empty() {
return Err("processor key cannot be empty".to_string());
}
if self.raw_handlers.contains_key(&key) {
return Err(format!("processor already registered, type: {key}"));
}
self.raw_handlers.insert(key, Arc::new(handler));
Ok(self)
}
fn extract_event_type(payload: &[u8]) -> Option<String> {
serde_json::from_slice::<RawEventEnvelope>(payload)
.ok()
.map(|event| event.header.event_type)
.filter(|event_type| !event_type.trim().is_empty())
}
fn dispatch_raw_handler(&self, key: &str, payload: &[u8]) -> Result<(), String> {
if let Some(handler) = self.raw_handlers.get(key) {
handler
.handle(payload)
.map_err(|err| format!("处理原始事件 {key} 失败: {err}"))?;
}
Ok(())
}
pub fn do_without_validation(&self, payload: &[u8]) -> Result<(), String> {
if let Some(payload_tx) = &self.payload_tx {
payload_tx
.send(payload.to_vec())
.map_err(|e| format!("转发事件负载失败: {e}"))?;
}
if let Some(event_type) = Self::extract_event_type(payload) {
self.dispatch_raw_handler(&event_type, payload)?;
}
self.dispatch_raw_handler(Self::RAW_EVENT_KEY, payload)?;
Ok(())
}
}