Skip to main content

backbone_messaging/
crud_event.rs

1//! CRUD domain events — generic event types for all entity operations.
2//!
3//! Every generated entity gets a type alias over these generics:
4//!
5//! ```rust,ignore
6//! // Generated (not hand-written):
7//! pub type StoredFileCrudEvent = CrudEvent<StoredFile>;
8//! pub type StoredFileCrudEventPublisher = Arc<dyn CrudEventPublisher<StoredFile>>;
9//! ```
10//!
11//! Custom services can inject `Arc<dyn CrudEventPublisher<E>>` directly
12//! or use `NoOpCrudEventPublisher<E>` when no external subscribers exist.
13
14use async_trait::async_trait;
15use chrono::{DateTime, Utc};
16use std::sync::Arc;
17
18use crate::error::EventError;
19use crate::event::DomainEvent;
20
21// ─── EventMetadata ────────────────────────────────────────────────────────────
22
23/// Infrastructure metadata attached to every CRUD event.
24///
25/// Carries correlation tracking, timestamps, and entity identification
26/// without polluting the entity type itself.
27#[derive(Debug, Clone)]
28pub struct EventMetadata {
29    /// Unique ID for this event instance (UUID v4).
30    pub event_id: String,
31    /// When the event was created (UTC).
32    pub timestamp: DateTime<Utc>,
33    /// Correlation ID for distributed tracing across services.
34    pub correlation_id: Option<String>,
35    /// The static entity type name (e.g. `"order"`, `"stored_file"`).
36    pub entity_type: &'static str,
37    /// The entity's aggregate root ID at the time of the event.
38    pub aggregate_id: String,
39}
40
41impl EventMetadata {
42    /// Construct metadata for an event on `aggregate_id` of type `entity_type`.
43    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    /// Attach a correlation ID (builder).
54    pub fn with_correlation(mut self, id: impl Into<String>) -> Self {
55        self.correlation_id = Some(id.into());
56        self
57    }
58}
59
60// ─── CrudEvent ────────────────────────────────────────────────────────────────
61
62/// The standard set of CRUD domain events for any entity type `E`.
63///
64/// All variants carry the full entity snapshot so subscribers never
65/// need to fetch back to the database.
66#[derive(Debug, Clone)]
67pub enum CrudEvent<E: Clone + Send + Sync + 'static> {
68    /// Entity was created for the first time.
69    Created { entity: E, metadata: EventMetadata },
70    /// Entity fields were updated (full replacement).
71    Updated { before: E, after: E, metadata: EventMetadata },
72    /// Entity was partially patched.
73    Patched { before: E, after: E, metadata: EventMetadata },
74    /// Entity was soft-deleted (moved to trash).
75    SoftDeleted { entity: E, metadata: EventMetadata },
76    /// Soft-deleted entity was restored.
77    Restored { entity: E, metadata: EventMetadata },
78    /// Entity was permanently deleted.
79    HardDeleted { entity_id: String, metadata: EventMetadata },
80    /// Bulk create completed.
81    BulkCreated { entities: Vec<E>, metadata: EventMetadata },
82}
83
84impl<E: Clone + Send + Sync + 'static> CrudEvent<E> {
85    /// Access the event metadata from any variant.
86    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
99// ─── DomainEvent impl ─────────────────────────────────────────────────────────
100
101impl<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// ─── CrudEventPublisher ───────────────────────────────────────────────────────
128
129/// Contract for publishing `CrudEvent<E>` instances.
130///
131/// Implemented by:
132/// - `NoOpCrudEventPublisher<E>` — in tests and modules without subscribers
133/// - Real bus adapters (Kafka, in-memory, etc.) provided by infrastructure
134///
135/// Services always require `Arc<dyn CrudEventPublisher<E>>` — never `Option<...>`.
136#[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    /// Convenience method: publish a Created event.
148    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    /// Convenience method: publish an Updated event.
154    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    /// Convenience method: publish a soft-deleted (or hard-deleted) event.
160    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
166/// No-op implementation — discards all events.  Default for services
167/// that don't need cross-module integration.
168pub 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    /// Convenience constructor that returns the publisher wrapped in an `Arc`,
180    /// ready to inject into a service.
181    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}