Skip to main content

flare_core/client/
manager.rs

1//! 客户端连接管理器
2//!
3//! 统一管理客户端连接、心跳、重连等生命周期
4//! 提供自动重连、心跳管理、连接状态监控等功能
5
6use crate::client::config::ClientConfig;
7use crate::client::connection::ConnectionStateManager;
8use crate::client::heartbeat::HeartbeatManager;
9use crate::client::transports::Client;
10use crate::common::MessageParser;
11use crate::common::error::Result;
12use crate::common::platform::sleep;
13use crate::common::protocol::Frame;
14use crate::transport::events::{ArcObserver, ConnectionEvent};
15
16use std::sync::{Arc, Mutex as StdMutex};
17use std::time::Duration;
18use tokio::sync::Mutex;
19use tracing::{debug, error, info, warn};
20
21/// 客户端连接管理器
22///
23/// 管理客户端连接的生命周期,包括:
24/// - 自动连接和重连
25/// - 心跳管理
26/// - 连接状态监控
27/// - 消息观察者管理
28pub struct ClientConnectionManager {
29    /// 客户端配置
30    config: ClientConfig,
31    /// 客户端实例
32    client: Arc<Mutex<Box<dyn Client>>>,
33    /// 连接状态管理器
34    state_manager: Arc<ConnectionStateManager>,
35    /// 心跳管理器
36    heartbeat_manager: Arc<Mutex<Option<HeartbeatManager>>>,
37    /// 消息解析器
38    #[allow(dead_code)] // 保留用于未来扩展
39    parser: MessageParser,
40    /// 观察者列表
41    observers: Arc<StdMutex<Vec<ArcObserver>>>,
42    /// 重连任务句柄
43    reconnect_handle: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>,
44    /// 是否正在重连
45    is_reconnecting: Arc<Mutex<bool>>,
46}
47
48impl ClientConnectionManager {
49    /// 创建新的客户端连接管理器
50    ///
51    /// # 参数
52    /// - `client`: 客户端实例(可以是 WebSocketClient、QUICClient 等)
53    /// - `config`: 客户端配置
54    pub fn new(client: Box<dyn Client>, config: ClientConfig) -> Self {
55        let parser = MessageParser::new(
56            config.serialization_format,
57            config.compression.clone(),
58            crate::common::encryption::EncryptionAlgorithm::None,
59        );
60
61        Self {
62            config,
63            client: Arc::new(Mutex::new(client)),
64            state_manager: Arc::new(ConnectionStateManager::new()),
65            heartbeat_manager: Arc::new(Mutex::new(None)),
66            parser,
67            observers: Arc::new(StdMutex::new(Vec::new())),
68            reconnect_handle: Arc::new(Mutex::new(None)),
69            is_reconnecting: Arc::new(Mutex::new(false)),
70        }
71    }
72
73    /// 连接到服务器
74    ///
75    /// 自动处理连接、心跳启动等
76    pub async fn connect(&self) -> Result<()> {
77        info!("Connecting to server: {}", self.config.server_url);
78
79        let mut client = self.client.lock().await;
80
81        // 检查是否已连接
82        if client.is_connected() {
83            debug!("Already connected, skipping connect");
84            return Ok(());
85        }
86
87        // 设置状态为连接中
88        self.state_manager
89            .set_state(crate::client::connection::ConnectionState::Connecting);
90
91        // 执行连接
92        match client.connect().await {
93            Ok(()) => {
94                info!("Successfully connected to server");
95                self.state_manager
96                    .set_state(crate::client::connection::ConnectionState::Connected);
97
98                // 启动心跳(如果启用)
99                if self.config.heartbeat.enabled {
100                    self.start_heartbeat().await?;
101                }
102
103                // 通知观察者
104                self.notify_observers(&ConnectionEvent::Connected);
105
106                Ok(())
107            }
108            Err(e) => {
109                error!("Failed to connect: {}", e);
110                self.state_manager
111                    .set_state(crate::client::connection::ConnectionState::Disconnected);
112                Err(e)
113            }
114        }
115    }
116
117    /// 断开连接
118    pub async fn disconnect(&self) -> Result<()> {
119        info!("Disconnecting from server");
120
121        // 停止重连
122        self.stop_reconnect().await;
123
124        // 停止心跳
125        self.stop_heartbeat().await;
126
127        // 断开连接
128        let mut client = self.client.lock().await;
129        let result = client.disconnect().await;
130
131        self.state_manager
132            .set_state(crate::client::connection::ConnectionState::Disconnected);
133        self.notify_observers(&ConnectionEvent::Disconnected(String::new()));
134
135        result
136    }
137
138    /// 发送消息
139    pub async fn send_frame(&self, frame: &Frame) -> Result<()> {
140        if !self.state_manager.get_state().can_send() {
141            return Err(crate::common::error::FlareError::connection_failed(
142                "Not connected".to_string(),
143            ));
144        }
145
146        let mut client = self.client.lock().await;
147        client.send_frame(frame).await
148    }
149
150    /// 启动心跳
151    async fn start_heartbeat(&self) -> Result<()> {
152        if !self.config.heartbeat.enabled {
153            return Ok(());
154        }
155
156        debug!(
157            "Starting heartbeat: interval={:?}, timeout={:?}",
158            self.config.heartbeat.interval, self.config.heartbeat.timeout
159        );
160
161        let heartbeat = HeartbeatManager::new(
162            self.config.heartbeat.interval,
163            self.config.heartbeat.timeout,
164        );
165
166        // 需要从 Client 获取 Connection,但 Client trait 没有提供这个方法
167        // 这里我们需要一个不同的设计
168        // 暂时跳过,因为心跳应该在 Client 实现内部管理
169        // 或者我们需要扩展 Client trait
170
171        let mut hb_mgr = self.heartbeat_manager.lock().await;
172        *hb_mgr = Some(heartbeat);
173
174        Ok(())
175    }
176
177    /// 停止心跳
178    async fn stop_heartbeat(&self) {
179        let mut hb_mgr = self.heartbeat_manager.lock().await;
180        if let Some(mut hb) = hb_mgr.take() {
181            hb.stop();
182        }
183    }
184
185    /// 启动自动重连
186    pub async fn start_auto_reconnect(&self) {
187        let mut is_reconnecting = self.is_reconnecting.lock().await;
188        if *is_reconnecting {
189            return; // 已经启动了
190        }
191        *is_reconnecting = true;
192        drop(is_reconnecting);
193
194        let client = Arc::clone(&self.client);
195        let state_mgr = Arc::clone(&self.state_manager);
196        let config = self.config.clone();
197        let heartbeat_cfg = self.config.heartbeat.clone();
198        let heartbeat_mgr = Arc::clone(&self.heartbeat_manager);
199        let observers = Arc::clone(&self.observers);
200        let reconnect_handle = Arc::clone(&self.reconnect_handle);
201        let is_reconnecting_flag = Arc::clone(&self.is_reconnecting);
202
203        let handle = tokio::spawn(async move {
204            loop {
205                // 检查连接状态
206                let should_reconnect = {
207                    let client_guard = client.lock().await;
208                    !client_guard.is_connected()
209                        && matches!(
210                            state_mgr.get_state(),
211                            crate::client::connection::ConnectionState::Disconnected
212                        )
213                };
214
215                if !should_reconnect {
216                    // 检查是否超过最大重连次数
217                    // 这里简化处理,实际应该在连接失败时触发
218                    sleep(Duration::from_secs(1)).await;
219                    continue;
220                }
221
222                // 检查重连次数限制
223                // 注意:这里简化了,实际应该在连接失败时计数
224                // 这里我们假设只要状态是断开就重连
225
226                info!("Attempting to reconnect...");
227                state_mgr.set_state(crate::client::connection::ConnectionState::Connecting);
228
229                // 尝试重连
230                let reconnect_result = {
231                    let mut client_guard = client.lock().await;
232                    client_guard.connect().await
233                };
234
235                match reconnect_result {
236                    Ok(()) => {
237                        info!("Reconnected successfully");
238                        state_mgr.set_state(crate::client::connection::ConnectionState::Connected);
239
240                        // 重新启动心跳
241                        if heartbeat_cfg.enabled {
242                            let _hb_mgr = heartbeat_mgr.lock().await;
243                            // 重新创建心跳管理器
244                            // 注意:这里简化了,实际应该从 Client 获取 Connection
245                        }
246
247                        // 通知观察者
248                        {
249                            let observers_guard = observers.lock().unwrap();
250                            for observer in observers_guard.iter() {
251                                observer.on_event(&ConnectionEvent::Connected);
252                            }
253                        }
254
255                        // 重连成功,退出重连循环
256                        *is_reconnecting_flag.lock().await = false;
257                        break;
258                    }
259                    Err(e) => {
260                        warn!(
261                            "Reconnect failed: {}, retrying in {:?}",
262                            e, config.reconnect_interval
263                        );
264                        state_mgr
265                            .set_state(crate::client::connection::ConnectionState::Disconnected);
266                        sleep(config.reconnect_interval).await;
267                    }
268                }
269            }
270
271            // 清理句柄
272            let mut handle_guard = reconnect_handle.lock().await;
273            *handle_guard = None;
274        });
275
276        let mut handle_guard = self.reconnect_handle.lock().await;
277        *handle_guard = Some(handle);
278    }
279
280    /// 停止自动重连
281    async fn stop_reconnect(&self) {
282        let mut handle_guard = self.reconnect_handle.lock().await;
283        if let Some(handle) = handle_guard.take() {
284            handle.abort();
285        }
286
287        *self.is_reconnecting.lock().await = false;
288    }
289
290    /// 添加观察者
291    pub fn add_observer(&self, observer: ArcObserver) {
292        let mut observers = self.observers.lock().unwrap();
293        observers.push(observer);
294    }
295
296    /// 移除观察者
297    pub fn remove_observer(&self, observer: ArcObserver) {
298        let mut observers = self.observers.lock().unwrap();
299        observers.retain(|o| !Arc::ptr_eq(o, &observer));
300    }
301
302    /// 通知所有观察者
303    fn notify_observers(&self, event: &ConnectionEvent) {
304        let observers = self.observers.lock().unwrap();
305        for observer in observers.iter() {
306            observer.on_event(event);
307        }
308    }
309
310    /// 检查是否已连接
311    pub async fn is_connected(&self) -> bool {
312        let client = self.client.lock().await;
313        client.is_connected()
314    }
315
316    /// 获取连接 ID
317    pub async fn connection_id(&self) -> Option<String> {
318        let client = self.client.lock().await;
319        client.connection_id()
320    }
321
322    /// 获取连接状态
323    pub fn state(&self) -> crate::client::connection::ConnectionState {
324        self.state_manager.get_state()
325    }
326}