backbone_messaging/
handler.rs1use async_trait::async_trait;
4
5use crate::{DomainEvent, EventEnvelope, EventError};
6
7#[async_trait]
35pub trait EventHandler<E: DomainEvent>: Send + Sync {
36 async fn handle(&self, envelope: EventEnvelope<E>) -> Result<(), EventError>;
41
42 fn event_types(&self) -> Vec<&'static str>;
46
47 fn name(&self) -> &'static str {
49 std::any::type_name::<Self>()
50 }
51
52 fn should_retry(&self) -> bool {
54 true
55 }
56
57 fn max_retries(&self) -> u32 {
59 3
60 }
61}
62
63pub struct LoggingHandler {
65 event_types: Vec<&'static str>,
66}
67
68impl LoggingHandler {
69 pub fn new(event_types: Vec<&'static str>) -> Self {
71 Self { event_types }
72 }
73
74 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#[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 pub fn new() -> Self {
111 Self {
112 events: std::sync::Arc::new(tokio::sync::RwLock::new(Vec::new())),
113 }
114 }
115
116 pub async fn events(&self) -> Vec<EventEnvelope<E>> {
118 self.events.read().await.clone()
119 }
120
121 pub async fn clear(&self) {
123 self.events.write().await.clear();
124 }
125
126 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![] }
142
143 fn name(&self) -> &'static str {
144 "CollectingHandler"
145 }
146}