use crate::client::connection::ConnectionStateManager;
use crate::client::transports::ClientCore;
use crate::common::error::{FlareError, Result};
use crate::common::platform::sleep;
use crate::common::protocol::Frame;
use crate::transport::connection::Connection;
use crate::transport::events::{ConnectionEvent, ConnectionObserver};
use std::sync::Arc;
use tokio::sync::Mutex;
pub struct ClientMessageObserver {
core: Arc<ClientCore>,
}
impl ClientMessageObserver {
pub fn new(core: Arc<ClientCore>) -> Self {
Self { core }
}
}
impl ConnectionObserver for ClientMessageObserver {
fn on_event(&self, event: &ConnectionEvent) {
match event {
ConnectionEvent::Message(data) => {
let core = Arc::clone(&self.core);
let data_clone = data.clone();
#[cfg(target_arch = "wasm32")]
{
core.push_wasm_inbound(data_clone);
let core_drain = Arc::clone(&core);
crate::client::wasm_tokio::spawn_detached(async move {
core_drain.drain_wasm_inbound().await;
});
}
#[cfg(not(target_arch = "wasm32"))]
{
crate::client::runtime::spawn_client_task(async move {
core.handle_message(data_clone).await;
});
}
}
ConnectionEvent::Connected
| ConnectionEvent::Disconnected(_)
| ConnectionEvent::Error(_) => {
self.core.handle_connection_event(event);
}
}
}
}
pub struct ClientConnectionHelper;
impl ClientConnectionHelper {
pub async fn send_frame_internal(
core: &ClientCore,
connection: Option<&Arc<Mutex<Box<dyn Connection>>>>,
frame: &Frame,
) -> Result<()> {
if !core.can_send() {
return Err(FlareError::connection_failed(
"Cannot send: connection state is not ready".to_string(),
));
}
let negotiation_completed = core.is_negotiation_completed();
let parser = core.parser.lock().await;
tracing::trace!(
"[ClientConnectionHelper] 发送消息: message_id={}, format={:?}",
frame.message_id,
parser.default_format()
);
if !negotiation_completed {
tracing::warn!(
"[ClientConnectionHelper] ⚠️ 协商未完成但尝试发送消息: message_id={}, format={:?}, compression={:?}, encryption={:?}",
frame.message_id,
parser.default_format(),
parser.default_compression(),
parser.default_encryption()
);
}
let data = parser.serialize(frame)?;
drop(parser);
let conn =
connection.ok_or_else(|| FlareError::connection_failed("Not connected".to_string()))?;
let mut c = conn.lock().await;
c.send(&data).await?;
Ok(())
}
#[allow(dead_code)] pub async fn try_reconnect<F, Fut>(
reconnect_attempts: &mut u32,
max_attempts: Option<u32>,
reconnect_interval: std::time::Duration,
old_connection: Option<Arc<Mutex<Box<dyn Connection>>>>,
state_manager: &Arc<ConnectionStateManager>,
connect_fn: F,
) -> Result<()>
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = Result<()>>,
{
if let Some(max) = max_attempts
&& *reconnect_attempts >= max
{
return Err(FlareError::connection_failed(format!(
"Max reconnect attempts ({}) exceeded",
max
)));
}
state_manager.start_connecting();
*reconnect_attempts += 1;
sleep(reconnect_interval).await;
if let Some(conn) = old_connection {
let mut c = conn.lock().await;
let _ = c.close().await;
}
connect_fn().await
}
pub async fn disconnect_internal(
connection: Option<Arc<Mutex<Box<dyn Connection>>>>,
core: &mut ClientCore,
) -> Result<()> {
core.state_manager
.set_state(crate::client::connection::ConnectionState::Disconnecting);
core.set_disconnect_requested(true);
core.stop_heartbeat();
core.clear_client_connection();
let close_result = if let Some(conn) = connection {
let mut c = conn.lock().await;
c.close().await
} else {
Ok(())
};
core.cancel_all_pending_responses().await;
core.handle_connection_event(&ConnectionEvent::Disconnected(
"Client disconnected".to_string(),
));
close_result
}
pub fn can_reconnect(max_attempts: Option<u32>) -> bool {
max_attempts.map(|n| n > 0).unwrap_or(true)
}
pub async fn setup_connection_and_send_connect(
connection: Arc<Mutex<Box<dyn Connection>>>,
core: &mut ClientCore,
observer: Arc<dyn ConnectionObserver>,
) -> Result<()> {
core.set_client_connection(Arc::clone(&connection));
{
let mut conn = connection.lock().await;
conn.add_observer(observer);
}
core.send_connect_message(Arc::clone(&connection)).await?;
Ok(())
}
}