flare_core/transport/
events.rs1use crate::common::error::FlareError;
6use std::sync::{Arc, Mutex};
7
8#[derive(Debug, Clone)]
12pub enum ConnectionEvent {
13 Connected,
15 Disconnected(String),
19 Message(Vec<u8>),
23 Error(FlareError),
25}
26
27impl ConnectionEvent {
28 pub fn is_connected(&self) -> bool {
30 matches!(self, Self::Connected)
31 }
32
33 pub fn is_disconnected(&self) -> bool {
35 matches!(self, Self::Disconnected(_))
36 }
37
38 pub fn is_message(&self) -> bool {
40 matches!(self, Self::Message(_))
41 }
42
43 pub fn is_error(&self) -> bool {
45 matches!(self, Self::Error(_))
46 }
47
48 pub fn disconnect_reason(&self) -> Option<&str> {
50 match self {
51 Self::Disconnected(reason) => Some(reason),
52 _ => None,
53 }
54 }
55
56 pub fn message_data(&self) -> Option<&[u8]> {
58 match self {
59 Self::Message(data) => Some(data),
60 _ => None,
61 }
62 }
63
64 pub fn error(&self) -> Option<&FlareError> {
66 match self {
67 Self::Error(err) => Some(err),
68 _ => None,
69 }
70 }
71}
72
73pub trait ConnectionObserver: Send + Sync {
78 fn on_event(&self, event: &ConnectionEvent);
83}
84
85pub 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
124pub struct NoOpObserver;
132
133impl ConnectionObserver for NoOpObserver {
134 fn on_event(&self, _event: &ConnectionEvent) {
135 }
137}
138
139impl NoOpObserver {
140 #[allow(clippy::new_ret_no_self)]
142 pub fn new() -> ArcObserver {
143 Arc::new(Self)
144 }
145}
146
147pub struct LoggingObserver {
151 prefix: String,
152}
153
154impl LoggingObserver {
155 #[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
186pub struct CompositeObserver {
190 observers: Vec<ArcObserver>,
191}
192
193impl CompositeObserver {
194 pub fn new() -> Self {
196 Self {
197 observers: Vec::new(),
198 }
199 }
200
201 pub fn add(&mut self, observer: ArcObserver) {
203 self.observers.push(observer);
204 }
205
206 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}