backbone_messaging/integration.rs
1//! Integration Events - Cross-Bounded Context Communication
2//!
3//! Integration events differ from domain events:
4//! - **Domain events**: Internal to a bounded context, may contain domain-specific types
5//! - **Integration events**: Published for other bounded contexts to consume, contain only primitive/common types
6//!
7//! # Key Differences
8//!
9//! | Aspect | Domain Event | Integration Event |
10//! |--------|--------------|-------------------|
11//! | Scope | Single bounded context | Cross-context |
12//! | Types | Can use domain types | Primitives only |
13//! | Coupling | Internal | Loose coupling |
14//! | Serialization | Optional | Required (JSON) |
15//!
16//! # Example
17//!
18//! ```rust,ignore
19//! use backbone_messaging::IntegrationEvent;
20//! use chrono::{DateTime, Utc};
21//! use serde::{Deserialize, Serialize};
22//!
23//! #[derive(Clone, Debug, Serialize, Deserialize)]
24//! pub struct UserCreatedIntegrationEvent {
25//! pub user_id: String,
26//! pub email: String,
27//! pub occurred_at: DateTime<Utc>,
28//! }
29//!
30//! impl IntegrationEvent for UserCreatedIntegrationEvent {
31//! fn event_type(&self) -> &'static str { "sapiens.user.created" }
32//! fn source_context(&self) -> &'static str { "sapiens" }
33//! fn aggregate_id(&self) -> &str { &self.user_id }
34//! fn occurred_at(&self) -> DateTime<Utc> { self.occurred_at }
35//! }
36//! ```
37
38use chrono::{DateTime, Utc};
39use serde::{Deserialize, Serialize};
40use uuid::Uuid;
41
42/// Trait for integration events that cross bounded context boundaries
43///
44/// Integration events are serializable and contain only primitive/common types
45/// to avoid coupling between modules. They follow a naming convention:
46/// `{source_context}.{aggregate}.{action}` (e.g., "sapiens.user.created")
47///
48/// # Implementation Requirements
49///
50/// - Must be `Clone + Send + Sync + 'static` for async handling
51/// - Must implement `Serialize + Deserialize` for JSON transport
52/// - Should only use primitive types (String, numbers, booleans)
53/// - Event type should follow dot-notation naming
54///
55/// # Example
56///
57/// ```rust,ignore
58/// impl IntegrationEvent for UserCreatedEvent {
59/// fn event_type(&self) -> &'static str { "sapiens.user.created" }
60/// fn source_context(&self) -> &'static str { "sapiens" }
61/// fn aggregate_id(&self) -> &str { &self.user_id }
62/// fn occurred_at(&self) -> DateTime<Utc> { self.occurred_at }
63/// }
64/// ```
65pub trait IntegrationEvent: Clone + Send + Sync + Serialize + for<'de> Deserialize<'de> + 'static {
66 /// Event type identifier using dot notation
67 ///
68 /// Format: `{source_context}.{aggregate}.{action}`
69 /// Examples: "sapiens.user.created", "postman.email.sent"
70 fn event_type(&self) -> &'static str;
71
72 /// Source bounded context that publishes this event
73 ///
74 /// Examples: "sapiens", "postman", "bucket"
75 fn source_context(&self) -> &'static str;
76
77 /// Aggregate ID in the source context
78 ///
79 /// This identifies which aggregate instance the event belongs to
80 fn aggregate_id(&self) -> &str;
81
82 /// When the event occurred in the domain
83 fn occurred_at(&self) -> DateTime<Utc>;
84
85 /// Event schema version for evolution
86 ///
87 /// Increment when making breaking changes to the event structure.
88 /// Consumers can use this to handle multiple versions.
89 fn version(&self) -> u32 {
90 1
91 }
92
93 /// Correlation ID for distributed tracing
94 ///
95 /// Used to trace a request across multiple bounded contexts.
96 /// Returns `None` by default; override if your event carries correlation info.
97 fn correlation_id(&self) -> Option<&str> {
98 None
99 }
100}
101
102/// Type-erased integration event envelope for cross-module transport
103///
104/// The envelope wraps any integration event with metadata and serializes
105/// the event payload to JSON for transport between bounded contexts.
106///
107/// # Fields
108///
109/// - `id`: Unique envelope ID (UUID)
110/// - `event_type`: Dot-notation event type
111/// - `source_context`: Origin bounded context
112/// - `payload`: JSON-serialized event data
113///
114/// # Example
115///
116/// ```rust,ignore
117/// let event = UserCreatedIntegrationEvent { ... };
118/// let envelope = IntegrationEventEnvelope::from_event(&event)?;
119///
120/// // Later, deserialize back to the typed event
121/// let restored: UserCreatedIntegrationEvent = envelope.deserialize()?;
122/// ```
123#[derive(Clone, Debug, Serialize, Deserialize)]
124pub struct IntegrationEventEnvelope {
125 /// Unique envelope ID
126 pub id: String,
127 /// Event type (e.g., "sapiens.user.created")
128 pub event_type: String,
129 /// Source bounded context
130 pub source_context: String,
131 /// Aggregate ID in source context
132 pub aggregate_id: String,
133 /// When the event occurred
134 pub occurred_at: DateTime<Utc>,
135 /// When the envelope was created/published
136 pub published_at: DateTime<Utc>,
137 /// Event schema version
138 pub version: u32,
139 /// Correlation ID for distributed tracing
140 pub correlation_id: Option<String>,
141 /// Causation ID (parent event/envelope that caused this one)
142 pub causation_id: Option<String>,
143 /// JSON-serialized event payload
144 pub payload: serde_json::Value,
145}
146
147impl IntegrationEventEnvelope {
148 /// Create an envelope from an integration event
149 ///
150 /// Serializes the event to JSON and wraps it with metadata.
151 ///
152 /// # Errors
153 ///
154 /// Returns `serde_json::Error` if the event cannot be serialized.
155 pub fn from_event<E: IntegrationEvent>(event: &E) -> Result<Self, serde_json::Error> {
156 Ok(Self {
157 id: Uuid::new_v4().to_string(),
158 event_type: event.event_type().to_string(),
159 source_context: event.source_context().to_string(),
160 aggregate_id: event.aggregate_id().to_string(),
161 occurred_at: event.occurred_at(),
162 published_at: Utc::now(),
163 version: event.version(),
164 correlation_id: event.correlation_id().map(String::from),
165 causation_id: None,
166 payload: serde_json::to_value(event)?,
167 })
168 }
169
170 /// Deserialize the payload back to a typed event
171 ///
172 /// # Errors
173 ///
174 /// Returns `serde_json::Error` if the payload doesn't match the expected type.
175 pub fn deserialize<E: IntegrationEvent>(&self) -> Result<E, serde_json::Error> {
176 serde_json::from_value(self.payload.clone())
177 }
178
179 /// Set the causation ID (for event chaining)
180 pub fn with_causation_id(mut self, causation_id: impl Into<String>) -> Self {
181 self.causation_id = Some(causation_id.into());
182 self
183 }
184
185 /// Set the correlation ID (for distributed tracing)
186 pub fn with_correlation_id(mut self, correlation_id: impl Into<String>) -> Self {
187 self.correlation_id = Some(correlation_id.into());
188 self
189 }
190
191 /// Check if this envelope matches a pattern
192 ///
193 /// Patterns support:
194 /// - Exact match: "sapiens.user.created"
195 /// - Wildcard suffix: "sapiens.user.*" matches "sapiens.user.created", "sapiens.user.deleted"
196 /// - Global wildcard: "*" matches everything
197 pub fn matches_pattern(&self, pattern: &str) -> bool {
198 if pattern == "*" {
199 return true;
200 }
201 if let Some(prefix) = pattern.strip_suffix(".*") {
202 return self.event_type.starts_with(prefix);
203 }
204 pattern == self.event_type
205 }
206}
207
208#[cfg(test)]
209mod tests {
210 use super::*;
211
212 #[derive(Clone, Debug, Serialize, Deserialize)]
213 struct TestIntegrationEvent {
214 user_id: String,
215 email: String,
216 occurred_at: DateTime<Utc>,
217 }
218
219 impl IntegrationEvent for TestIntegrationEvent {
220 fn event_type(&self) -> &'static str {
221 "test.user.created"
222 }
223
224 fn source_context(&self) -> &'static str {
225 "test"
226 }
227
228 fn aggregate_id(&self) -> &str {
229 &self.user_id
230 }
231
232 fn occurred_at(&self) -> DateTime<Utc> {
233 self.occurred_at
234 }
235 }
236
237 #[test]
238 fn test_envelope_from_event() {
239 let event = TestIntegrationEvent {
240 user_id: "user-123".to_string(),
241 email: "test@example.com".to_string(),
242 occurred_at: Utc::now(),
243 };
244
245 let envelope = IntegrationEventEnvelope::from_event(&event).unwrap();
246
247 assert_eq!(envelope.event_type, "test.user.created");
248 assert_eq!(envelope.source_context, "test");
249 assert_eq!(envelope.aggregate_id, "user-123");
250 assert_eq!(envelope.version, 1);
251 assert!(!envelope.id.is_empty());
252 }
253
254 #[test]
255 fn test_envelope_deserialize() {
256 let event = TestIntegrationEvent {
257 user_id: "user-123".to_string(),
258 email: "test@example.com".to_string(),
259 occurred_at: Utc::now(),
260 };
261
262 let envelope = IntegrationEventEnvelope::from_event(&event).unwrap();
263 let restored: TestIntegrationEvent = envelope.deserialize().unwrap();
264
265 assert_eq!(restored.user_id, "user-123");
266 assert_eq!(restored.email, "test@example.com");
267 }
268
269 #[test]
270 fn test_pattern_matching() {
271 let event = TestIntegrationEvent {
272 user_id: "user-123".to_string(),
273 email: "test@example.com".to_string(),
274 occurred_at: Utc::now(),
275 };
276
277 let envelope = IntegrationEventEnvelope::from_event(&event).unwrap();
278
279 // Exact match
280 assert!(envelope.matches_pattern("test.user.created"));
281 assert!(!envelope.matches_pattern("test.user.deleted"));
282
283 // Wildcard suffix
284 assert!(envelope.matches_pattern("test.user.*"));
285 assert!(envelope.matches_pattern("test.*"));
286 assert!(!envelope.matches_pattern("other.*"));
287
288 // Global wildcard
289 assert!(envelope.matches_pattern("*"));
290 }
291
292 #[test]
293 fn test_envelope_with_causation() {
294 let event = TestIntegrationEvent {
295 user_id: "user-123".to_string(),
296 email: "test@example.com".to_string(),
297 occurred_at: Utc::now(),
298 };
299
300 let envelope = IntegrationEventEnvelope::from_event(&event)
301 .unwrap()
302 .with_causation_id("parent-event-id")
303 .with_correlation_id("trace-123");
304
305 assert_eq!(envelope.causation_id, Some("parent-event-id".to_string()));
306 assert_eq!(envelope.correlation_id, Some("trace-123".to_string()));
307 }
308}