flare_core/client/
manager.rs1use 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
21pub struct ClientConnectionManager {
29 config: ClientConfig,
31 client: Arc<Mutex<Box<dyn Client>>>,
33 state_manager: Arc<ConnectionStateManager>,
35 heartbeat_manager: Arc<Mutex<Option<HeartbeatManager>>>,
37 #[allow(dead_code)] parser: MessageParser,
40 observers: Arc<StdMutex<Vec<ArcObserver>>>,
42 reconnect_handle: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>,
44 is_reconnecting: Arc<Mutex<bool>>,
46}
47
48impl ClientConnectionManager {
49 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 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 if client.is_connected() {
83 debug!("Already connected, skipping connect");
84 return Ok(());
85 }
86
87 self.state_manager
89 .set_state(crate::client::connection::ConnectionState::Connecting);
90
91 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 if self.config.heartbeat.enabled {
100 self.start_heartbeat().await?;
101 }
102
103 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 pub async fn disconnect(&self) -> Result<()> {
119 info!("Disconnecting from server");
120
121 self.stop_reconnect().await;
123
124 self.stop_heartbeat().await;
126
127 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 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 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 let mut hb_mgr = self.heartbeat_manager.lock().await;
172 *hb_mgr = Some(heartbeat);
173
174 Ok(())
175 }
176
177 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 pub async fn start_auto_reconnect(&self) {
187 let mut is_reconnecting = self.is_reconnecting.lock().await;
188 if *is_reconnecting {
189 return; }
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 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 sleep(Duration::from_secs(1)).await;
219 continue;
220 }
221
222 info!("Attempting to reconnect...");
227 state_mgr.set_state(crate::client::connection::ConnectionState::Connecting);
228
229 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 if heartbeat_cfg.enabled {
242 let _hb_mgr = heartbeat_mgr.lock().await;
243 }
246
247 {
249 let observers_guard = observers.lock().unwrap();
250 for observer in observers_guard.iter() {
251 observer.on_event(&ConnectionEvent::Connected);
252 }
253 }
254
255 *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 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 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 pub fn add_observer(&self, observer: ArcObserver) {
292 let mut observers = self.observers.lock().unwrap();
293 observers.push(observer);
294 }
295
296 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 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 pub async fn is_connected(&self) -> bool {
312 let client = self.client.lock().await;
313 client.is_connected()
314 }
315
316 pub async fn connection_id(&self) -> Option<String> {
318 let client = self.client.lock().await;
319 client.connection_id()
320 }
321
322 pub fn state(&self) -> crate::client::connection::ConnectionState {
324 self.state_manager.get_state()
325 }
326}