Skip to main content

flare_core/client/events/
observer.rs

1//! 默认客户端消息观察者
2//!
3//! 提供通用的客户端消息和事件处理逻辑
4
5use crate::client::connection::ConnectionStateManager;
6use crate::client::events::handler::ClientEventHandler;
7use crate::client::transports::ClientCore;
8use crate::transport::events::{ConnectionEvent, ConnectionObserver};
9use std::sync::Arc;
10use tracing::{debug, error};
11
12/// 默认客户端消息观察者
13///
14/// 处理常见的系统命令(CONNECT_ACK, PONG, KICKED)和连接事件
15/// 其他命令类型可以委托给可选的 `ClientEventHandler`
16///
17/// 注意:协商前的消息使用全局共享的 `PRE_NEGOTIATION_PARSER`,
18/// 协商后的消息使用 `ClientCore` 中动态更新的 parser
19pub struct DefaultClientMessageObserver {
20    /// 客户端核心
21    core: Arc<ClientCore>,
22    /// 状态管理器
23    state_manager: Arc<ConnectionStateManager>,
24    /// 事件处理器(可选,用于自定义业务逻辑)
25    event_handler: Option<Arc<dyn ClientEventHandler>>,
26}
27
28impl Clone for DefaultClientMessageObserver {
29    fn clone(&self) -> Self {
30        Self {
31            core: Arc::clone(&self.core),
32            state_manager: Arc::clone(&self.state_manager),
33            event_handler: self.event_handler.clone(),
34        }
35    }
36}
37
38impl DefaultClientMessageObserver {
39    /// 创建新的默认客户端消息观察者
40    pub fn new(
41        core: Arc<ClientCore>,
42        state_manager: Arc<ConnectionStateManager>,
43        event_handler: Option<Arc<dyn ClientEventHandler>>,
44    ) -> Self {
45        Self {
46            core,
47            state_manager,
48            event_handler,
49        }
50    }
51}
52
53impl DefaultClientMessageObserver {
54    /// 处理消息事件
55    async fn handle_message_event(core: &Arc<ClientCore>, data: Vec<u8>) {
56        // 直接转发给 ClientCore 处理
57        // ClientCore 会根据协商状态使用正确的 parser
58        core.handle_message(data).await;
59    }
60
61    /// 处理连接建立事件
62    async fn handle_connected_event(
63        state_manager: &Arc<ConnectionStateManager>,
64        event_handler: Option<Arc<dyn ClientEventHandler>>,
65    ) {
66        debug!("[DefaultClientObserver] Connection established");
67        state_manager.set_connected();
68
69        // 通知事件处理器
70        if let Some(handler) = event_handler {
71            let event = ConnectionEvent::Connected;
72            crate::client::runtime::spawn_client_task(async move {
73                let _ = handler.handle_connection_event(&event).await;
74            });
75        }
76    }
77
78    /// 处理连接断开事件
79    async fn handle_disconnected_event(
80        state_manager: &Arc<ConnectionStateManager>,
81        event_handler: Option<Arc<dyn ClientEventHandler>>,
82        reason: String,
83    ) {
84        debug!(
85            "[DefaultClientObserver] Connection disconnected: {}",
86            reason
87        );
88        state_manager.set_disconnected();
89
90        // 通知事件处理器
91        if let Some(handler) = event_handler {
92            let event = ConnectionEvent::Disconnected(reason);
93            crate::client::runtime::spawn_client_task(async move {
94                let _ = handler.handle_connection_event(&event).await;
95            });
96        }
97    }
98
99    /// 处理错误事件
100    async fn handle_error_event(
101        state_manager: &Arc<ConnectionStateManager>,
102        event_handler: Option<Arc<dyn ClientEventHandler>>,
103        error: crate::common::error::FlareError,
104    ) {
105        error!("[DefaultClientObserver] Connection error: {:?}", error);
106        state_manager.set_failed();
107
108        // 通知事件处理器
109        if let Some(handler) = event_handler {
110            let event = ConnectionEvent::Error(error);
111            crate::client::runtime::spawn_client_task(async move {
112                let _ = handler.handle_connection_event(&event).await;
113            });
114        }
115    }
116}
117
118impl ConnectionObserver for DefaultClientMessageObserver {
119    fn on_event(&self, event: &ConnectionEvent) {
120        match event {
121            ConnectionEvent::Message(data) => {
122                let core = Arc::clone(&self.core);
123                let data_clone = data.clone();
124                crate::client::runtime::spawn_client_task(async move {
125                    Self::handle_message_event(&core, data_clone).await;
126                });
127            }
128            ConnectionEvent::Connected => {
129                let state_manager = Arc::clone(&self.state_manager);
130                let event_handler = self.event_handler.clone();
131                crate::client::runtime::spawn_client_task(async move {
132                    Self::handle_connected_event(&state_manager, event_handler).await;
133                });
134            }
135            ConnectionEvent::Disconnected(reason) => {
136                let state_manager = Arc::clone(&self.state_manager);
137                let event_handler = self.event_handler.clone();
138                let reason_clone = reason.clone();
139                crate::client::runtime::spawn_client_task(async move {
140                    Self::handle_disconnected_event(&state_manager, event_handler, reason_clone)
141                        .await;
142                });
143            }
144            ConnectionEvent::Error(e) => {
145                let state_manager = Arc::clone(&self.state_manager);
146                let event_handler = self.event_handler.clone();
147                let error_clone = e.clone();
148                crate::client::runtime::spawn_client_task(async move {
149                    Self::handle_error_event(&state_manager, event_handler, error_clone).await;
150                });
151            }
152        }
153    }
154}