Skip to main content

flare_core/transport/
events.rs

1//! 传输层事件模块
2//!
3//! 定义连接事件类型和观察者模式接口,用于在传输层和上层之间传递事件
4
5use crate::common::error::FlareError;
6use std::sync::{Arc, Mutex};
7
8/// 连接事件
9///
10/// 表示连接上发生的各种事件,由 `ConnectionObserver` 使用以响应连接状态和数据接收。
11#[derive(Debug, Clone)]
12pub enum ConnectionEvent {
13    /// 连接成功建立时触发
14    Connected,
15    /// 连接关闭时触发
16    ///
17    /// 参数提供断开连接的原因
18    Disconnected(String),
19    /// 收到新消息时触发
20    ///
21    /// 负载是字节向量
22    Message(Vec<u8>),
23    /// 连接上发生非致命错误时触发
24    Error(FlareError),
25}
26
27impl ConnectionEvent {
28    /// 检查是否为连接事件
29    pub fn is_connected(&self) -> bool {
30        matches!(self, Self::Connected)
31    }
32
33    /// 检查是否为断开连接事件
34    pub fn is_disconnected(&self) -> bool {
35        matches!(self, Self::Disconnected(_))
36    }
37
38    /// 检查是否为消息事件
39    pub fn is_message(&self) -> bool {
40        matches!(self, Self::Message(_))
41    }
42
43    /// 检查是否为错误事件
44    pub fn is_error(&self) -> bool {
45        matches!(self, Self::Error(_))
46    }
47
48    /// 获取断开连接的原因(如果是断开连接事件)
49    pub fn disconnect_reason(&self) -> Option<&str> {
50        match self {
51            Self::Disconnected(reason) => Some(reason),
52            _ => None,
53        }
54    }
55
56    /// 获取消息数据(如果是消息事件)
57    pub fn message_data(&self) -> Option<&[u8]> {
58        match self {
59            Self::Message(data) => Some(data),
60            _ => None,
61        }
62    }
63
64    /// 获取错误(如果是错误事件)
65    pub fn error(&self) -> Option<&FlareError> {
66        match self {
67            Self::Error(err) => Some(err),
68            _ => None,
69        }
70    }
71}
72
73/// 连接事件观察者
74///
75/// 实现此 trait 以响应连接建立、断开和接收消息等事件。
76/// 观察者通过 `Connection` 实例注册。
77pub trait ConnectionObserver: Send + Sync {
78    /// 当连接上发生事件时由 `Connection` 调用
79    ///
80    /// # 参数
81    /// - `event`: 发生的事件
82    fn on_event(&self, event: &ConnectionEvent);
83}
84
85/// 线程安全的、引用计数的观察者类型别名
86pub type ArcObserver = Arc<dyn ConnectionObserver>;
87
88pub(crate) fn notify_observers(
89    observers_arc: &Arc<Mutex<Vec<ArcObserver>>>,
90    event: &ConnectionEvent,
91    lock_name: &str,
92) {
93    let observers = match observers_arc.lock() {
94        Ok(observers) => observers.clone(),
95        Err(e) => {
96            tracing::warn!("{lock_name} lock poisoned: {e}");
97            return;
98        }
99    };
100
101    for observer in observers {
102        observer.on_event(event);
103    }
104}
105
106pub(crate) fn notify_observers_and_clear(
107    observers_arc: &Arc<Mutex<Vec<ArcObserver>>>,
108    event: &ConnectionEvent,
109    lock_name: &str,
110) {
111    let observers = match observers_arc.lock() {
112        Ok(mut observers) => std::mem::take(&mut *observers),
113        Err(e) => {
114            tracing::warn!("{lock_name} lock poisoned: {e}");
115            return;
116        }
117    };
118
119    for observer in observers {
120        observer.on_event(event);
121    }
122}
123
124// ============================================================================
125// 便利观察者实现
126// ============================================================================
127
128/// 空观察者(不执行任何操作)
129///
130/// 用于需要观察者但不需要实际处理的情况
131pub struct NoOpObserver;
132
133impl ConnectionObserver for NoOpObserver {
134    fn on_event(&self, _event: &ConnectionEvent) {
135        // 不执行任何操作
136    }
137}
138
139impl NoOpObserver {
140    /// 创建新的空观察者
141    #[allow(clippy::new_ret_no_self)]
142    pub fn new() -> ArcObserver {
143        Arc::new(Self)
144    }
145}
146
147/// 日志观察者(记录所有事件到日志)
148///
149/// 用于调试和监控连接事件
150pub struct LoggingObserver {
151    prefix: String,
152}
153
154impl LoggingObserver {
155    /// 创建新的日志观察者
156    ///
157    /// # 参数
158    /// - `prefix`: 日志前缀,用于标识不同的观察者实例
159    #[allow(clippy::new_ret_no_self)]
160    pub fn new(prefix: impl Into<String>) -> ArcObserver {
161        Arc::new(Self {
162            prefix: prefix.into(),
163        })
164    }
165}
166
167impl ConnectionObserver for LoggingObserver {
168    fn on_event(&self, event: &ConnectionEvent) {
169        match event {
170            ConnectionEvent::Connected => {
171                tracing::info!("[{}] Connection established", self.prefix);
172            }
173            ConnectionEvent::Disconnected(reason) => {
174                tracing::info!("[{}] Connection disconnected: {}", self.prefix, reason);
175            }
176            ConnectionEvent::Message(data) => {
177                tracing::debug!("[{}] Message received: {} bytes", self.prefix, data.len());
178            }
179            ConnectionEvent::Error(err) => {
180                tracing::error!("[{}] Connection error: {:?}", self.prefix, err);
181            }
182        }
183    }
184}
185
186/// 组合观察者(将事件转发给多个观察者)
187///
188/// 用于需要多个观察者处理同一事件的情况
189pub struct CompositeObserver {
190    observers: Vec<ArcObserver>,
191}
192
193impl CompositeObserver {
194    /// 创建新的组合观察者
195    pub fn new() -> Self {
196        Self {
197            observers: Vec::new(),
198        }
199    }
200
201    /// 添加观察者
202    pub fn add(&mut self, observer: ArcObserver) {
203        self.observers.push(observer);
204    }
205
206    /// 创建 Arc 包装的组合观察者
207    pub fn into_arc(self) -> ArcObserver {
208        Arc::new(self)
209    }
210}
211
212impl ConnectionObserver for CompositeObserver {
213    fn on_event(&self, event: &ConnectionEvent) {
214        for observer in &self.observers {
215            observer.on_event(event);
216        }
217    }
218}
219
220impl Default for CompositeObserver {
221    fn default() -> Self {
222        Self::new()
223    }
224}
225
226#[cfg(test)]
227mod tests {
228    use super::*;
229    use std::sync::atomic::{AtomicUsize, Ordering};
230
231    struct CountingObserver {
232        calls: Arc<AtomicUsize>,
233        drops: Arc<AtomicUsize>,
234    }
235
236    impl ConnectionObserver for CountingObserver {
237        fn on_event(&self, _event: &ConnectionEvent) {
238            self.calls.fetch_add(1, Ordering::SeqCst);
239        }
240    }
241
242    impl Drop for CountingObserver {
243        fn drop(&mut self) {
244            self.drops.fetch_add(1, Ordering::SeqCst);
245        }
246    }
247
248    #[test]
249    fn terminal_notify_clears_and_drops_observers() {
250        let calls = Arc::new(AtomicUsize::new(0));
251        let drops = Arc::new(AtomicUsize::new(0));
252        let observers = Arc::new(Mutex::new(Vec::<ArcObserver>::new()));
253        observers.lock().unwrap().push(Arc::new(CountingObserver {
254            calls: Arc::clone(&calls),
255            drops: Arc::clone(&drops),
256        }));
257
258        notify_observers_and_clear(
259            &observers,
260            &ConnectionEvent::Disconnected("test".to_string()),
261            "test observers",
262        );
263
264        assert_eq!(calls.load(Ordering::SeqCst), 1);
265        assert_eq!(drops.load(Ordering::SeqCst), 1);
266        assert!(observers.lock().unwrap().is_empty());
267    }
268
269    #[test]
270    fn non_terminal_notify_keeps_observers_registered() {
271        let calls = Arc::new(AtomicUsize::new(0));
272        let drops = Arc::new(AtomicUsize::new(0));
273        let observers = Arc::new(Mutex::new(Vec::<ArcObserver>::new()));
274        observers.lock().unwrap().push(Arc::new(CountingObserver {
275            calls: Arc::clone(&calls),
276            drops: Arc::clone(&drops),
277        }));
278
279        notify_observers(&observers, &ConnectionEvent::Connected, "test observers");
280
281        assert_eq!(calls.load(Ordering::SeqCst), 1);
282        assert_eq!(drops.load(Ordering::SeqCst), 0);
283        assert_eq!(observers.lock().unwrap().len(), 1);
284    }
285}