use crate::event::Event;
use async_trait::async_trait;
use std::sync::Arc;
use thiserror::Error;
use tokio::sync::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HandlerLane {
Sync,
Async,
}
#[derive(Debug, Error)]
pub enum EventBusError {
#[error("Event handler failed for event '{event_type}' (event_id: {event_id}): {message}")]
HandlerFailed {
event_type: String,
event_id: String,
message: String,
},
#[error("No handlers registered for event type '{event_type}'")]
NoHandlers { event_type: String },
#[error("Failed to subscribe handler: {message}")]
SubscriptionFailed { message: String },
#[error("Event bus error: {message}")]
Other { message: String },
}
impl EventBusError {
pub fn handler_failed(
event_type: impl Into<String>,
event_id: impl Into<String>,
message: impl Into<String>,
) -> Self {
EventBusError::HandlerFailed {
event_type: event_type.into(),
event_id: event_id.into(),
message: message.into(),
}
}
pub fn no_handlers(event_type: impl Into<String>) -> Self {
EventBusError::NoHandlers {
event_type: event_type.into(),
}
}
pub fn subscription_failed(message: impl Into<String>) -> Self {
EventBusError::SubscriptionFailed {
message: message.into(),
}
}
pub fn other(message: impl Into<String>) -> Self {
EventBusError::Other {
message: message.into(),
}
}
}
pub type EventBusResult<T> = Result<T, EventBusError>;
#[async_trait]
pub trait EventHandler: Send + Sync {
fn handles(&self) -> Vec<String>;
async fn handle(&self, event: &Event) -> Result<(), Box<dyn std::error::Error + Send + Sync>>;
fn lane(&self) -> HandlerLane {
HandlerLane::Sync
}
}
#[async_trait]
pub trait EventBus: Send + Sync {
async fn publish(&self, events: Vec<Event>) -> EventBusResult<()>;
async fn subscribe(&mut self, handler: Box<dyn EventHandler>) -> EventBusResult<()>;
}
#[derive(Clone)]
pub struct InProcessEventBus {
handlers: Arc<Mutex<Vec<Box<dyn EventHandler>>>>,
}
impl InProcessEventBus {
pub fn new() -> Self {
Self {
handlers: Arc::new(Mutex::new(Vec::new())),
}
}
pub async fn handler_count(&self) -> usize {
self.handlers.lock().await.len()
}
}
impl Default for InProcessEventBus {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
impl EventBus for InProcessEventBus {
async fn publish(&self, events: Vec<Event>) -> EventBusResult<()> {
let handlers = self.handlers.lock().await;
for event in &events {
for handler in handlers.iter() {
let handled_types = handler.handles();
if handled_types.contains(&event.event_type) {
handler.handle(event).await.map_err(|e| {
EventBusError::handler_failed(
&event.event_type,
event.event_id.to_string(),
e.to_string(),
)
})?;
}
}
}
Ok(())
}
async fn subscribe(&mut self, handler: Box<dyn EventHandler>) -> EventBusResult<()> {
let mut handlers = self.handlers.lock().await;
handlers.push(handler);
Ok(())
}
}
#[derive(Clone)]
pub struct TwoLaneEventBus {
sync: InProcessEventBus,
async_lane: AsyncLaneEventBus,
}
#[derive(Clone)]
enum AsyncLaneEventBus {
InProcess(InProcessEventBus),
External(Arc<dyn EventBus>),
}
impl TwoLaneEventBus {
pub fn new() -> Self {
Self {
sync: InProcessEventBus::new(),
async_lane: AsyncLaneEventBus::InProcess(InProcessEventBus::new()),
}
}
pub fn with_async_bus(async_bus: Arc<dyn EventBus>) -> Self {
Self {
sync: InProcessEventBus::new(),
async_lane: AsyncLaneEventBus::External(async_bus),
}
}
pub async fn sync_handler_count(&self) -> usize {
self.sync.handler_count().await
}
pub async fn async_handler_count(&self) -> usize {
match &self.async_lane {
AsyncLaneEventBus::InProcess(async_lane) => async_lane.handler_count().await,
AsyncLaneEventBus::External(_) => 0,
}
}
}
impl Default for TwoLaneEventBus {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
impl EventBus for TwoLaneEventBus {
async fn publish(&self, events: Vec<Event>) -> EventBusResult<()> {
self.sync.publish(events.clone()).await?;
let async_lane = self.async_lane.clone();
match async_lane {
AsyncLaneEventBus::InProcess(async_lane) => {
let handle = tokio::spawn(async move {
if let Err(error) = async_lane.publish(events).await {
tracing::warn!(error = ?error, "async event handler failed");
}
});
drop(handle);
}
AsyncLaneEventBus::External(async_bus) => async_bus.publish(events).await?,
}
Ok(())
}
async fn subscribe(&mut self, handler: Box<dyn EventHandler>) -> EventBusResult<()> {
match handler.lane() {
HandlerLane::Sync => self.sync.subscribe(handler).await,
HandlerLane::Async => match &mut self.async_lane {
AsyncLaneEventBus::InProcess(async_lane) => async_lane.subscribe(handler).await,
AsyncLaneEventBus::External(_) => Ok(()),
},
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use std::sync::Arc;
use tokio::sync::Mutex as TokioMutex;
struct CountingHandler {
count: Arc<TokioMutex<usize>>,
event_types: Vec<String>,
}
impl CountingHandler {
fn new(event_types: Vec<String>) -> Self {
Self {
count: Arc::new(TokioMutex::new(0)),
event_types,
}
}
}
struct LaneCountingHandler {
count: Arc<TokioMutex<usize>>,
event_types: Vec<String>,
lane: HandlerLane,
}
#[async_trait]
impl EventHandler for LaneCountingHandler {
fn handles(&self) -> Vec<String> {
self.event_types.clone()
}
async fn handle(
&self,
_event: &Event,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let mut count = self.count.lock().await;
*count += 1;
Ok(())
}
fn lane(&self) -> HandlerLane {
self.lane
}
}
#[async_trait]
impl EventHandler for CountingHandler {
fn handles(&self) -> Vec<String> {
self.event_types.clone()
}
async fn handle(
&self,
_event: &Event,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let mut count = self.count.lock().await;
*count += 1;
Ok(())
}
}
struct FailingHandler {
fail_on: String,
}
struct LaneFailingHandler {
fail_on: String,
lane: HandlerLane,
}
#[async_trait]
impl EventHandler for FailingHandler {
fn handles(&self) -> Vec<String> {
vec![self.fail_on.clone()]
}
async fn handle(
&self,
_event: &Event,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
Err("Intentional test failure".into())
}
}
#[async_trait]
impl EventHandler for LaneFailingHandler {
fn handles(&self) -> Vec<String> {
vec![self.fail_on.clone()]
}
async fn handle(
&self,
_event: &Event,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
Err("Intentional test failure".into())
}
fn lane(&self) -> HandlerLane {
self.lane
}
}
#[tokio::test]
async fn test_new_event_bus() {
let bus = InProcessEventBus::new();
assert_eq!(bus.handler_count().await, 0);
}
#[tokio::test]
async fn test_subscribe_handler() {
let mut bus = InProcessEventBus::new();
let handler = Box::new(CountingHandler::new(vec!["UserCreated".to_string()]));
bus.subscribe(handler).await.unwrap();
assert_eq!(bus.handler_count().await, 1);
}
#[tokio::test]
async fn test_subscribe_multiple_handlers() {
let mut bus = InProcessEventBus::new();
bus.subscribe(Box::new(CountingHandler::new(vec![
"UserCreated".to_string()
])))
.await
.unwrap();
bus.subscribe(Box::new(CountingHandler::new(vec![
"UserUpdated".to_string()
])))
.await
.unwrap();
assert_eq!(bus.handler_count().await, 2);
}
#[tokio::test]
async fn test_publish_single_event() {
let mut bus = InProcessEventBus::new();
let counter = Arc::new(TokioMutex::new(0));
let counter_clone = counter.clone();
let handler = CountingHandler {
count: counter_clone,
event_types: vec!["UserCreated".to_string()],
};
bus.subscribe(Box::new(handler)).await.unwrap();
let event = Event::new(
"User",
"user-123",
1,
"UserCreated",
json!({ "name": "Alice" }),
);
bus.publish(vec![event]).await.unwrap();
let count = *counter.lock().await;
assert_eq!(count, 1);
}
#[tokio::test]
async fn test_publish_multiple_events() {
let mut bus = InProcessEventBus::new();
let counter = Arc::new(TokioMutex::new(0));
let counter_clone = counter.clone();
let handler = CountingHandler {
count: counter_clone,
event_types: vec!["UserCreated".to_string(), "UserUpdated".to_string()],
};
bus.subscribe(Box::new(handler)).await.unwrap();
let events = vec![
Event::new(
"User",
"user-1",
1,
"UserCreated",
json!({ "name": "Alice" }),
),
Event::new(
"User",
"user-1",
2,
"UserUpdated",
json!({ "name": "Alice Smith" }),
),
Event::new("User", "user-2", 1, "UserCreated", json!({ "name": "Bob" })),
];
bus.publish(events).await.unwrap();
let count = *counter.lock().await;
assert_eq!(count, 3);
}
#[tokio::test]
async fn test_handler_filters_event_types() {
let mut bus = InProcessEventBus::new();
let counter = Arc::new(TokioMutex::new(0));
let counter_clone = counter.clone();
let handler = CountingHandler {
count: counter_clone,
event_types: vec!["UserCreated".to_string()],
};
bus.subscribe(Box::new(handler)).await.unwrap();
let events = vec![
Event::new("User", "user-1", 1, "UserCreated", json!({})),
Event::new("User", "user-1", 2, "UserUpdated", json!({})),
Event::new("User", "user-1", 3, "UserDeleted", json!({})),
];
bus.publish(events).await.unwrap();
let count = *counter.lock().await;
assert_eq!(count, 1);
}
#[tokio::test]
async fn test_multiple_handlers_same_event() {
let mut bus = InProcessEventBus::new();
let counter1 = Arc::new(TokioMutex::new(0));
let counter2 = Arc::new(TokioMutex::new(0));
let handler1 = CountingHandler {
count: counter1.clone(),
event_types: vec!["UserCreated".to_string()],
};
let handler2 = CountingHandler {
count: counter2.clone(),
event_types: vec!["UserCreated".to_string()],
};
bus.subscribe(Box::new(handler1)).await.unwrap();
bus.subscribe(Box::new(handler2)).await.unwrap();
let event = Event::new("User", "user-1", 1, "UserCreated", json!({}));
bus.publish(vec![event]).await.unwrap();
assert_eq!(*counter1.lock().await, 1);
assert_eq!(*counter2.lock().await, 1);
}
#[tokio::test]
async fn test_handler_failure_propagates() {
let mut bus = InProcessEventBus::new();
let failing_handler = Box::new(FailingHandler {
fail_on: "UserCreated".to_string(),
});
bus.subscribe(failing_handler).await.unwrap();
let event = Event::new("User", "user-1", 1, "UserCreated", json!({}));
let result = bus.publish(vec![event]).await;
assert!(result.is_err());
match result.unwrap_err() {
EventBusError::HandlerFailed {
event_type,
event_id,
message,
} => {
assert_eq!(event_type, "UserCreated");
assert!(!event_id.is_empty());
assert!(message.contains("Intentional test failure"));
}
_ => panic!("Expected HandlerFailed error"),
}
}
#[tokio::test]
async fn test_event_handler_default_lane_is_sync() {
let handler = CountingHandler::new(vec!["UserCreated".to_string()]);
assert_eq!(handler.lane(), HandlerLane::Sync);
}
#[tokio::test]
async fn test_two_lane_sync_handler_failure_propagates() {
let mut bus = TwoLaneEventBus::new();
bus.subscribe(Box::new(LaneFailingHandler {
fail_on: "UserCreated".to_string(),
lane: HandlerLane::Sync,
}))
.await
.unwrap();
let event = Event::new("User", "user-1", 1, "UserCreated", json!({}));
let result = bus.publish(vec![event]).await;
assert!(matches!(result, Err(EventBusError::HandlerFailed { .. })));
}
#[tokio::test]
async fn test_two_lane_async_handler_failure_does_not_propagate() {
let mut bus = TwoLaneEventBus::new();
bus.subscribe(Box::new(LaneFailingHandler {
fail_on: "UserCreated".to_string(),
lane: HandlerLane::Async,
}))
.await
.unwrap();
let event = Event::new("User", "user-1", 1, "UserCreated", json!({}));
let result = bus.publish(vec![event]).await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_two_lane_routes_handlers_by_lane() {
let mut bus = TwoLaneEventBus::new();
let sync_count = Arc::new(TokioMutex::new(0));
let async_count = Arc::new(TokioMutex::new(0));
bus.subscribe(Box::new(LaneCountingHandler {
count: sync_count,
event_types: vec!["UserCreated".to_string()],
lane: HandlerLane::Sync,
}))
.await
.unwrap();
bus.subscribe(Box::new(LaneCountingHandler {
count: async_count,
event_types: vec!["UserCreated".to_string()],
lane: HandlerLane::Async,
}))
.await
.unwrap();
assert_eq!(bus.sync_handler_count().await, 1);
assert_eq!(bus.async_handler_count().await, 1);
}
#[tokio::test]
async fn test_no_handlers_for_event_type() {
let bus = InProcessEventBus::new();
let event = Event::new("User", "user-1", 1, "UserCreated", json!({}));
let result = bus.publish(vec![event]).await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_handler_called_in_order() {
let mut bus = InProcessEventBus::new();
let order = Arc::new(TokioMutex::new(Vec::new()));
struct OrderTracker {
id: usize,
order: Arc<TokioMutex<Vec<usize>>>,
}
#[async_trait]
impl EventHandler for OrderTracker {
fn handles(&self) -> Vec<String> {
vec!["TestEvent".to_string()]
}
async fn handle(
&self,
_event: &Event,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
self.order.lock().await.push(self.id);
Ok(())
}
}
for i in 1..=3 {
bus.subscribe(Box::new(OrderTracker {
id: i,
order: order.clone(),
}))
.await
.unwrap();
}
let event = Event::new("Test", "test-1", 1, "TestEvent", json!({}));
bus.publish(vec![event]).await.unwrap();
let call_order = order.lock().await;
assert_eq!(*call_order, vec![1, 2, 3]);
}
#[tokio::test]
async fn test_two_lane_sync_handlers_called_in_order() {
let mut bus = TwoLaneEventBus::new();
let order = Arc::new(TokioMutex::new(Vec::new()));
struct OrderTracker {
id: usize,
order: Arc<TokioMutex<Vec<usize>>>,
}
#[async_trait]
impl EventHandler for OrderTracker {
fn handles(&self) -> Vec<String> {
vec!["TestEvent".to_string()]
}
async fn handle(
&self,
_event: &Event,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
self.order.lock().await.push(self.id);
Ok(())
}
}
for i in 1..=3 {
bus.subscribe(Box::new(OrderTracker {
id: i,
order: order.clone(),
}))
.await
.unwrap();
}
let event = Event::new("Test", "test-1", 1, "TestEvent", json!({}));
bus.publish(vec![event]).await.unwrap();
let call_order = order.lock().await;
assert_eq!(*call_order, vec![1, 2, 3]);
}
#[tokio::test]
async fn test_event_bus_clone() {
let mut bus1 = InProcessEventBus::new();
let counter = Arc::new(TokioMutex::new(0));
let handler = CountingHandler {
count: counter.clone(),
event_types: vec!["UserCreated".to_string()],
};
bus1.subscribe(Box::new(handler)).await.unwrap();
let bus2 = bus1.clone();
assert_eq!(bus1.handler_count().await, 1);
assert_eq!(bus2.handler_count().await, 1);
let event = Event::new("User", "user-1", 1, "UserCreated", json!({}));
bus2.publish(vec![event]).await.unwrap();
assert_eq!(*counter.lock().await, 1);
}
#[test]
fn test_error_messages() {
let error = EventBusError::handler_failed("UserCreated", "event-123", "Connection timeout");
let msg = error.to_string();
assert!(msg.contains("UserCreated"));
assert!(msg.contains("event-123"));
assert!(msg.contains("Connection timeout"));
let error = EventBusError::no_handlers("UnknownEvent");
assert!(error.to_string().contains("UnknownEvent"));
let error = EventBusError::subscription_failed("Handler invalid");
assert!(error.to_string().contains("Handler invalid"));
}
}