1use async_trait::async_trait;
15use chrono::{DateTime, Utc};
16use std::sync::Arc;
17
18use crate::error::EventError;
19use crate::event::DomainEvent;
20
21#[derive(Debug, Clone)]
28pub struct EventMetadata {
29 pub event_id: String,
31 pub timestamp: DateTime<Utc>,
33 pub correlation_id: Option<String>,
35 pub entity_type: &'static str,
37 pub aggregate_id: String,
39}
40
41impl EventMetadata {
42 pub fn new(aggregate_id: impl Into<String>, entity_type: &'static str) -> Self {
44 Self {
45 event_id: uuid::Uuid::new_v4().to_string(),
46 timestamp: Utc::now(),
47 correlation_id: None,
48 entity_type,
49 aggregate_id: aggregate_id.into(),
50 }
51 }
52
53 pub fn with_correlation(mut self, id: impl Into<String>) -> Self {
55 self.correlation_id = Some(id.into());
56 self
57 }
58}
59
60#[derive(Debug, Clone)]
67pub enum CrudEvent<E: Clone + Send + Sync + 'static> {
68 Created { entity: E, metadata: EventMetadata },
70 Updated { before: E, after: E, metadata: EventMetadata },
72 Patched { before: E, after: E, metadata: EventMetadata },
74 SoftDeleted { entity: E, metadata: EventMetadata },
76 Restored { entity: E, metadata: EventMetadata },
78 HardDeleted { entity_id: String, metadata: EventMetadata },
80 BulkCreated { entities: Vec<E>, metadata: EventMetadata },
82}
83
84impl<E: Clone + Send + Sync + 'static> CrudEvent<E> {
85 pub fn metadata(&self) -> &EventMetadata {
87 match self {
88 CrudEvent::Created { metadata, .. } => metadata,
89 CrudEvent::Updated { metadata, .. } => metadata,
90 CrudEvent::Patched { metadata, .. } => metadata,
91 CrudEvent::SoftDeleted { metadata, .. } => metadata,
92 CrudEvent::Restored { metadata, .. } => metadata,
93 CrudEvent::HardDeleted { metadata, .. } => metadata,
94 CrudEvent::BulkCreated { metadata, .. } => metadata,
95 }
96 }
97}
98
99impl<E: Clone + Send + Sync + 'static> DomainEvent for CrudEvent<E> {
102 fn event_type(&self) -> &'static str {
103 match self {
104 CrudEvent::Created { .. } => "entity.created",
105 CrudEvent::Updated { .. } => "entity.updated",
106 CrudEvent::Patched { .. } => "entity.patched",
107 CrudEvent::SoftDeleted { .. } => "entity.soft_deleted",
108 CrudEvent::Restored { .. } => "entity.restored",
109 CrudEvent::HardDeleted { .. } => "entity.hard_deleted",
110 CrudEvent::BulkCreated { .. } => "entity.bulk_created",
111 }
112 }
113
114 fn aggregate_id(&self) -> &str {
115 &self.metadata().aggregate_id
116 }
117
118 fn occurred_at(&self) -> DateTime<Utc> {
119 self.metadata().timestamp
120 }
121
122 fn aggregate_type(&self) -> &'static str {
123 self.metadata().entity_type
124 }
125}
126
127#[async_trait]
137pub trait CrudEventPublisher<E: Clone + Send + Sync + 'static>: Send + Sync {
138 async fn publish(&self, event: CrudEvent<E>) -> Result<(), EventError>;
139
140 async fn publish_many(&self, events: Vec<CrudEvent<E>>) -> Result<(), EventError> {
141 for event in events {
142 self.publish(event).await?;
143 }
144 Ok(())
145 }
146
147 async fn publish_created(&self, entity: E, _user_id: Option<String>) -> Result<(), EventError> {
149 let meta = EventMetadata::new("", std::any::type_name::<E>());
150 self.publish(CrudEvent::Created { entity, metadata: meta }).await
151 }
152
153 async fn publish_updated(&self, entity: E, _user_id: Option<String>) -> Result<(), EventError> {
155 let meta = EventMetadata::new("", std::any::type_name::<E>());
156 self.publish(CrudEvent::Updated { before: entity.clone(), after: entity, metadata: meta }).await
157 }
158
159 async fn publish_deleted(&self, entity_id: String, _user_id: Option<String>) -> Result<(), EventError> {
161 let meta = EventMetadata::new(entity_id.clone(), std::any::type_name::<E>());
162 self.publish(CrudEvent::HardDeleted { entity_id, metadata: meta }).await
163 }
164}
165
166pub struct NoOpCrudEventPublisher<E: Clone + Send + Sync + 'static> {
169 _phantom: std::marker::PhantomData<E>,
170}
171
172impl<E: Clone + Send + Sync + 'static> NoOpCrudEventPublisher<E> {
173 pub fn new() -> Self {
174 Self {
175 _phantom: std::marker::PhantomData,
176 }
177 }
178
179 pub fn arc() -> Arc<dyn CrudEventPublisher<E>> {
182 Arc::new(Self::new())
183 }
184}
185
186impl<E: Clone + Send + Sync + 'static> Default for NoOpCrudEventPublisher<E> {
187 fn default() -> Self {
188 Self::new()
189 }
190}
191
192#[async_trait]
193impl<E: Clone + Send + Sync + 'static> CrudEventPublisher<E> for NoOpCrudEventPublisher<E> {
194 async fn publish(&self, _event: CrudEvent<E>) -> Result<(), EventError> {
195 Ok(())
196 }
197
198 async fn publish_many(&self, _events: Vec<CrudEvent<E>>) -> Result<(), EventError> {
199 Ok(())
200 }
201}
202
203#[cfg(test)]
204mod tests {
205 use super::*;
206
207 #[derive(Debug, Clone)]
208 struct FakeEntity {
209 id: String,
210 }
211
212 fn meta(id: &str) -> EventMetadata {
213 EventMetadata::new(id.to_string(), "fake_entity")
214 }
215
216 #[tokio::test]
217 async fn noop_publisher_returns_ok_for_all_variants() {
218 let publisher = NoOpCrudEventPublisher::<FakeEntity>::new();
219
220 let events = vec![
221 CrudEvent::Created {
222 entity: FakeEntity { id: "1".into() },
223 metadata: meta("1"),
224 },
225 CrudEvent::SoftDeleted {
226 entity: FakeEntity { id: "1".into() },
227 metadata: meta("1"),
228 },
229 CrudEvent::HardDeleted {
230 entity_id: "1".into(),
231 metadata: meta("1"),
232 },
233 ];
234
235 for event in events {
236 assert!(publisher.publish(event).await.is_ok());
237 }
238 }
239
240 #[test]
241 fn event_type_names_are_stable() {
242 let created: CrudEvent<FakeEntity> = CrudEvent::Created {
243 entity: FakeEntity { id: "1".into() },
244 metadata: meta("1"),
245 };
246 assert_eq!(created.event_type(), "entity.created");
247
248 let deleted: CrudEvent<FakeEntity> = CrudEvent::HardDeleted {
249 entity_id: "1".into(),
250 metadata: meta("1"),
251 };
252 assert_eq!(deleted.event_type(), "entity.hard_deleted");
253 }
254
255 #[test]
256 fn aggregate_id_roundtrips() {
257 let event: CrudEvent<FakeEntity> = CrudEvent::Created {
258 entity: FakeEntity { id: "abc".into() },
259 metadata: meta("abc"),
260 };
261 assert_eq!(event.aggregate_id(), "abc");
262 }
263
264 #[test]
265 fn metadata_correlation_id() {
266 let m = EventMetadata::new("e1", "entity").with_correlation("corr-42");
267 assert_eq!(m.correlation_id.as_deref(), Some("corr-42"));
268 assert_eq!(m.aggregate_id, "e1");
269 assert_eq!(m.entity_type, "entity");
270 }
271}