Skip to main content

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}