1use 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
46type HandlerMap = HashMap<String, Vec<Arc<dyn IntegrationEventHandler>>>;
53type HandlerMapRef = Arc<RwLock<HandlerMap>>;
54#[async_trait]
94pub trait IntegrationEventHandler: Send + Sync {
95 async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError>;
100
101 fn event_patterns(&self) -> Vec<&'static str>;
108
109 fn name(&self) -> &'static str;
111
112 fn should_retry(&self) -> bool {
114 true
115 }
116
117 fn max_retries(&self) -> u32 {
119 3
120 }
121}
122
123#[derive(Clone, Debug)]
125pub struct IntegrationBusConfig {
126 pub buffer_size: usize,
128 pub persist_events: bool,
130 pub max_history_size: usize,
132 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 pub fn with_persistence() -> Self {
150 Self {
151 persist_events: true,
152 ..Default::default()
153 }
154 }
155
156 pub fn buffer_size(mut self, size: usize) -> Self {
158 self.buffer_size = size;
159 self
160 }
161
162 pub fn max_history_size(mut self, size: usize) -> Self {
164 self.max_history_size = size;
165 self
166 }
167}
168
169#[derive(Clone, Debug)]
171pub struct DeadLetterEntry {
172 pub envelope: IntegrationEventEnvelope,
174 pub handler_name: String,
176 pub error: String,
178 pub retry_count: u32,
180 pub failed_at: chrono::DateTime<chrono::Utc>,
182}
183
184pub struct IntegrationEventBus {
201 sender: broadcast::Sender<IntegrationEventEnvelope>,
203 handlers: HandlerMapRef,
205 history: Arc<RwLock<Vec<IntegrationEventEnvelope>>>,
207 dead_letter_queue: Arc<RwLock<Vec<DeadLetterEntry>>>,
209 config: IntegrationBusConfig,
211}
212
213impl IntegrationEventBus {
214 pub fn new() -> Self {
216 Self::with_config(IntegrationBusConfig::default())
217 }
218
219 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 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 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 if self.config.persist_events {
258 self.store_event(&envelope).await;
259 }
260
261 let _ = self.sender.send(envelope.clone());
263
264 self.dispatch(envelope).await
266 }
267
268 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 pub fn subscribe(&self) -> broadcast::Receiver<IntegrationEventEnvelope> {
293 self.sender.subscribe()
294 }
295
296 pub async fn history(&self) -> Vec<IntegrationEventEnvelope> {
298 self.history.read().await.clone()
299 }
300
301 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 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 pub async fn dead_letters(&self) -> Vec<DeadLetterEntry> {
325 self.dead_letter_queue.read().await.clone()
326 }
327
328 pub async fn clear_dead_letters(&self) {
330 self.dead_letter_queue.write().await.clear();
331 }
332
333 pub async fn clear_history(&self) {
335 self.history.write().await.clear();
336 }
337
338 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 pub async fn registered_patterns(&self) -> Vec<String> {
350 self.handlers.read().await.keys().cloned().collect()
351 }
352
353 async fn store_event(&self, envelope: &IntegrationEventEnvelope) {
358 let mut history = self.history.write().await;
359 history.push(envelope.clone());
360
361 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 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); 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 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 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 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, failed_at: chrono::Utc::now(),
456 };
457
458 self.dead_letter_queue.write().await.push(entry);
459 }
460
461 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
491pub struct IntegrationLoggingHandler {
493 patterns: Vec<&'static str>,
494}
495
496impl IntegrationLoggingHandler {
497 pub fn new(patterns: Vec<&'static str>) -> Self {
499 Self { patterns }
500 }
501
502 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 }
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 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 assert!(IntegrationEventBus::matches_pattern("test.user.created", "test.user.created"));
699 assert!(!IntegrationEventBus::matches_pattern("test.user.created", "test.user.deleted"));
700
701 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 assert!(IntegrationEventBus::matches_pattern("test.*", "test.user.created"));
708 assert!(IntegrationEventBus::matches_pattern("test.*", "test.role.deleted"));
709
710 assert!(IntegrationEventBus::matches_pattern("*", "anything.here"));
712 }
713}