Skip to main content

backbone_messaging/
handler.rs

1//! Event Handler trait definition
2
3use async_trait::async_trait;
4
5use crate::{DomainEvent, EventEnvelope, EventError};
6
7/// Trait for event handlers
8///
9/// Event handlers process domain events asynchronously. They can be used
10/// for side effects like sending notifications, updating read models,
11/// or triggering workflows.
12///
13/// # Example
14///
15/// ```rust,ignore
16/// use backbone_messaging::{EventHandler, EventEnvelope, EventError, DomainEvent};
17/// use async_trait::async_trait;
18///
19/// struct EmailNotificationHandler;
20///
21/// #[async_trait]
22/// impl<E: DomainEvent> EventHandler<E> for EmailNotificationHandler {
23///     async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError> {
24///         // Send email notification
25///         println!("Event {} occurred for {}", envelope.event_type, envelope.aggregate_id);
26///         Ok(())
27///     }
28///
29///     fn event_types(&self) -> Vec<&'static str> {
30///         vec!["UserCreated", "OrderPlaced"]
31///     }
32/// }
33/// ```
34#[async_trait]
35pub trait EventHandler<E: DomainEvent>: Send + Sync {
36    /// Handle a domain event
37    ///
38    /// This method is called when an event matching one of the types
39    /// returned by `event_types()` is published.
40    async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError>;
41
42    /// Event types this handler is interested in
43    ///
44    /// Return an empty slice to receive all events.
45    fn event_types(&self) -> Vec<&'static str>;
46
47    /// Handler name for logging and debugging
48    fn name(&self) -> &'static str {
49        std::any::type_name::<Self>()
50    }
51
52    /// Whether this handler should be retried on failure
53    fn should_retry(&self) -> bool {
54        true
55    }
56
57    /// Maximum retry attempts
58    fn max_retries(&self) -> u32 {
59        3
60    }
61}
62
63/// A handler that logs all events for debugging
64pub struct LoggingHandler {
65    event_types: Vec<&'static str>,
66}
67
68impl LoggingHandler {
69    /// Create a handler that logs specific event types
70    pub fn new(event_types: Vec<&'static str>) -> Self {
71        Self { event_types }
72    }
73
74    /// Create a handler that logs all events
75    pub fn all() -> Self {
76        Self { event_types: vec![] }
77    }
78}
79
80#[async_trait]
81impl<E: DomainEvent> EventHandler<E> for LoggingHandler {
82    async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError> {
83        tracing::info!(
84            event_type = %envelope.event_type,
85            aggregate_id = %envelope.aggregate_id,
86            event_id = %envelope.id,
87            correlation_id = ?envelope.correlation_id,
88            "Domain event received"
89        );
90        Ok(())
91    }
92
93    fn event_types(&self) -> Vec<&'static str> {
94        self.event_types.clone()
95    }
96
97    fn name(&self) -> &'static str {
98        "LoggingHandler"
99    }
100}
101
102/// A handler that collects events for testing
103#[derive(Default)]
104pub struct CollectingHandler<E: DomainEvent> {
105    events: std::sync::Arc<tokio::sync::RwLock<Vec<EventEnvelope<E>>>>,
106}
107
108impl<E: DomainEvent> CollectingHandler<E> {
109    /// Create a new collecting handler
110    pub fn new() -> Self {
111        Self {
112            events: std::sync::Arc::new(tokio::sync::RwLock::new(Vec::new())),
113        }
114    }
115
116    /// Get all collected events
117    pub async fn events(&self) -> Vec<EventEnvelope<E>> {
118        self.events.read().await.clone()
119    }
120
121    /// Clear collected events
122    pub async fn clear(&self) {
123        self.events.write().await.clear();
124    }
125
126    /// Get event count
127    pub async fn count(&self) -> usize {
128        self.events.read().await.len()
129    }
130}
131
132#[async_trait]
133impl<E: DomainEvent> EventHandler<E> for CollectingHandler<E> {
134    async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError> {
135        self.events.write().await.push(envelope);
136        Ok(())
137    }
138
139    fn event_types(&self) -> Vec<&'static str> {
140        vec![] // Collect all events
141    }
142
143    fn name(&self) -> &'static str {
144        "CollectingHandler"
145    }
146}