use async_trait::async_trait;
use crate::{DomainEvent, EventEnvelope, EventError};
#[async_trait]
pub trait EventHandler<E: DomainEvent>: Send + Sync {
async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError>;
fn event_types(&self) -> Vec<&'static str>;
fn name(&self) -> &'static str {
std::any::type_name::<Self>()
}
fn should_retry(&self) -> bool {
true
}
fn max_retries(&self) -> u32 {
3
}
}
pub struct LoggingHandler {
event_types: Vec<&'static str>,
}
impl LoggingHandler {
pub fn new(event_types: Vec<&'static str>) -> Self {
Self { event_types }
}
pub fn all() -> Self {
Self { event_types: vec![] }
}
}
#[async_trait]
impl<E: DomainEvent> EventHandler<E> for LoggingHandler {
async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError> {
tracing::info!(
event_type = %envelope.event_type,
aggregate_id = %envelope.aggregate_id,
event_id = %envelope.id,
correlation_id = ?envelope.correlation_id,
"Domain event received"
);
Ok(())
}
fn event_types(&self) -> Vec<&'static str> {
self.event_types.clone()
}
fn name(&self) -> &'static str {
"LoggingHandler"
}
}
#[derive(Default)]
pub struct CollectingHandler<E: DomainEvent> {
events: std::sync::Arc<tokio::sync::RwLock<Vec<EventEnvelope<E>>>>,
}
impl<E: DomainEvent> CollectingHandler<E> {
pub fn new() -> Self {
Self {
events: std::sync::Arc::new(tokio::sync::RwLock::new(Vec::new())),
}
}
pub async fn events(&self) -> Vec<EventEnvelope<E>> {
self.events.read().await.clone()
}
pub async fn clear(&self) {
self.events.write().await.clear();
}
pub async fn count(&self) -> usize {
self.events.read().await.len()
}
}
#[async_trait]
impl<E: DomainEvent> EventHandler<E> for CollectingHandler<E> {
async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError> {
self.events.write().await.push(envelope);
Ok(())
}
fn event_types(&self) -> Vec<&'static str> {
vec![] }
fn name(&self) -> &'static str {
"CollectingHandler"
}
}