use crate::{
net_params::{ITcpConnectionParams, STcpClientParams, SConnectionParams},
net_pub::{EConnectionState, INetwork, ServiceType},
smart_pub::{ArcMutex, StdArc, StdMutex, StdResult, WatchReceiver, WatchSender},
tcp_connection::CTcpConnection,
};
#[derive(Debug)]
pub struct CTcpClient {
exit_event_sender: WatchSender<bool>,
exit_event_receiver: WatchReceiver<bool>,
params: STcpClientParams,
}
impl CTcpClient {
#[tracing::instrument(name = "TCP Client New")]
pub fn new(params: STcpClientParams) -> Self {
let (exit_sender, exit_receiver) = tokio::sync::watch::channel(false);
CTcpClient {
exit_event_sender: exit_sender,
exit_event_receiver: exit_receiver,
params,
}
}
async fn connect_to(self: StdArc<Self>) -> StdResult<()> {
let duration = tokio::time::Duration::from_secs(1);
self.clone().execute_state(EConnectionState::Connecting);
let client_connection_state: ArcMutex<bool> = StdArc::new(StdMutex::new(false)); loop {
match self.clone().exit_event_receiver.has_changed() {
Ok(state) if state => {
tracing::debug!("Exit Signal Recved");
break;
}
Err(e) => {
tracing::error!("{:?}", e);
}
_ => {}
}
let cnn_state = {
let state_arc = client_connection_state.clone();
let state_mutex = state_arc.lock().unwrap();
*state_mutex
};
if cnn_state {
tokio::time::sleep(duration).await;
continue;
}
let ret = tokio::net::TcpStream::connect(self.params.address_inner()).await;
match ret {
Ok(stream) => {
let conn_params = SConnectionParams {
call_back: self.params.call_back_inner(),
serv_type: ServiceType::Client,
extra_params: self.params.handle_inner(),
};
tracing::info!("Start new connection");
let _conn =
CTcpConnection::new(stream, conn_params, self.exit_event_receiver.clone(), client_connection_state.clone());
tracing::info!("Start connection OK");
{
let state_arc = client_connection_state.clone();
let mut state_mutex = state_arc.lock().unwrap();
*state_mutex = true;
};
}
Err(e) => {
tracing::error!(
"Connection Failure [{:?}], Try reconnecting after 1 second",
e
);
tokio::time::sleep(duration).await;
}
}
}
self.clone().execute_state(EConnectionState::Stoped);
Ok(())
}
fn execute_state(self: StdArc<Self>, state: EConnectionState<Vec<u8>>) {
let call_back_option = self.params.call_back_inner();
call_back_option.clone().state_callback(
self.params.handle_inner(),
state,
ServiceType::Client,
);
}
}
impl INetwork for CTcpClient {
#[tracing::instrument(name = "TCP Client Start!", skip(self), fields(
address = %self.params.address_ref()
))]
async fn start(self: StdArc<Self>) -> StdResult<()> {
tracing::info!("Start Tcp Client");
loop {
match self.clone().exit_event_receiver.has_changed() {
Ok(state) if state => {
tracing::debug!("Exited State Recved");
break;
}
Err(e) => {
tracing::error!("{:?}", e);
}
_ => {}
}
let ret = self.clone().connect_to().await;
match ret {
_ => {}
}
}
tracing::info!("Stop Tcp Client");
Ok(())
}
#[tracing::instrument(name = "TCP Client Stop!", skip(self), fields(
address = %self.params.address_ref()
))]
async fn stop(self: StdArc<Self>) -> StdResult<()> {
tracing::info!("TCP Client will stop");
let _ = self.exit_event_sender.send(true);
Ok(())
}
}