use crate::client::connection::ConnectionStateManager;
use crate::client::events::handler::ClientEventHandler;
use crate::client::transports::ClientCore;
use crate::transport::events::{ConnectionEvent, ConnectionObserver};
use std::sync::Arc;
use tracing::{debug, error};
pub struct DefaultClientMessageObserver {
core: Arc<ClientCore>,
state_manager: Arc<ConnectionStateManager>,
event_handler: Option<Arc<dyn ClientEventHandler>>,
}
impl Clone for DefaultClientMessageObserver {
fn clone(&self) -> Self {
Self {
core: Arc::clone(&self.core),
state_manager: Arc::clone(&self.state_manager),
event_handler: self.event_handler.clone(),
}
}
}
impl DefaultClientMessageObserver {
pub fn new(
core: Arc<ClientCore>,
state_manager: Arc<ConnectionStateManager>,
event_handler: Option<Arc<dyn ClientEventHandler>>,
) -> Self {
Self {
core,
state_manager,
event_handler,
}
}
}
impl DefaultClientMessageObserver {
async fn handle_message_event(core: &Arc<ClientCore>, data: Vec<u8>) {
core.handle_message(data).await;
}
async fn handle_connected_event(
state_manager: &Arc<ConnectionStateManager>,
event_handler: Option<Arc<dyn ClientEventHandler>>,
) {
debug!("[DefaultClientObserver] Connection established");
state_manager.set_connected();
if let Some(handler) = event_handler {
let event = ConnectionEvent::Connected;
crate::client::runtime::spawn_client_task(async move {
let _ = handler.handle_connection_event(&event).await;
});
}
}
async fn handle_disconnected_event(
state_manager: &Arc<ConnectionStateManager>,
event_handler: Option<Arc<dyn ClientEventHandler>>,
reason: String,
) {
debug!(
"[DefaultClientObserver] Connection disconnected: {}",
reason
);
state_manager.set_disconnected();
if let Some(handler) = event_handler {
let event = ConnectionEvent::Disconnected(reason);
crate::client::runtime::spawn_client_task(async move {
let _ = handler.handle_connection_event(&event).await;
});
}
}
async fn handle_error_event(
state_manager: &Arc<ConnectionStateManager>,
event_handler: Option<Arc<dyn ClientEventHandler>>,
error: crate::common::error::FlareError,
) {
error!("[DefaultClientObserver] Connection error: {:?}", error);
state_manager.set_failed();
if let Some(handler) = event_handler {
let event = ConnectionEvent::Error(error);
crate::client::runtime::spawn_client_task(async move {
let _ = handler.handle_connection_event(&event).await;
});
}
}
}
impl ConnectionObserver for DefaultClientMessageObserver {
fn on_event(&self, event: &ConnectionEvent) {
match event {
ConnectionEvent::Message(data) => {
let core = Arc::clone(&self.core);
let data_clone = data.clone();
crate::client::runtime::spawn_client_task(async move {
Self::handle_message_event(&core, data_clone).await;
});
}
ConnectionEvent::Connected => {
let state_manager = Arc::clone(&self.state_manager);
let event_handler = self.event_handler.clone();
crate::client::runtime::spawn_client_task(async move {
Self::handle_connected_event(&state_manager, event_handler).await;
});
}
ConnectionEvent::Disconnected(reason) => {
let state_manager = Arc::clone(&self.state_manager);
let event_handler = self.event_handler.clone();
let reason_clone = reason.clone();
crate::client::runtime::spawn_client_task(async move {
Self::handle_disconnected_event(&state_manager, event_handler, reason_clone)
.await;
});
}
ConnectionEvent::Error(e) => {
let state_manager = Arc::clone(&self.state_manager);
let event_handler = self.event_handler.clone();
let error_clone = e.clone();
crate::client::runtime::spawn_client_task(async move {
Self::handle_error_event(&state_manager, event_handler, error_clone).await;
});
}
}
}
}