use tokio::net::{TcpListener, TcpStream};
use crate::{
net_params::{ITcpConnectionParams, SConnectionParams, STcpServerParams},
smart_pub::{ArcMutex, StdMutex},
};
use super::{
net_pub::{INetwork, ServiceType},
smart_pub::{StdArc, StdResult, WatchReceiver, WatchSender},
tcp_connection,
};
#[derive(Debug)]
pub struct CTcpServer {
exit_event_sender: WatchSender<bool>,
exit_event_receiver: WatchReceiver<bool>,
params: STcpServerParams,
}
impl CTcpServer {
pub fn new(params: STcpServerParams) -> Self {
let (event_sender, event_receiver) = tokio::sync::watch::channel(false);
CTcpServer {
exit_event_sender: event_sender,
exit_event_receiver: event_receiver,
params,
}
}
#[tracing::instrument(name = "TCP Server Start", skip(self)/* , fields(cleint_address=tracing::field::Empty) */)]
async fn connect_new(
self: StdArc<Self>,
stream: TcpStream,
) -> StdResult<()> {
let remote_addr = stream.peer_addr().unwrap();
let connection_params = SConnectionParams {
call_back: self.params.call_back_inner(),
serv_type: ServiceType::Server,
extra_params: remote_addr.to_string()
};
let connection_state: ArcMutex<bool> = StdArc::new(StdMutex::new(true));
tracing::info!("New Client");
let _conn = tcp_connection::CTcpConnection::new(
stream,
connection_params,
self.exit_event_receiver.clone(),
connection_state,
);
tracing::info!("New Client OK");
Ok(())
}
}
impl INetwork for CTcpServer {
#[tracing::instrument(name = "TCP Server Start", skip(self)/* , fields(address=%self.params.handle_ref()) */)]
async fn start(self: StdArc<Self>) -> StdResult<()> {
tracing::info!("Start Server! {}", self.params.handle_ref());
let listener = TcpListener::bind(self.params.address_inner()).await?;
let mut exit_receiver = self.clone().exit_event_receiver.clone();
loop {
tokio::select! {
flag = exit_receiver.changed() =>
{
match flag {
Ok(_exit_flag) => { break;}
Err(_) => {}
}
}
ret = listener.accept() =>
{
let (stream, _addr) = ret?;
tokio::spawn(self.clone().connect_new(stream));
}
};
}
tracing::info!("Stop Server!");
Ok(())
}
#[tracing::instrument(name = "TCP Server Stop", skip(self), fields(address=%self.params.handle_ref()))]
async fn stop(self: StdArc<Self>) -> StdResult<()> {
let r = self.exit_event_sender.send(true);
match r {
Ok(()) => Ok(()),
Err(e) => Err(Box::new(e)),
}
}
}