use std::sync::Arc;
use crate::bus::{EventBus, RetryPolicy, DeadLetterHandler};
use crate::execution_log::ExecutionLog;
use crate::registry::EventRegistry;
use crate::telemetry::Telemetry;
pub struct EventBusBuilder {
capacity: usize,
telemetry: Option<Arc<dyn Telemetry>>,
registry: Option<Arc<EventRegistry>>,
execution_log: Option<Arc<ExecutionLog>>,
retry_policy: RetryPolicy,
dead_letter: Option<Arc<dyn DeadLetterHandler>>,
}
impl EventBusBuilder {
pub fn new() -> Self {
Self {
capacity: 1024,
telemetry: None,
registry: None,
execution_log: None,
retry_policy: RetryPolicy::default(),
dead_letter: None,
}
}
pub fn capacity(mut self, capacity: usize) -> Self {
self.capacity = capacity;
self
}
pub fn telemetry<T: Telemetry>(mut self, telemetry: T) -> Self {
self.telemetry = Some(Arc::new(telemetry));
self
}
pub fn registry(mut self, registry: Arc<EventRegistry>) -> Self {
self.registry = Some(registry);
self
}
pub fn execution_log(mut self, log: Arc<ExecutionLog>) -> Self {
self.execution_log = Some(log);
self
}
pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
self.retry_policy = policy;
self
}
pub fn dead_letter_handler<H: DeadLetterHandler>(mut self, handler: H) -> Self {
self.dead_letter = Some(Arc::new(handler));
self
}
pub fn build(self) -> EventBus {
let registry = self.registry.unwrap_or_else(|| Arc::new(EventRegistry::new()));
EventBus::from_builder(
self.capacity,
self.telemetry,
registry,
self.execution_log,
self.retry_policy,
self.dead_letter,
)
}
}
impl Default for EventBusBuilder {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::telemetry::TracingTelemetry;
#[test]
fn builder_default_capacity() {
let builder = EventBusBuilder::new();
assert_eq!(builder.capacity, 1024);
assert!(builder.telemetry.is_none());
assert!(builder.registry.is_none());
}
#[test]
fn builder_custom_capacity() {
let builder = EventBusBuilder::new().capacity(2048);
assert_eq!(builder.capacity, 2048);
}
#[test]
fn builder_with_telemetry() {
let builder = EventBusBuilder::new().telemetry(TracingTelemetry);
assert!(builder.telemetry.is_some());
}
#[test]
fn builder_with_registry() {
let registry = Arc::new(EventRegistry::new());
let builder = EventBusBuilder::new().registry(registry);
assert!(builder.registry.is_some());
}
#[test]
fn builder_builds_event_bus() {
let _bus = EventBusBuilder::new().capacity(512).build();
}
#[test]
fn builder_builds_with_custom_registry() {
let registry = Arc::new(EventRegistry::new());
let bus = EventBusBuilder::new().registry(registry.clone()).build();
assert!(Arc::ptr_eq(bus.registry(), ®istry));
}
#[test]
fn builder_with_execution_log() {
let log = Arc::new(ExecutionLog::in_memory());
let builder = EventBusBuilder::new().execution_log(log);
assert!(builder.execution_log.is_some());
}
#[test]
fn builder_builds_with_execution_log() {
let log = Arc::new(ExecutionLog::in_memory());
let bus = EventBusBuilder::new().execution_log(log.clone()).build();
assert!(bus.execution_log().is_some());
assert!(Arc::ptr_eq(bus.execution_log().unwrap(), &log));
}
}