use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use tokio::sync::{broadcast, RwLock};
use tracing::{debug, error, info, warn};
use crate::integration::{IntegrationEvent, IntegrationEventEnvelope};
use crate::EventError;
type HandlerMap = HashMap<String, Vec<Arc<dyn IntegrationEventHandler>>>;
type HandlerMapRef = Arc<RwLock<HandlerMap>>;
#[async_trait]
pub trait IntegrationEventHandler: Send + Sync {
async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError>;
fn event_patterns(&self) -> Vec<&'static str>;
fn name(&self) -> &'static str;
fn should_retry(&self) -> bool {
true
}
fn max_retries(&self) -> u32 {
3
}
}
#[derive(Clone, Debug)]
pub struct IntegrationBusConfig {
pub buffer_size: usize,
pub persist_events: bool,
pub max_history_size: usize,
pub enable_dead_letter_queue: bool,
}
impl Default for IntegrationBusConfig {
fn default() -> Self {
Self {
buffer_size: 10000,
persist_events: true,
max_history_size: 100000,
enable_dead_letter_queue: true,
}
}
}
impl IntegrationBusConfig {
pub fn with_persistence() -> Self {
Self {
persist_events: true,
..Default::default()
}
}
pub fn buffer_size(mut self, size: usize) -> Self {
self.buffer_size = size;
self
}
pub fn max_history_size(mut self, size: usize) -> Self {
self.max_history_size = size;
self
}
}
#[derive(Clone, Debug)]
pub struct DeadLetterEntry {
pub envelope: IntegrationEventEnvelope,
pub handler_name: String,
pub error: String,
pub retry_count: u32,
pub failed_at: chrono::DateTime<chrono::Utc>,
}
pub struct IntegrationEventBus {
sender: broadcast::Sender<IntegrationEventEnvelope>,
handlers: HandlerMapRef,
history: Arc<RwLock<Vec<IntegrationEventEnvelope>>>,
dead_letter_queue: Arc<RwLock<Vec<DeadLetterEntry>>>,
config: IntegrationBusConfig,
}
impl IntegrationEventBus {
pub fn new() -> Self {
Self::with_config(IntegrationBusConfig::default())
}
pub fn with_config(config: IntegrationBusConfig) -> Self {
let (sender, _) = broadcast::channel(config.buffer_size);
Self {
sender,
handlers: Arc::new(RwLock::new(HashMap::new())),
history: Arc::new(RwLock::new(Vec::new())),
dead_letter_queue: Arc::new(RwLock::new(Vec::new())),
config,
}
}
pub async fn publish<E: IntegrationEvent>(&self, event: E) -> Result<(), EventError> {
let envelope = IntegrationEventEnvelope::from_event(&event)
.map_err(|e| EventError::SerializationError(e.to_string()))?;
self.publish_envelope(envelope).await
}
pub async fn publish_envelope(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
debug!(
event_type = %envelope.event_type,
source = %envelope.source_context,
aggregate_id = %envelope.aggregate_id,
"Publishing integration event"
);
if self.config.persist_events {
self.store_event(&envelope).await;
}
let _ = self.sender.send(envelope.clone());
self.dispatch(envelope).await
}
pub async fn register_handler(&self, handler: Arc<dyn IntegrationEventHandler>) {
let patterns = handler.event_patterns();
let handler_name = handler.name();
let mut handlers = self.handlers.write().await;
for pattern in patterns {
info!(
handler = %handler_name,
pattern = %pattern,
"Registering integration event handler"
);
handlers
.entry(pattern.to_string())
.or_default()
.push(Arc::clone(&handler));
}
}
pub fn subscribe(&self) -> broadcast::Receiver<IntegrationEventEnvelope> {
self.sender.subscribe()
}
pub async fn history(&self) -> Vec<IntegrationEventEnvelope> {
self.history.read().await.clone()
}
pub async fn events_by_pattern(&self, pattern: &str) -> Vec<IntegrationEventEnvelope> {
self.history
.read()
.await
.iter()
.filter(|e| e.matches_pattern(pattern))
.cloned()
.collect()
}
pub async fn events_for_aggregate(&self, aggregate_id: &str) -> Vec<IntegrationEventEnvelope> {
self.history
.read()
.await
.iter()
.filter(|e| e.aggregate_id == aggregate_id)
.cloned()
.collect()
}
pub async fn dead_letters(&self) -> Vec<DeadLetterEntry> {
self.dead_letter_queue.read().await.clone()
}
pub async fn clear_dead_letters(&self) {
self.dead_letter_queue.write().await.clear();
}
pub async fn clear_history(&self) {
self.history.write().await.clear();
}
pub async fn handler_count(&self) -> usize {
self.handlers
.read()
.await
.values()
.map(|v| v.len())
.sum()
}
pub async fn registered_patterns(&self) -> Vec<String> {
self.handlers.read().await.keys().cloned().collect()
}
async fn store_event(&self, envelope: &IntegrationEventEnvelope) {
let mut history = self.history.write().await;
history.push(envelope.clone());
while history.len() > self.config.max_history_size {
history.remove(0);
}
}
async fn dispatch(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
let handlers = self.handlers.read().await;
let mut handlers_to_call = Vec::new();
for (pattern, pattern_handlers) in handlers.iter() {
if Self::matches_pattern(pattern, &envelope.event_type) {
handlers_to_call.extend(pattern_handlers.iter().cloned());
}
}
drop(handlers);
let mut seen = std::collections::HashSet::new();
handlers_to_call.retain(|h| seen.insert(h.name()));
debug!(
event_type = %envelope.event_type,
handler_count = handlers_to_call.len(),
"Dispatching integration event"
);
for handler in handlers_to_call {
if let Err(e) = self.call_handler_with_retry(&handler, &envelope).await {
error!(
handler = %handler.name(),
event_type = %envelope.event_type,
error = ?e,
"Integration event handler failed after retries"
);
if self.config.enable_dead_letter_queue {
self.add_to_dead_letter(&envelope, handler.name(), &e.to_string()).await;
}
}
}
Ok(())
}
async fn call_handler_with_retry(
&self,
handler: &Arc<dyn IntegrationEventHandler>,
envelope: &IntegrationEventEnvelope,
) -> Result<(), EventError> {
let max_retries = if handler.should_retry() {
handler.max_retries()
} else {
1
};
let mut last_error = None;
for attempt in 0..max_retries {
match handler.handle(envelope.clone()).await {
Ok(()) => return Ok(()),
Err(e) => {
if attempt < max_retries - 1 {
warn!(
handler = %handler.name(),
attempt = attempt + 1,
max_retries = max_retries,
error = ?e,
"Handler failed, retrying"
);
tokio::time::sleep(tokio::time::Duration::from_millis(100 * (attempt as u64 + 1))).await;
}
last_error = Some(e);
}
}
}
Err(last_error.unwrap_or_else(|| EventError::HandlerError {
handler: handler.name().to_string(),
message: "Unknown error".to_string(),
}))
}
async fn add_to_dead_letter(&self, envelope: &IntegrationEventEnvelope, handler_name: &str, error: &str) {
let entry = DeadLetterEntry {
envelope: envelope.clone(),
handler_name: handler_name.to_string(),
error: error.to_string(),
retry_count: 3, failed_at: chrono::Utc::now(),
};
self.dead_letter_queue.write().await.push(entry);
}
fn matches_pattern(pattern: &str, event_type: &str) -> bool {
if pattern == "*" {
return true;
}
if let Some(prefix) = pattern.strip_suffix(".*") {
return event_type.starts_with(prefix);
}
pattern == event_type
}
}
impl Default for IntegrationEventBus {
fn default() -> Self {
Self::new()
}
}
impl Clone for IntegrationEventBus {
fn clone(&self) -> Self {
Self {
sender: self.sender.clone(),
handlers: Arc::clone(&self.handlers),
history: Arc::clone(&self.history),
dead_letter_queue: Arc::clone(&self.dead_letter_queue),
config: self.config.clone(),
}
}
}
pub struct IntegrationLoggingHandler {
patterns: Vec<&'static str>,
}
impl IntegrationLoggingHandler {
pub fn new(patterns: Vec<&'static str>) -> Self {
Self { patterns }
}
pub fn all() -> Self {
Self { patterns: vec!["*"] }
}
}
#[async_trait]
impl IntegrationEventHandler for IntegrationLoggingHandler {
async fn handle(&self, envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
info!(
event_type = %envelope.event_type,
source = %envelope.source_context,
aggregate_id = %envelope.aggregate_id,
correlation_id = ?envelope.correlation_id,
"Integration event received"
);
Ok(())
}
fn event_patterns(&self) -> Vec<&'static str> {
self.patterns.clone()
}
fn name(&self) -> &'static str {
"IntegrationLoggingHandler"
}
fn should_retry(&self) -> bool {
false }
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Utc;
use serde::{Deserialize, Serialize};
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone, Debug, Serialize, Deserialize)]
struct TestEvent {
id: String,
data: String,
occurred_at: chrono::DateTime<Utc>,
}
impl IntegrationEvent for TestEvent {
fn event_type(&self) -> &'static str {
"test.entity.created"
}
fn source_context(&self) -> &'static str {
"test"
}
fn aggregate_id(&self) -> &str {
&self.id
}
fn occurred_at(&self) -> chrono::DateTime<Utc> {
self.occurred_at
}
}
struct CountingHandler {
count: Arc<AtomicUsize>,
patterns: Vec<&'static str>,
}
impl CountingHandler {
fn new(patterns: Vec<&'static str>) -> Self {
Self {
count: Arc::new(AtomicUsize::new(0)),
patterns,
}
}
fn count(&self) -> usize {
self.count.load(Ordering::SeqCst)
}
}
#[async_trait]
impl IntegrationEventHandler for CountingHandler {
async fn handle(&self, _envelope: IntegrationEventEnvelope) -> Result<(), EventError> {
self.count.fetch_add(1, Ordering::SeqCst);
Ok(())
}
fn event_patterns(&self) -> Vec<&'static str> {
self.patterns.clone()
}
fn name(&self) -> &'static str {
"CountingHandler"
}
}
#[tokio::test]
async fn test_bus_publish_and_handle() {
let bus = IntegrationEventBus::new();
let handler = Arc::new(CountingHandler::new(vec!["test.entity.created"]));
bus.register_handler(handler.clone()).await;
let event = TestEvent {
id: "test-123".to_string(),
data: "Hello".to_string(),
occurred_at: Utc::now(),
};
bus.publish(event).await.unwrap();
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
assert_eq!(handler.count(), 1);
}
#[tokio::test]
async fn test_bus_wildcard_pattern() {
let bus = IntegrationEventBus::new();
let handler = Arc::new(CountingHandler::new(vec!["test.*"]));
bus.register_handler(handler.clone()).await;
let event = TestEvent {
id: "test-123".to_string(),
data: "Hello".to_string(),
occurred_at: Utc::now(),
};
bus.publish(event).await.unwrap();
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
assert_eq!(handler.count(), 1);
}
#[tokio::test]
async fn test_bus_global_wildcard() {
let bus = IntegrationEventBus::new();
let handler = Arc::new(CountingHandler::new(vec!["*"]));
bus.register_handler(handler.clone()).await;
let event = TestEvent {
id: "test-123".to_string(),
data: "Hello".to_string(),
occurred_at: Utc::now(),
};
bus.publish(event).await.unwrap();
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
assert_eq!(handler.count(), 1);
}
#[tokio::test]
async fn test_bus_history() {
let bus = IntegrationEventBus::with_config(IntegrationBusConfig::with_persistence());
let event = TestEvent {
id: "test-123".to_string(),
data: "Hello".to_string(),
occurred_at: Utc::now(),
};
bus.publish(event).await.unwrap();
let history = bus.history().await;
assert_eq!(history.len(), 1);
assert_eq!(history[0].event_type, "test.entity.created");
}
#[tokio::test]
async fn test_bus_subscribe() {
let bus = IntegrationEventBus::new();
let mut rx = bus.subscribe();
let event = TestEvent {
id: "test-123".to_string(),
data: "Hello".to_string(),
occurred_at: Utc::now(),
};
bus.publish(event).await.unwrap();
let envelope = rx.recv().await.unwrap();
assert_eq!(envelope.event_type, "test.entity.created");
}
#[test]
fn test_pattern_matching() {
assert!(IntegrationEventBus::matches_pattern("test.user.created", "test.user.created"));
assert!(!IntegrationEventBus::matches_pattern("test.user.created", "test.user.deleted"));
assert!(IntegrationEventBus::matches_pattern("test.user.*", "test.user.created"));
assert!(IntegrationEventBus::matches_pattern("test.user.*", "test.user.deleted"));
assert!(!IntegrationEventBus::matches_pattern("test.user.*", "test.role.created"));
assert!(IntegrationEventBus::matches_pattern("test.*", "test.user.created"));
assert!(IntegrationEventBus::matches_pattern("test.*", "test.role.deleted"));
assert!(IntegrationEventBus::matches_pattern("*", "anything.here"));
}
}