1use std::sync::Arc;
4
5use crate::bus::{EventBus, RetryPolicy, DeadLetterHandler};
6use crate::execution_log::ExecutionLog;
7use crate::registry::EventRegistry;
8use crate::telemetry::Telemetry;
9
10pub struct EventBusBuilder {
27 capacity: usize,
28 telemetry: Option<Arc<dyn Telemetry>>,
29 registry: Option<Arc<EventRegistry>>,
30 execution_log: Option<Arc<ExecutionLog>>,
31 retry_policy: RetryPolicy,
32 dead_letter: Option<Arc<dyn DeadLetterHandler>>,
33}
34
35impl EventBusBuilder {
36 pub fn new() -> Self {
40 Self {
41 capacity: 1024,
42 telemetry: None,
43 registry: None,
44 execution_log: None,
45 retry_policy: RetryPolicy::default(),
46 dead_letter: None,
47 }
48 }
49
50 pub fn capacity(mut self, capacity: usize) -> Self {
54 self.capacity = capacity;
55 self
56 }
57
58 pub fn telemetry<T: Telemetry>(mut self, telemetry: T) -> Self {
62 self.telemetry = Some(Arc::new(telemetry));
63 self
64 }
65
66 pub fn registry(mut self, registry: Arc<EventRegistry>) -> Self {
70 self.registry = Some(registry);
71 self
72 }
73
74 pub fn execution_log(mut self, log: Arc<ExecutionLog>) -> Self {
79 self.execution_log = Some(log);
80 self
81 }
82
83 pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
87 self.retry_policy = policy;
88 self
89 }
90
91 pub fn dead_letter_handler<H: DeadLetterHandler>(mut self, handler: H) -> Self {
95 self.dead_letter = Some(Arc::new(handler));
96 self
97 }
98
99 pub fn build(self) -> EventBus {
101 let registry = self.registry.unwrap_or_else(|| Arc::new(EventRegistry::new()));
102 EventBus::from_builder(
103 self.capacity,
104 self.telemetry,
105 registry,
106 self.execution_log,
107 self.retry_policy,
108 self.dead_letter,
109 )
110 }
111}
112
113impl Default for EventBusBuilder {
114 fn default() -> Self {
115 Self::new()
116 }
117}
118
119#[cfg(test)]
120mod tests {
121 use super::*;
122 use crate::telemetry::TracingTelemetry;
123
124 #[test]
125 fn builder_default_capacity() {
126 let builder = EventBusBuilder::new();
127 assert_eq!(builder.capacity, 1024);
128 assert!(builder.telemetry.is_none());
129 assert!(builder.registry.is_none());
130 }
131
132 #[test]
133 fn builder_custom_capacity() {
134 let builder = EventBusBuilder::new().capacity(2048);
135 assert_eq!(builder.capacity, 2048);
136 }
137
138 #[test]
139 fn builder_with_telemetry() {
140 let builder = EventBusBuilder::new().telemetry(TracingTelemetry);
141 assert!(builder.telemetry.is_some());
142 }
143
144 #[test]
145 fn builder_with_registry() {
146 let registry = Arc::new(EventRegistry::new());
147 let builder = EventBusBuilder::new().registry(registry);
148 assert!(builder.registry.is_some());
149 }
150
151 #[test]
152 fn builder_builds_event_bus() {
153 let _bus = EventBusBuilder::new().capacity(512).build();
154 }
155
156 #[test]
157 fn builder_builds_with_custom_registry() {
158 let registry = Arc::new(EventRegistry::new());
159 let bus = EventBusBuilder::new().registry(registry.clone()).build();
160 assert!(Arc::ptr_eq(bus.registry(), ®istry));
161 }
162
163 #[test]
164 fn builder_with_execution_log() {
165 let log = Arc::new(ExecutionLog::in_memory());
166 let builder = EventBusBuilder::new().execution_log(log);
167 assert!(builder.execution_log.is_some());
168 }
169
170 #[test]
171 fn builder_builds_with_execution_log() {
172 let log = Arc::new(ExecutionLog::in_memory());
173 let bus = EventBusBuilder::new().execution_log(log.clone()).build();
174 assert!(bus.execution_log().is_some());
175 assert!(Arc::ptr_eq(bus.execution_log().unwrap(), &log));
176 }
177}