Skip to main content

anycms_event/
builder.rs

1//! [`EventBusBuilder`] 用于构建带配置的 EventBus。
2
3use std::sync::Arc;
4
5use crate::bus::{EventBus, RetryPolicy, DeadLetterHandler};
6use crate::execution_log::ExecutionLog;
7use crate::registry::EventRegistry;
8use crate::telemetry::Telemetry;
9
10/// EventBus 构建器,提供流畅的配置 API。
11///
12/// # Example
13///
14/// ```ignore
15/// use anycms_event::prelude::*;
16/// use anycms_event::telemetry::TracingTelemetry;
17/// use anycms_event::registry::EventRegistry;
18///
19/// let registry = Arc::new(EventRegistry::new());
20/// let bus = EventBus::builder()
21///     .capacity(2048)
22///     .telemetry(TracingTelemetry)
23///     .registry(registry)
24///     .build();
25/// ```
26pub 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    /// 创建一个新的构建器,使用默认值。
37    ///
38    /// 默认容量为 1024,无遥测,自动创建空注册表。
39    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    /// 设置广播通道容量。
51    ///
52    /// 容量控制慢订阅者开始被滞后(丢弃旧消息)之前可以缓冲多少消息。
53    pub fn capacity(mut self, capacity: usize) -> Self {
54        self.capacity = capacity;
55        self
56    }
57
58    /// 设置遥测层。
59    ///
60    /// 遥测回调将在事件总线的发布/订阅生命周期中被调用。
61    pub fn telemetry<T: Telemetry>(mut self, telemetry: T) -> Self {
62        self.telemetry = Some(Arc::new(telemetry));
63        self
64    }
65
66    /// 设置自定义事件注册表。
67    ///
68    /// 如果不设置,构建时会自动创建一个空注册表。
69    pub fn registry(mut self, registry: Arc<EventRegistry>) -> Self {
70        self.registry = Some(registry);
71        self
72    }
73
74    /// 设置执行日志。
75    ///
76    /// 执行日志记录事件发布和 Handler 执行的历史记录。
77    /// 通常与 [`ExecutionLogTelemetry`](crate::execution_log::ExecutionLogTelemetry) 一起使用。
78    pub fn execution_log(mut self, log: Arc<ExecutionLog>) -> Self {
79        self.execution_log = Some(log);
80        self
81    }
82
83    /// 设置 Handler 重试策略。
84    ///
85    /// 默认不重试(`max_retries: 0`)。
86    pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
87        self.retry_policy = policy;
88        self
89    }
90
91    /// 设置死信处理器。
92    ///
93    /// 当 Handler 重试耗尽后,将调用死信处理器。
94    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    /// 消费构建器并返回配置好的 [`EventBus`]。
100    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(), &registry));
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}