Skip to main content

backbone_messaging/
integration_bus.rs

1//! Integration Event Bus - Central event dispatcher for cross-module communication
2//!
3//! The `IntegrationEventBus` provides a type-erased event bus for publishing and
4//! subscribing to integration events across bounded contexts.
5//!
6//! # Architecture
7//!
8//! ```text
9//! ┌─────────────┐     ┌─────────────────────┐     ┌─────────────┐
10//! │   Sapiens   │────▶│ IntegrationEventBus │────▶│   Postman   │
11//! │   Module    │     │  (Type-erased)      │     │   Module    │
12//! └─────────────┘     └─────────────────────┘     └─────────────┘
13//!                              │
14//!                              ▼
15//!                     ┌─────────────┐
16//!                     │  Bucket   │
17//!                     │   Module    │
18//!                     └─────────────┘
19//! ```
20//!
21//! # Example
22//!
23//! ```rust,ignore
24//! use backbone_messaging::{IntegrationEventBus, IntegrationEventHandler};
25//!
26//! // Create the bus
27//! let bus = IntegrationEventBus::new();
28//!
29//! // Register a handler for user events
30//! bus.register_handler(Arc::new(UserEventHandler)).await;
31//!
32//! // Publish an event
33//! bus.publish(UserCreatedIntegrationEvent { ... }).await?;
34//! ```
35
36use std::collections::HashMap;
37use std::sync::Arc;
38
39use async_trait::async_trait;
40use tokio::sync::{broadcast, RwLock};
41use tracing::{debug, error, info, warn};
42
43use crate::integration::{IntegrationEvent, IntegrationEventEnvelope};
44use crate::EventError;
45
46/// Handler for integration events (type-erased)
47///
48/// Implement this trait to create handlers that react to integration events
49/// from other bounded contexts.
50///
51/// Type alias for the handler map to reduce type complexity
52type HandlerMap = HashMap<String, Vec<Arc<dyn IntegrationEventHandler>>>;
53type HandlerMapRef = Arc<RwLock<HandlerMap>>;
54///
55/// # Pattern Matching
56///
57/// Handlers specify which events they're interested in via `event_patterns()`.
58/// Patterns support:
59/// - Exact match: `"sapiens.user.created"`
60/// - Wildcard suffix: `"sapiens.user.*"` matches all user events
61/// - Global wildcard: `"*"` matches all events
62///
63/// # Example
64///
65/// ```rust,ignore
66/// struct EmailNotificationHandler {
67///     email_service: Arc<EmailService>,
68/// }
69///
70/// #[async_trait]
71/// impl IntegrationEventHandler for EmailNotificationHandler {
72///     async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
73///         match envelope.event_type.as_str() {
74///             "sapiens.user.created" => {
75///                 // Deserialize and send welcome email
76///                 let event: UserCreatedEvent = envelope.deserialize()?;
77///                 self.email_service.send_welcome(&event.email).await?;
78///             }
79///             _ => {}
80///         }
81///         Ok(())
82///     }
83///
84///     fn event_patterns(&self) -> Vec<&'static str> {
85///         vec!["sapiens.user.*"]
86///     }
87///
88///     fn name(&self) -> &'static str {
89///         "EmailNotificationHandler"
90///     }
91/// }
92/// ```
93#[async_trait]
94pub trait IntegrationEventHandler: Send + Sync {
95    /// Handle an integration event envelope
96    ///
97    /// Called when an event matching one of the patterns from `event_patterns()`
98    /// is published to the bus.
99    async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError>;
100
101    /// Event patterns this handler is interested in
102    ///
103    /// Patterns support wildcards:
104    /// - `"sapiens.user.created"` - exact match
105    /// - `"sapiens.user.*"` - matches all sapiens.user.* events
106    /// - `"*"` - matches all events
107    fn event_patterns(&self) -> Vec<&'static str>;
108
109    /// Handler name for logging and debugging
110    fn name(&self) -> &'static str;
111
112    /// Whether this handler should be retried on failure
113    fn should_retry(&self) -> bool {
114        true
115    }
116
117    /// Maximum retry attempts
118    fn max_retries(&self) -> u32 {
119        3
120    }
121}
122
123/// Configuration for the integration event bus
124#[derive(Clone, Debug)]
125pub struct IntegrationBusConfig {
126    /// Maximum number of events to buffer in the broadcast channel
127    pub buffer_size: usize,
128    /// Enable event persistence for replay/audit
129    pub persist_events: bool,
130    /// Maximum events to keep in history
131    pub max_history_size: usize,
132    /// Enable dead letter queue for failed handlers
133    pub enable_dead_letter_queue: bool,
134}
135
136impl Default for IntegrationBusConfig {
137    fn default() -> Self {
138        Self {
139            buffer_size: 10000,
140            persist_events: true,
141            max_history_size: 100000,
142            enable_dead_letter_queue: true,
143        }
144    }
145}
146
147impl IntegrationBusConfig {
148    /// Create config with persistence enabled
149    pub fn with_persistence() -> Self {
150        Self {
151            persist_events: true,
152            ..Default::default()
153        }
154    }
155
156    /// Set buffer size
157    pub fn buffer_size(mut self, size: usize) -> Self {
158        self.buffer_size = size;
159        self
160    }
161
162    /// Set max history size
163    pub fn max_history_size(mut self, size: usize) -> Self {
164        self.max_history_size = size;
165        self
166    }
167}
168
169/// Dead letter entry for failed event handling
170#[derive(Clone, Debug)]
171pub struct DeadLetterEntry {
172    /// The envelope that failed
173    pub envelope: IntegrationEventEnvelope,
174    /// Handler that failed
175    pub handler_name: String,
176    /// Error message
177    pub error: String,
178    /// Number of retry attempts
179    pub retry_count: u32,
180    /// When the failure occurred
181    pub failed_at: chrono::DateTime<chrono::Utc>,
182}
183
184/// Central event bus for integration events across all bounded contexts
185///
186/// The `IntegrationEventBus` is the heart of cross-module communication.
187/// It receives events from publishers and dispatches them to registered handlers
188/// based on pattern matching.
189///
190/// # Features
191///
192/// - **Type-erased**: Events are serialized to JSON, allowing loose coupling
193/// - **Pattern matching**: Handlers subscribe to event patterns with wildcards
194/// - **Persistence**: Optional event history for replay/audit
195/// - **Dead letter queue**: Failed events are captured for debugging
196///
197/// # Thread Safety
198///
199/// The bus is fully thread-safe and can be cloned and shared across tasks.
200pub struct IntegrationEventBus {
201    /// Broadcast channel sender
202    sender: broadcast::Sender<IntegrationEventEnvelope>,
203    /// Registered handlers by pattern
204    handlers: HandlerMapRef,
205    /// Event history (if persistence enabled)
206    history: Arc<RwLock<Vec<IntegrationEventEnvelope>>>,
207    /// Dead letter queue for failed handlers
208    dead_letter_queue: Arc<RwLock<Vec<DeadLetterEntry>>>,
209    /// Configuration
210    config: IntegrationBusConfig,
211}
212
213impl IntegrationEventBus {
214    /// Create a new integration event bus with default configuration
215    pub fn new() -> Self {
216        Self::with_config(IntegrationBusConfig::default())
217    }
218
219    /// Create a new integration event bus with custom configuration
220    pub fn with_config(config: IntegrationBusConfig) -> Self {
221        let (sender, _) = broadcast::channel(config.buffer_size);
222        Self {
223            sender,
224            handlers: Arc::new(RwLock::new(HashMap::new())),
225            history: Arc::new(RwLock::new(Vec::new())),
226            dead_letter_queue: Arc::new(RwLock::new(Vec::new())),
227            config,
228        }
229    }
230
231    /// Publish an integration event
232    ///
233    /// The event is serialized to JSON and wrapped in an envelope.
234    /// It's then broadcast to all subscribers and dispatched to matching handlers.
235    ///
236    /// # Errors
237    ///
238    /// Returns `EventError::SerializationError` if the event cannot be serialized.
239    pub async fn publish<E: IntegrationEvent>(&self, event: E) -> Result<(), EventError> {
240        let envelope = IntegrationEventEnvelope::from_event(&event)
241            .map_err(|e| EventError::SerializationError(e.to_string()))?;
242        self.publish_envelope(envelope).await
243    }
244
245    /// Publish a pre-built envelope
246    ///
247    /// Use this when you already have an envelope (e.g., forwarding from another bus).
248    pub async fn publish_envelope(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
249        debug!(
250            event_type = %envelope.event_type,
251            source = %envelope.source_context,
252            aggregate_id = %envelope.aggregate_id,
253            "Publishing integration event"
254        );
255
256        // Store in history if persistence enabled
257        if self.config.persist_events {
258            self.store_event(&envelope).await;
259        }
260
261        // Broadcast to subscribers
262        let _ = self.sender.send(envelope.clone());
263
264        // Dispatch to handlers
265        self.dispatch(envelope).await
266    }
267
268    /// Register an integration event handler
269    ///
270    /// The handler will be called for events matching any of its patterns.
271    pub async fn register_handler(&self, handler: Arc<dyn IntegrationEventHandler>) {
272        let patterns = handler.event_patterns();
273        let handler_name = handler.name();
274        let mut handlers = self.handlers.write().await;
275
276        for pattern in patterns {
277            info!(
278                handler = %handler_name,
279                pattern = %pattern,
280                "Registering integration event handler"
281            );
282            handlers
283                .entry(pattern.to_string())
284                .or_default()
285                .push(Arc::clone(&handler));
286        }
287    }
288
289    /// Subscribe to all integration events (returns a broadcast receiver)
290    ///
291    /// Use this for monitoring, logging, or custom event processing.
292    pub fn subscribe(&self) -> broadcast::Receiver<IntegrationEventEnvelope> {
293        self.sender.subscribe()
294    }
295
296    /// Get event history (if persistence enabled)
297    pub async fn history(&self) -> Vec<IntegrationEventEnvelope> {
298        self.history.read().await.clone()
299    }
300
301    /// Get events by type pattern
302    pub async fn events_by_pattern(&self, pattern: &str) -> Vec<IntegrationEventEnvelope> {
303        self.history
304            .read()
305            .await
306            .iter()
307            .filter(|e| e.matches_pattern(pattern))
308            .cloned()
309            .collect()
310    }
311
312    /// Get events for a specific aggregate
313    pub async fn events_for_aggregate(&self, aggregate_id: &str) -> Vec<IntegrationEventEnvelope> {
314        self.history
315            .read()
316            .await
317            .iter()
318            .filter(|e| e.aggregate_id == aggregate_id)
319            .cloned()
320            .collect()
321    }
322
323    /// Get dead letter queue entries
324    pub async fn dead_letters(&self) -> Vec<DeadLetterEntry> {
325        self.dead_letter_queue.read().await.clone()
326    }
327
328    /// Clear dead letter queue
329    pub async fn clear_dead_letters(&self) {
330        self.dead_letter_queue.write().await.clear();
331    }
332
333    /// Clear event history
334    pub async fn clear_history(&self) {
335        self.history.write().await.clear();
336    }
337
338    /// Get handler count
339    pub async fn handler_count(&self) -> usize {
340        self.handlers
341            .read()
342            .await
343            .values()
344            .map(|v| v.len())
345            .sum()
346    }
347
348    /// Get registered patterns
349    pub async fn registered_patterns(&self) -> Vec<String> {
350        self.handlers.read().await.keys().cloned().collect()
351    }
352
353    // ========================================
354    // Private Methods
355    // ========================================
356
357    async fn store_event(&self, envelope: &IntegrationEventEnvelope) {
358        let mut history = self.history.write().await;
359        history.push(envelope.clone());
360
361        // Trim by max size
362        while history.len() > self.config.max_history_size {
363            history.remove(0);
364        }
365    }
366
367    async fn dispatch(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
368        let handlers = self.handlers.read().await;
369        let mut handlers_to_call = Vec::new();
370
371        // Find handlers matching the event type
372        for (pattern, pattern_handlers) in handlers.iter() {
373            if Self::matches_pattern(pattern, &envelope.event_type) {
374                handlers_to_call.extend(pattern_handlers.iter().cloned());
375            }
376        }
377
378        drop(handlers); // Release lock before calling handlers
379
380        // Deduplicate handlers (same handler might match multiple patterns)
381        let mut seen = std::collections::HashSet::new();
382        handlers_to_call.retain(|h| seen.insert(h.name()));
383
384        debug!(
385            event_type = %envelope.event_type,
386            handler_count = handlers_to_call.len(),
387            "Dispatching integration event"
388        );
389
390        // Call all matching handlers
391        for handler in handlers_to_call {
392            if let Err(e) = self.call_handler_with_retry(&handler, &envelope).await {
393                error!(
394                    handler = %handler.name(),
395                    event_type = %envelope.event_type,
396                    error = ?e,
397                    "Integration event handler failed after retries"
398                );
399
400                // Add to dead letter queue
401                if self.config.enable_dead_letter_queue {
402                    self.add_to_dead_letter(&envelope, handler.name(), &e.to_string()).await;
403                }
404            }
405        }
406
407        Ok(())
408    }
409
410    async fn call_handler_with_retry(
411        &self,
412        handler: &Arc<dyn IntegrationEventHandler>,
413        envelope: &IntegrationEventEnvelope,
414    ) -> Result<(), EventError> {
415        let max_retries = if handler.should_retry() {
416            handler.max_retries()
417        } else {
418            1
419        };
420
421        let mut last_error = None;
422
423        for attempt in 0..max_retries {
424            match handler.handle(envelope.clone()).await {
425                Ok(()) => return Ok(()),
426                Err(e) => {
427                    if attempt < max_retries - 1 {
428                        warn!(
429                            handler = %handler.name(),
430                            attempt = attempt + 1,
431                            max_retries = max_retries,
432                            error = ?e,
433                            "Handler failed, retrying"
434                        );
435                        // Simple backoff
436                        tokio::time::sleep(tokio::time::Duration::from_millis(100 * (attempt as u64 + 1))).await;
437                    }
438                    last_error = Some(e);
439                }
440            }
441        }
442
443        Err(last_error.unwrap_or_else(|| EventError::HandlerError {
444            handler: handler.name().to_string(),
445            message: "Unknown error".to_string(),
446        }))
447    }
448
449    async fn add_to_dead_letter(&self, envelope: &IntegrationEventEnvelope, handler_name: &str, error: &str) {
450        let entry = DeadLetterEntry {
451            envelope: envelope.clone(),
452            handler_name: handler_name.to_string(),
453            error: error.to_string(),
454            retry_count: 3, // Already exhausted retries
455            failed_at: chrono::Utc::now(),
456        };
457
458        self.dead_letter_queue.write().await.push(entry);
459    }
460
461    /// Check if event type matches pattern
462    fn matches_pattern(pattern: &str, event_type: &str) -> bool {
463        if pattern == "*" {
464            return true;
465        }
466        if let Some(prefix) = pattern.strip_suffix(".*") {
467            return event_type.starts_with(prefix);
468        }
469        pattern == event_type
470    }
471}
472
473impl Default for IntegrationEventBus {
474    fn default() -> Self {
475        Self::new()
476    }
477}
478
479impl Clone for IntegrationEventBus {
480    fn clone(&self) -> Self {
481        Self {
482            sender: self.sender.clone(),
483            handlers: Arc::clone(&self.handlers),
484            history: Arc::clone(&self.history),
485            dead_letter_queue: Arc::clone(&self.dead_letter_queue),
486            config: self.config.clone(),
487        }
488    }
489}
490
491/// A handler that logs all integration events
492pub struct IntegrationLoggingHandler {
493    patterns: Vec<&'static str>,
494}
495
496impl IntegrationLoggingHandler {
497    /// Create a handler that logs specific patterns
498    pub fn new(patterns: Vec<&'static str>) -> Self {
499        Self { patterns }
500    }
501
502    /// Create a handler that logs all events
503    pub fn all() -> Self {
504        Self { patterns: vec!["*"] }
505    }
506}
507
508#[async_trait]
509impl IntegrationEventHandler for IntegrationLoggingHandler {
510    async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
511        info!(
512            event_type = %envelope.event_type,
513            source = %envelope.source_context,
514            aggregate_id = %envelope.aggregate_id,
515            correlation_id = ?envelope.correlation_id,
516            "Integration event received"
517        );
518        Ok(())
519    }
520
521    fn event_patterns(&self) -> Vec<&'static str> {
522        self.patterns.clone()
523    }
524
525    fn name(&self) -> &'static str {
526        "IntegrationLoggingHandler"
527    }
528
529    fn should_retry(&self) -> bool {
530        false // Logging shouldn't retry
531    }
532}
533
534#[cfg(test)]
535mod tests {
536    use super::*;
537    use chrono::Utc;
538    use serde::{Deserialize, Serialize};
539    use std::sync::atomic::{AtomicUsize, Ordering};
540
541    #[derive(Clone, Debug, Serialize, Deserialize)]
542    struct TestEvent {
543        id: String,
544        data: String,
545        occurred_at: chrono::DateTime<Utc>,
546    }
547
548    impl IntegrationEvent for TestEvent {
549        fn event_type(&self) -> &'static str {
550            "test.entity.created"
551        }
552
553        fn source_context(&self) -> &'static str {
554            "test"
555        }
556
557        fn aggregate_id(&self) -> &str {
558            &self.id
559        }
560
561        fn occurred_at(&self) -> chrono::DateTime<Utc> {
562            self.occurred_at
563        }
564    }
565
566    struct CountingHandler {
567        count: Arc<AtomicUsize>,
568        patterns: Vec<&'static str>,
569    }
570
571    impl CountingHandler {
572        fn new(patterns: Vec<&'static str>) -> Self {
573            Self {
574                count: Arc::new(AtomicUsize::new(0)),
575                patterns,
576            }
577        }
578
579        fn count(&self) -> usize {
580            self.count.load(Ordering::SeqCst)
581        }
582    }
583
584    #[async_trait]
585    impl IntegrationEventHandler for CountingHandler {
586        async fn handle(&self, _envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
587            self.count.fetch_add(1, Ordering::SeqCst);
588            Ok(())
589        }
590
591        fn event_patterns(&self) -> Vec<&'static str> {
592            self.patterns.clone()
593        }
594
595        fn name(&self) -> &'static str {
596            "CountingHandler"
597        }
598    }
599
600    #[tokio::test]
601    async fn test_bus_publish_and_handle() {
602        let bus = IntegrationEventBus::new();
603        let handler = Arc::new(CountingHandler::new(vec!["test.entity.created"]));
604
605        bus.register_handler(handler.clone()).await;
606
607        let event = TestEvent {
608            id: "test-123".to_string(),
609            data: "Hello".to_string(),
610            occurred_at: Utc::now(),
611        };
612
613        bus.publish(event).await.unwrap();
614
615        // Give handler time to process
616        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
617
618        assert_eq!(handler.count(), 1);
619    }
620
621    #[tokio::test]
622    async fn test_bus_wildcard_pattern() {
623        let bus = IntegrationEventBus::new();
624        let handler = Arc::new(CountingHandler::new(vec!["test.*"]));
625
626        bus.register_handler(handler.clone()).await;
627
628        let event = TestEvent {
629            id: "test-123".to_string(),
630            data: "Hello".to_string(),
631            occurred_at: Utc::now(),
632        };
633
634        bus.publish(event).await.unwrap();
635
636        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
637
638        assert_eq!(handler.count(), 1);
639    }
640
641    #[tokio::test]
642    async fn test_bus_global_wildcard() {
643        let bus = IntegrationEventBus::new();
644        let handler = Arc::new(CountingHandler::new(vec!["*"]));
645
646        bus.register_handler(handler.clone()).await;
647
648        let event = TestEvent {
649            id: "test-123".to_string(),
650            data: "Hello".to_string(),
651            occurred_at: Utc::now(),
652        };
653
654        bus.publish(event).await.unwrap();
655
656        tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
657
658        assert_eq!(handler.count(), 1);
659    }
660
661    #[tokio::test]
662    async fn test_bus_history() {
663        let bus = IntegrationEventBus::with_config(IntegrationBusConfig::with_persistence());
664
665        let event = TestEvent {
666            id: "test-123".to_string(),
667            data: "Hello".to_string(),
668            occurred_at: Utc::now(),
669        };
670
671        bus.publish(event).await.unwrap();
672
673        let history = bus.history().await;
674        assert_eq!(history.len(), 1);
675        assert_eq!(history[0].event_type, "test.entity.created");
676    }
677
678    #[tokio::test]
679    async fn test_bus_subscribe() {
680        let bus = IntegrationEventBus::new();
681        let mut rx = bus.subscribe();
682
683        let event = TestEvent {
684            id: "test-123".to_string(),
685            data: "Hello".to_string(),
686            occurred_at: Utc::now(),
687        };
688
689        bus.publish(event).await.unwrap();
690
691        let envelope = rx.recv().await.unwrap();
692        assert_eq!(envelope.event_type, "test.entity.created");
693    }
694
695    #[test]
696    fn test_pattern_matching() {
697        // Exact match
698        assert!(IntegrationEventBus::matches_pattern("test.user.created", "test.user.created"));
699        assert!(!IntegrationEventBus::matches_pattern("test.user.created", "test.user.deleted"));
700
701        // Wildcard suffix
702        assert!(IntegrationEventBus::matches_pattern("test.user.*", "test.user.created"));
703        assert!(IntegrationEventBus::matches_pattern("test.user.*", "test.user.deleted"));
704        assert!(!IntegrationEventBus::matches_pattern("test.user.*", "test.role.created"));
705
706        // Multi-level wildcard
707        assert!(IntegrationEventBus::matches_pattern("test.*", "test.user.created"));
708        assert!(IntegrationEventBus::matches_pattern("test.*", "test.role.deleted"));
709
710        // Global wildcard
711        assert!(IntegrationEventBus::matches_pattern("*", "anything.here"));
712    }
713}