use crate::{
command::TransportStats,
event::ServerEvent,
transport::{
config::TransportConfig,
connection_state::ConnectionStateManager,
lockfree::LockFreeHashMap,
session_actor::{
create_session_actor, SessionHandle, SessionHandler, DEFAULT_ACTOR_BUFFER_SIZE,
},
},
Packet, SessionId, TransportError,
};
use std::sync::Arc;
use tokio::sync::broadcast;
const LISTENER_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(200);
const DEFAULT_REQUEST_LIFECYCLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
pub struct TransportServer {
config: TransportConfig,
context: Arc<crate::transport::context::TransportContext>,
transports: Arc<LockFreeHashMap<SessionId, Arc<crate::transport::transport::Transport>>>,
session_handles: Arc<LockFreeHashMap<SessionId, SessionHandle>>,
session_id_generator: Arc<std::sync::atomic::AtomicU64>,
stats: Arc<LockFreeHashMap<SessionId, TransportStats>>,
event_sender: broadcast::Sender<ServerEvent>,
is_running: Arc<std::sync::atomic::AtomicBool>,
protocol_configs:
std::collections::HashMap<String, Box<dyn crate::protocol::adapter::DynServerConfig>>,
state_manager: ConnectionStateManager,
request_tracker: Arc<crate::transport::transport::RequestTracker>,
request_registry: Arc<crate::transport::request_registry::RequestRegistry>,
message_id_counter: std::sync::atomic::AtomicU32,
session_handler: Option<Arc<dyn SessionHandler>>,
actor_buffer_size: usize,
frame_policy: crate::packet::FramePolicy,
}
impl TransportServer {
pub async fn new(config: TransportConfig) -> Result<Self, TransportError> {
let ctx = Arc::new(crate::transport::context::TransportContext::new().await?);
let (event_sender, _) = broadcast::channel(8192);
Ok(Self {
config,
context: ctx,
transports: Arc::new(LockFreeHashMap::new()),
session_handles: Arc::new(LockFreeHashMap::new()),
session_id_generator: Arc::new(std::sync::atomic::AtomicU64::new(1)),
stats: Arc::new(LockFreeHashMap::new()),
event_sender,
is_running: Arc::new(std::sync::atomic::AtomicBool::new(false)),
protocol_configs: std::collections::HashMap::new(),
state_manager: ConnectionStateManager::new(),
request_tracker: Arc::new(
crate::transport::transport::RequestTracker::new_with_start_id(10000),
),
request_registry: Arc::new(crate::transport::request_registry::RequestRegistry::new()),
message_id_counter: std::sync::atomic::AtomicU32::new(20000),
session_handler: None,
actor_buffer_size: DEFAULT_ACTOR_BUFFER_SIZE,
frame_policy: crate::packet::FramePolicy::Lenient,
})
}
pub async fn new_with_protocols(
config: TransportConfig,
protocol_configs: std::collections::HashMap<
String,
Box<dyn crate::protocol::adapter::DynServerConfig>,
>,
) -> Result<Self, TransportError> {
let ctx = Arc::new(crate::transport::context::TransportContext::new().await?);
let (event_sender, _) = broadcast::channel(8192);
Ok(Self {
config,
context: ctx,
transports: Arc::new(LockFreeHashMap::new()),
session_handles: Arc::new(LockFreeHashMap::new()),
session_id_generator: Arc::new(std::sync::atomic::AtomicU64::new(1)),
stats: Arc::new(LockFreeHashMap::new()),
event_sender,
is_running: Arc::new(std::sync::atomic::AtomicBool::new(false)),
protocol_configs,
state_manager: ConnectionStateManager::new(),
request_tracker: Arc::new(
crate::transport::transport::RequestTracker::new_with_start_id(10000),
),
request_registry: Arc::new(crate::transport::request_registry::RequestRegistry::new()),
message_id_counter: std::sync::atomic::AtomicU32::new(20000),
session_handler: None,
actor_buffer_size: DEFAULT_ACTOR_BUFFER_SIZE,
frame_policy: crate::packet::FramePolicy::Lenient,
})
}
pub async fn new_with_protocols_and_handler(
config: TransportConfig,
protocol_configs: std::collections::HashMap<
String,
Box<dyn crate::protocol::adapter::DynServerConfig>,
>,
handler: Arc<dyn SessionHandler>,
buffer_size: Option<usize>,
) -> Result<Self, TransportError> {
let ctx = Arc::new(crate::transport::context::TransportContext::new().await?);
let (event_sender, _) = broadcast::channel(64);
Ok(Self {
config,
context: ctx,
transports: Arc::new(LockFreeHashMap::new()),
session_handles: Arc::new(LockFreeHashMap::new()),
session_id_generator: Arc::new(std::sync::atomic::AtomicU64::new(1)),
stats: Arc::new(LockFreeHashMap::new()),
event_sender,
is_running: Arc::new(std::sync::atomic::AtomicBool::new(false)),
protocol_configs,
state_manager: ConnectionStateManager::new(),
request_tracker: Arc::new(
crate::transport::transport::RequestTracker::new_with_start_id(10000),
),
request_registry: Arc::new(crate::transport::request_registry::RequestRegistry::new()),
message_id_counter: std::sync::atomic::AtomicU32::new(20000),
session_handler: Some(handler),
actor_buffer_size: buffer_size.unwrap_or(DEFAULT_ACTOR_BUFFER_SIZE),
frame_policy: crate::packet::FramePolicy::Lenient,
})
}
pub fn is_actor_mode(&self) -> bool {
self.session_handler.is_some()
}
pub(crate) fn with_frame_policy(mut self, policy: crate::packet::FramePolicy) -> Self {
self.frame_policy = policy;
self
}
pub async fn send_to_session(
&self,
session_id: SessionId,
packet: Packet,
) -> Result<(), TransportError> {
tracing::debug!(
"[SEND] TransportServer sending packet to session {} (ID: {}, size: {} bytes)",
session_id,
packet.header.message_id,
packet.payload.len()
);
if let Some(handle) = self.session_handles.get(&session_id) {
match handle.send_packet_with_reply(packet).await {
Ok(()) => {
tracing::debug!(
"[SUCCESS] Session {} send successful (via actor)",
session_id
);
return Ok(());
}
Err(e) => {
let error_msg = format!("{:?}", e);
if error_msg.contains("Broken pipe")
|| error_msg.contains("Connection reset")
|| error_msg.contains("Connection closed")
|| error_msg.contains("ECONNRESET")
|| error_msg.contains("EPIPE")
|| error_msg.contains("Actor channel closed")
|| error_msg.contains("Actor dropped")
{
tracing::warn!(
"[WARN] Session {} connection closed: {}",
session_id,
error_msg
);
let _ = self.remove_session(session_id).await;
return Err(TransportError::connection_error(
"Connection closed during send",
false,
));
} else {
tracing::error!("[ERROR] Session {} send failed: {:?}", session_id, e);
return Err(e);
}
}
}
}
if let Some(transport) = self.transports.get(&session_id) {
if !transport.is_connected().await {
tracing::warn!(
"[WARN] Session {} connection closed, skipping send",
session_id
);
let _ = self.remove_session(session_id).await;
return Err(TransportError::connection_error("Connection closed", false));
}
match transport.send(packet).await {
Ok(()) => {
tracing::debug!(
"[SUCCESS] Session {} send successful (TransportServer layer confirmation)",
session_id
);
Ok(())
}
Err(e) => {
tracing::error!("[ERROR] Session {} send failed: {:?}", session_id, e);
let error_msg = format!("{:?}", e);
if error_msg.contains("Broken pipe")
|| error_msg.contains("Connection reset")
|| error_msg.contains("Connection closed")
|| error_msg.contains("ECONNRESET")
|| error_msg.contains("EPIPE")
{
tracing::warn!(
"[WARN] Session {} connection closed: {}",
session_id,
error_msg
);
let _ = self.remove_session(session_id).await;
Err(TransportError::connection_error(
"Connection closed during send",
false,
))
} else {
tracing::error!(
"[ERROR] Session {} send failed (non-connection error): {:?}",
session_id,
e
);
Err(e)
}
}
}
} else {
tracing::warn!(
"[WARN] Session {} does not exist in connection mapping",
session_id
);
Err(TransportError::connection_error("Session not found", false))
}
}
pub async fn request_to_session(
&self,
session_id: SessionId,
packet: Packet,
) -> Result<Packet, TransportError> {
tracing::debug!(
"[REQUEST] TransportServer sending request to session {} (ID: {})",
session_id,
packet.header.message_id
);
if self.transports.get(&session_id).is_none() {
tracing::warn!(
"[WARN] Session {} does not exist in connection mapping",
session_id
);
return Err(TransportError::connection_error("Session not found", false));
}
if packet.header.packet_type != crate::packet::PacketType::Request {
return Err(TransportError::connection_error(
"Not a Request packet",
false,
));
}
let message_id = packet.header.message_id;
let (_, rx) = self
.request_tracker
.register_with_session(session_id, message_id);
if let Err(e) = self.send_to_session(session_id, packet).await {
self.request_tracker
.remove_with_session(session_id, message_id);
tracing::error!(
"[ERROR] Session {} request send failed: {:?}",
session_id,
e
);
return Err(e);
}
match tokio::time::timeout(std::time::Duration::from_secs(10), rx).await {
Ok(Ok(response)) => {
tracing::debug!(
"[SUCCESS] Session {} received response (response ID: {})",
session_id,
response.header.message_id
);
Ok(response)
}
Ok(Err(_)) => {
self.request_tracker
.remove_with_session(session_id, message_id);
Err(TransportError::connection_error("Connection closed", true))
}
Err(_) => {
self.request_tracker
.remove_with_session(session_id, message_id);
Err(TransportError::timeout_error(
"server request",
std::time::Duration::from_secs(10),
))
}
}
}
pub async fn send(
&self,
session_id: SessionId,
data: &[u8],
) -> Result<crate::event::TransportResult, TransportError> {
let message_id = self
.message_id_counter
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let packet = crate::packet::Packet::one_way(message_id, data.to_vec());
tracing::debug!(
"TransportServer sending data to session {}: {} bytes (ID: {})",
session_id,
data.len(),
message_id
);
match self.send_to_session(session_id, packet).await {
Ok(()) => {
Ok(crate::event::TransportResult::new_sent(
Some(session_id),
message_id,
))
}
Err(e) => Err(e),
}
}
pub async fn request(
&self,
session_id: SessionId,
data: &[u8],
) -> Result<crate::event::TransportResult, TransportError> {
let message_id = self
.message_id_counter
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let packet = crate::packet::Packet::request(message_id, data.to_vec());
tracing::debug!(
"TransportServer sending request to session {}: {} bytes (ID: {})",
session_id,
data.len(),
message_id
);
match self.request_to_session(session_id, packet).await {
Ok(response_packet) => {
tracing::debug!(
"TransportServer received response from session {}: {} bytes (ID: {})",
session_id,
response_packet.payload.len(),
response_packet.header.message_id
);
Ok(crate::event::TransportResult::new_completed(
Some(session_id),
message_id,
response_packet.payload.clone(),
))
}
Err(e) => {
if matches!(e, TransportError::Timeout { .. }) {
Ok(crate::event::TransportResult::new_timeout(
Some(session_id),
message_id,
))
} else {
Err(e)
}
}
}
}
pub async fn add_session(&self, connection: Box<dyn crate::Connection>) -> SessionId {
let session_id = connection.session_id();
let mut connection = connection;
connection.set_frame_policy(self.frame_policy);
let transport = Arc::new(crate::transport::transport::Transport::with_context(
self.config.clone(),
&self.context,
));
let event_receiver_opt = connection.event_stream();
transport
.set_connection_no_consumer(connection, session_id)
.await;
if let Err(e) = self.transports.insert(session_id, transport.clone()) {
tracing::error!(
"[ERROR] Failed to register transport for session {}: {:?}",
session_id,
e
);
let _ = transport.disconnect().await;
return session_id;
}
if let Err(e) = self.stats.insert(session_id, TransportStats::new()) {
tracing::error!(
"[ERROR] Failed to register stats for session {}: {:?}",
session_id,
e
);
let _ = self.transports.remove(&session_id);
let _ = transport.disconnect().await;
return session_id;
}
self.state_manager.add_connection(session_id);
let actor_handle = if let Some(handler) = &self.session_handler {
let (handle, actor) = create_session_actor(
session_id,
transport.clone(),
handler.clone(),
self.actor_buffer_size,
);
let actor = actor.with_inbound_registry(Some(self.request_registry.clone()));
match self.session_handles.insert(session_id, handle.clone()) {
Ok(_) => {
tokio::spawn(actor.run());
Some(handle)
}
Err(e) => {
tracing::error!(
"[ERROR] Failed to register session actor handle for {}: {:?}",
session_id,
e
);
None
}
}
} else {
None
};
if let Some(mut event_receiver) = event_receiver_opt {
let server_clone = self.clone();
tokio::spawn(async move {
tracing::info!("[LISTENER] TransportServer starting event consumption loop for session {} (pre-subscribed)", session_id);
while let Ok(transport_event) = event_receiver.recv().await {
tracing::trace!(
"[EVENT] TransportServer received event from session {}: {:?}",
session_id,
transport_event
);
if let Some(handle) = &actor_handle {
if let crate::event::TransportEvent::MessageReceived(packet) =
&transport_event
{
if packet.header.packet_type == crate::packet::PacketType::Response
&& server_clone.request_tracker.complete_with_session(
session_id,
packet.header.message_id,
packet.clone(),
)
{
continue;
}
if packet.header.packet_type == crate::packet::PacketType::Request {
server_clone.request_registry.register(
packet.header.message_id,
Some(session_id),
packet.header.biz_type,
DEFAULT_REQUEST_LIFECYCLE_TIMEOUT,
);
}
}
if let Err(e) = handle.send_event(transport_event).await {
tracing::warn!(
"[WARN] Failed to forward event to actor for session {}: {:?}",
session_id,
e
);
break;
}
} else {
server_clone
.handle_transport_event(session_id, transport_event)
.await;
}
}
let closed_inbound = server_clone
.request_registry
.close_session_pending(session_id);
let failed_outbound = server_clone.request_tracker.fail_session(Some(session_id));
if closed_inbound > 0 || failed_outbound > 0 {
tracing::debug!(
"[END] Session {} loop ended: closed {} inbound, failed {} outbound pending",
session_id,
closed_inbound,
failed_outbound
);
}
tracing::info!(
"[END] TransportServer event consumption loop ended for session {}",
session_id
);
});
} else {
tracing::warn!(
"[WARN] Session {} unable to get event stream before connection setup",
session_id
);
}
tracing::info!(
"[SUCCESS] TransportServer added session: {} (using Transport abstraction)",
session_id
);
session_id
}
pub async fn remove_session(&self, session_id: SessionId) -> Result<(), TransportError> {
let closed_pending = self.request_registry.close_session_pending(session_id);
if closed_pending > 0 {
tracing::debug!(
"[REQUEST] Session {} removed, marked {} pending requests as SessionClosed",
session_id,
closed_pending
);
}
let failed_pending = self.request_tracker.fail_session(Some(session_id));
if failed_pending > 0 {
tracing::debug!(
"[REQUEST] Session {} removed, failed {} server-side pending requests",
session_id,
failed_pending
);
}
if let Err(e) = self.transports.remove(&session_id) {
tracing::warn!(
"[WARN] Failed to remove transport for session {}: {:?}",
session_id,
e
);
}
if let Err(e) = self.session_handles.remove(&session_id) {
tracing::warn!(
"[WARN] Failed to remove session actor handle for session {}: {:?}",
session_id,
e
);
}
if let Err(e) = self.stats.remove(&session_id) {
tracing::warn!(
"[WARN] Failed to remove stats for session {}: {:?}",
session_id,
e
);
}
self.state_manager.remove_connection(session_id);
tracing::info!("[REMOVE] TransportServer removed session: {}", session_id);
Ok(())
}
pub async fn close_session(&self, session_id: SessionId) -> Result<(), TransportError> {
if !self.state_manager.try_start_closing(session_id).await {
tracing::debug!(
"Session {} already closing or closed, skipping close logic",
session_id
);
return Ok(());
}
tracing::info!(
"[CLOSE] Starting graceful close for session: {}",
session_id
);
let close_event = ServerEvent::ConnectionClosed {
session_id,
reason: crate::error::CloseReason::Normal,
};
let _ = self.event_sender.send(close_event);
self.do_close_session(session_id).await?;
self.state_manager.mark_closed(session_id).await;
self.remove_session(session_id).await?;
tracing::info!("[SUCCESS] Session {} close completed", session_id);
Ok(())
}
pub async fn force_close_session(&self, session_id: SessionId) -> Result<(), TransportError> {
if !self.state_manager.try_start_closing(session_id).await {
tracing::debug!(
"Session {} already closing or closed, skipping force close",
session_id
);
return Ok(());
}
tracing::info!("[FORCE] Force closing session: {}", session_id);
let close_event = ServerEvent::ConnectionClosed {
session_id,
reason: crate::error::CloseReason::Forced,
};
let _ = self.event_sender.send(close_event);
if let Some(transport) = self.transports.get(&session_id) {
let _ = transport.disconnect().await; }
self.state_manager.mark_closed(session_id).await;
self.remove_session(session_id).await?;
tracing::info!("[SUCCESS] Session {} force close completed", session_id);
Ok(())
}
pub async fn close_all_sessions(&self) -> Result<(), TransportError> {
let session_ids = self.active_sessions().await;
let total_sessions = session_ids.len();
if total_sessions == 0 {
tracing::info!("No active sessions to close");
return Ok(());
}
tracing::info!(
"[BATCH] Starting batch close of {} sessions",
total_sessions
);
let start_time = std::time::Instant::now();
let timeout = self.config.graceful_timeout;
let mut success_count = 0;
let mut error_count = 0;
for session_id in session_ids {
if start_time.elapsed() >= timeout {
tracing::warn!(
"[WARN] Batch close timeout, remaining sessions will be force closed"
);
let _ = self.force_close_session(session_id).await;
continue;
}
match self.close_session(session_id).await {
Ok(_) => success_count += 1,
Err(e) => {
error_count += 1;
tracing::warn!("[WARN] Failed to close session {}: {:?}", session_id, e);
}
}
}
tracing::info!(
"[SUCCESS] Batch close completed, success: {}, failed: {}",
success_count,
error_count
);
Ok(())
}
async fn do_close_session(&self, session_id: SessionId) -> Result<(), TransportError> {
if let Some(transport) = self.transports.get(&session_id) {
match tokio::time::timeout(self.config.graceful_timeout, transport.disconnect()).await {
Ok(Ok(_)) => {
tracing::debug!("[SUCCESS] Session {} graceful close successful", session_id);
}
Ok(Err(e)) => {
tracing::warn!(
"[WARN] Session {} graceful close failed: {:?}",
session_id,
e
);
}
Err(_) => {
tracing::warn!("[WARN] Session {} graceful close timeout", session_id);
}
}
}
Ok(())
}
pub async fn should_ignore_messages(&self, session_id: SessionId) -> bool {
self.state_manager.should_ignore_messages(session_id).await
}
pub async fn broadcast(&self, packet: Packet) -> Result<(), TransportError> {
let mut success_count = 0;
let mut error_count = 0;
if self.is_actor_mode() {
let session_ids: Vec<SessionId> = self.session_handles.keys().unwrap_or_default();
for session_id in session_ids {
if let Some(handle) = self.session_handles.get(&session_id) {
match handle.send_packet(packet.clone()).await {
Ok(()) => success_count += 1,
Err(e) => {
error_count += 1;
tracing::warn!(
"[WARN] Broadcast to session {} failed: {:?}",
session_id,
e
);
}
}
}
}
} else {
let session_ids: Vec<SessionId> = self.transports.keys().unwrap_or_default();
for session_id in session_ids {
if let Some(transport) = self.transports.get(&session_id) {
match transport.send(packet.clone()).await {
Ok(()) => success_count += 1,
Err(e) => {
error_count += 1;
tracing::warn!(
"[WARN] Broadcast to session {} failed: {:?}",
session_id,
e
);
}
}
}
}
}
if error_count > 0 {
tracing::warn!(
"[WARN] Broadcast completed, success: {}, failed: {}",
success_count,
error_count
);
} else {
tracing::info!("[SUCCESS] Broadcast completed, success: {}", success_count);
}
Ok(())
}
pub async fn active_sessions(&self) -> Vec<SessionId> {
self.transports.keys().unwrap_or_default()
}
pub async fn session_count(&self) -> usize {
self.transports.len()
}
fn generate_session_id(&self) -> SessionId {
let id = self
.session_id_generator
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
SessionId(id)
}
pub fn subscribe_events(&self) -> tokio::sync::broadcast::Receiver<crate::event::ServerEvent> {
self.event_sender.subscribe()
}
pub async fn serve(&self) -> Result<(), TransportError> {
self.is_running
.store(true, std::sync::atomic::Ordering::SeqCst);
if self.protocol_configs.is_empty() {
tracing::warn!("[WARN] No protocols configured, server cannot start listening");
self.is_running
.store(false, std::sync::atomic::Ordering::SeqCst);
return Err(TransportError::config_error(
"protocols",
"No protocols configured",
));
}
tracing::info!(
"[START] Starting {} protocol servers",
self.protocol_configs.len()
);
let mut listen_tasks = Vec::new();
for (protocol_name, protocol_config) in &self.protocol_configs {
tracing::info!("[CONFIG] Processing protocol: {}", protocol_name);
let address = self.get_protocol_bind_address(protocol_config);
tracing::info!(
"[BIND] Protocol {} bind address: {}",
protocol_name,
address
);
match protocol_config.build_server_dyn().await {
Ok(server) => {
match self
.start_protocol_listener(server, protocol_name.clone())
.await
{
Ok(listener_task) => {
listen_tasks.push(listener_task);
tracing::info!(
"[SUCCESS] {} server started successfully: {}",
protocol_name,
address
);
}
Err(e) => {
tracing::error!(
"[ERROR] {} listener task creation failed: {:?}",
protocol_name,
e
);
self.is_running
.store(false, std::sync::atomic::Ordering::SeqCst);
for task in listen_tasks {
task.abort();
}
return Err(e);
}
}
}
Err(e) => {
tracing::error!("[ERROR] {} server build failed: {:?}", protocol_name, e);
self.is_running
.store(false, std::sync::atomic::Ordering::SeqCst);
for task in listen_tasks {
task.abort();
}
return Err(e);
}
}
}
listen_tasks.push(self.start_request_timeout_scanner());
tracing::info!("[TARGET] All protocol servers started, waiting for connections...");
for (index, task) in listen_tasks.into_iter().enumerate() {
tracing::info!("[WAIT] Waiting for task {} to complete...", index + 1);
if let Err(e) = task.await {
tracing::error!("[ERROR] Task {} was cancelled: {:?}", index + 1, e);
return Err(TransportError::config_error(
"server",
"Listener task cancelled",
));
}
}
tracing::info!("[STOP] TransportServer stopped");
Ok(())
}
async fn start_protocol_listener(
&self,
mut server: Box<dyn crate::Server>,
protocol_name: String,
) -> Result<tokio::task::JoinHandle<()>, TransportError> {
let server_clone = self.clone();
let task = tokio::spawn(async move {
tracing::info!("[START] {} listener task started", protocol_name);
let mut accept_count = 0u64;
loop {
if !server_clone
.is_running
.load(std::sync::atomic::Ordering::SeqCst)
{
tracing::info!("[STOP] {} listener received stop signal", protocol_name);
break;
}
tracing::debug!(
"[LOOP] {} waiting for connections... (accept count: {})",
protocol_name,
accept_count
);
match tokio::time::timeout(LISTENER_POLL_INTERVAL, server.accept()).await {
Err(_) => {
continue;
}
Ok(accept_result) => match accept_result {
Ok(mut connection) => {
accept_count += 1;
tracing::info!(
"[SUCCESS] {} accept successful! Connection #{}",
protocol_name,
accept_count
);
let connection_info = connection.connection_info();
let peer_addr = connection_info.peer_addr;
tracing::info!(
"[CONNECT] New {} connection #{}: {}",
protocol_name,
accept_count,
peer_addr
);
let session_id = server_clone.generate_session_id();
connection.set_session_id(session_id);
tracing::info!(
"[ID] Generated session ID for {} connection: {}",
protocol_name,
session_id
);
let actual_session_id = server_clone.add_session(connection).await;
let connect_event = ServerEvent::ConnectionEstablished {
session_id: actual_session_id,
info: connection_info,
};
let _ = server_clone.event_sender.send(connect_event);
tracing::info!("[EVENT] {} connection event sent", protocol_name);
}
Err(e) => {
if !server_clone
.is_running
.load(std::sync::atomic::Ordering::SeqCst)
{
tracing::info!(
"[STOP] {} listener stopping after accept exit: {:?}",
protocol_name,
e
);
break;
}
tracing::warn!(
"[WARN] {} accept connection failed, continue listening: {:?}",
protocol_name,
e
);
tokio::time::sleep(LISTENER_POLL_INTERVAL).await;
continue;
}
},
}
}
if let Err(e) = server.shutdown().await {
tracing::warn!("[WARN] {} server shutdown failed: {:?}", protocol_name, e);
}
tracing::info!("[STOP] {} server stopped", protocol_name);
});
Ok(task)
}
fn start_request_timeout_scanner(&self) -> tokio::task::JoinHandle<()> {
let server_clone = self.clone();
tokio::spawn(async move {
let tick = server_clone.request_registry.tick_duration();
tracing::info!(
"[START] request timeout scanner started (tick={}ms)",
tick.as_millis()
);
loop {
if !server_clone
.is_running
.load(std::sync::atomic::Ordering::SeqCst)
{
break;
}
tokio::time::sleep(tick).await;
let timed_out = server_clone.request_registry.scan_timeout_bucket();
if timed_out > 0 {
tracing::warn!(
"[TIMEOUT] request timeout scanner marked {} requests as TimedOut",
timed_out
);
}
}
tracing::info!("[STOP] request timeout scanner stopped");
})
}
fn get_protocol_bind_address(
&self,
protocol_config: &Box<dyn crate::protocol::adapter::DynServerConfig>,
) -> std::net::SocketAddr {
protocol_config.get_bind_address()
}
async fn handle_transport_event(
&self,
session_id: SessionId,
transport_event: crate::event::TransportEvent,
) {
tracing::debug!(
"[HANDLER] handle_transport_event: session {}, event: {:?}",
session_id,
transport_event
);
match transport_event {
crate::event::TransportEvent::MessageReceived(packet) => {
tracing::debug!("[RECV] Received MessageReceived event, packet type: {:?}, ID: {}, biz_type: {}", packet.header.packet_type, packet.header.message_id, packet.header.biz_type);
match packet.header.packet_type {
crate::packet::PacketType::Request => {
tracing::debug!(
"[REQUEST] Handling Request type packet (ID: {}, biz_type: {})",
packet.header.message_id,
packet.header.biz_type
);
let registered = self.request_registry.register(
packet.header.message_id,
Some(session_id),
packet.header.biz_type,
DEFAULT_REQUEST_LIFECYCLE_TIMEOUT,
);
if !registered {
tracing::warn!(
"[REQUEST] Duplicate request registration detected: session={}, request_id={}",
session_id,
packet.header.message_id
);
}
let server_clone = self.clone();
let request_registry = self.request_registry.clone();
let mut context = crate::event::TransportContext::new_request_with_registry(
Some(session_id),
packet.header.message_id,
packet.header.biz_type,
if packet.ext_header.is_empty() {
None
} else {
Some(packet.ext_header.clone())
},
packet.payload.clone(),
std::sync::Arc::new(
move |response_data: Vec<u8>| -> futures::future::BoxFuture<
'static,
Result<(), crate::TransportError>,
> {
let server = server_clone.clone();
let request_registry = request_registry.clone();
let request_message_id = packet.header.message_id;
let request_biz_type = packet.header.biz_type;
Box::pin(async move {
tracing::debug!("[RESPOND] Creating response packet: request_id={}, request_biz_type={}, response_size={} bytes", request_message_id, request_biz_type, response_data.len());
let response_packet = crate::packet::Packet {
header: crate::packet::FixedHeader {
version: 1,
compression: crate::packet::CompressionType::None,
packet_type: crate::packet::PacketType::Response,
biz_type: 0,
message_id: request_message_id,
ext_header_len: 0,
payload_len: response_data.len() as u32,
reserved: crate::packet::ReservedFlags::new(),
},
ext_header: Vec::new(),
payload: response_data,
};
match server
.send_to_session(session_id, response_packet)
.await
{
Ok(_) => {
tracing::debug!("[SUCCESS] Response sent successfully: session={}, request_id={}", session_id, request_message_id);
Ok(())
}
Err(e) => {
request_registry.record_response_send_failed();
tracing::error!("[ERROR] Failed to send response: session={}, request_id={}, error={:?}", session_id, request_message_id, e);
Err(e)
}
}
})
},
),
Some(self.request_registry.clone()),
);
context.set_primary();
tracing::debug!(
"[SUCCESS] Set TransportContext as primary instance (ID: {})",
packet.header.message_id
);
let event = crate::event::ServerEvent::MessageReceived {
session_id,
context,
};
tracing::debug!("[SEND] Preparing to send ServerEvent::MessageReceived (session: {}, ID: {}, biz_type: {})", session_id, packet.header.message_id, packet.header.biz_type);
match self.event_sender.send(event) {
Ok(receivers) => {
tracing::debug!("[SUCCESS] ServerEvent sent successfully, receiver count: {} (session: {}, ID: {})", receivers, session_id, packet.header.message_id);
}
Err(e) => {
tracing::error!(
"[ERROR] ServerEvent send failed: {:?} (session: {}, ID: {})",
e,
session_id,
packet.header.message_id
);
}
}
}
crate::packet::PacketType::Response => {
let message_id = packet.header.message_id;
tracing::debug!("[RECV] TransportServer received response packet from session {} (ID: {})", session_id, message_id);
if self.request_tracker.complete_with_session(
session_id,
message_id,
packet.clone(),
) {
tracing::debug!(
"[SUCCESS] Successfully completed server request (ID: {})",
message_id
);
} else {
tracing::warn!("[WARN] Received unknown response packet (ID: {}), may be timeout or duplicate response", message_id);
let context = crate::event::TransportContext::new_oneway(
Some(session_id),
packet.header.message_id,
packet.header.biz_type,
if packet.ext_header.is_empty() {
None
} else {
Some(packet.ext_header.clone())
},
packet.payload.clone(),
);
let event = crate::event::ServerEvent::MessageReceived {
session_id,
context,
};
match self.event_sender.send(event) {
Ok(receivers) => {
tracing::debug!("[SUCCESS] Unknown response packet sent as normal message successfully, receiver count: {}", receivers);
}
Err(e) => {
tracing::error!("[ERROR] Unknown response packet as normal message send failed: {:?}", e);
}
}
}
}
_ => {
tracing::debug!(
"[PACKET] Processing other type packet (type: {:?}, ID: {})",
packet.header.packet_type,
packet.header.message_id
);
let context = crate::event::TransportContext::new_oneway(
Some(session_id),
packet.header.message_id,
packet.header.biz_type,
if packet.ext_header.is_empty() {
None
} else {
Some(packet.ext_header.clone())
},
packet.payload.clone(),
);
let event = crate::event::ServerEvent::MessageReceived {
session_id,
context,
};
match self.event_sender.send(event) {
Ok(receivers) => {
tracing::debug!("[SUCCESS] Other type packet sent successfully, receiver count: {}", receivers);
}
Err(e) => {
tracing::error!("[ERROR] Other type packet send failed: {:?}", e);
}
}
}
}
}
crate::event::TransportEvent::MessageSent { packet_id } => {
tracing::debug!("[SEND] Processing MessageSent event (ID: {})", packet_id);
let event = crate::event::ServerEvent::MessageSent {
session_id,
message_id: packet_id,
};
match self.event_sender.send(event) {
Ok(receivers) => {
tracing::debug!(
"[SUCCESS] MessageSent event sent successfully, receiver count: {}",
receivers
);
}
Err(e) => {
tracing::error!("[ERROR] MessageSent event send failed: {:?}", e);
}
}
}
crate::event::TransportEvent::ConnectionClosed { reason } => {
tracing::debug!(
"[CLOSE] Processing ConnectionClosed event, reason: {:?}",
reason
);
let event = crate::event::ServerEvent::ConnectionClosed { session_id, reason };
match self.event_sender.send(event) {
Ok(receivers) => {
tracing::debug!("[SUCCESS] ConnectionClosed event sent successfully, receiver count: {}", receivers);
}
Err(e) => {
tracing::error!("[ERROR] ConnectionClosed event send failed: {:?}", e);
}
}
let _ = self.remove_session(session_id).await;
}
crate::event::TransportEvent::TransportError { error } => {
tracing::debug!("[ERROR] Processing TransportError event: {:?}", error);
let event = crate::event::ServerEvent::TransportError {
session_id: Some(session_id),
error,
};
match self.event_sender.send(event) {
Ok(receivers) => {
tracing::debug!(
"[SUCCESS] TransportError event sent successfully, receiver count: {}",
receivers
);
}
Err(e) => {
tracing::error!("[ERROR] TransportError event send failed: {:?}", e);
}
}
}
_ => {
tracing::trace!(
"[IGNORE] TransportServer ignoring event: {:?}",
transport_event
);
}
}
}
pub async fn stop(&self) {
tracing::info!("[STOP] Stopping TransportServer");
self.is_running
.store(false, std::sync::atomic::Ordering::SeqCst);
}
}
impl Clone for TransportServer {
fn clone(&self) -> Self {
let mut cloned_configs = std::collections::HashMap::new();
for (name, config) in &self.protocol_configs {
cloned_configs.insert(name.clone(), config.clone_server_dyn());
}
Self {
config: self.config.clone(),
context: self.context.clone(),
transports: self.transports.clone(),
session_id_generator: self.session_id_generator.clone(),
stats: self.stats.clone(),
event_sender: self.event_sender.clone(),
is_running: self.is_running.clone(),
protocol_configs: cloned_configs,
state_manager: self.state_manager.clone(),
request_tracker: self.request_tracker.clone(),
request_registry: self.request_registry.clone(),
message_id_counter: std::sync::atomic::AtomicU32::new(20000),
session_handles: self.session_handles.clone(),
session_handler: self.session_handler.clone(),
actor_buffer_size: self.actor_buffer_size,
frame_policy: self.frame_policy,
}
}
}
impl std::fmt::Debug for TransportServer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TransportServer")
.field("session_count", &self.transports.len())
.field("config", &self.config)
.finish()
}
}