use tokio::net::UdpSocket;
use crate::net_params::{IUdpConnectionParams, SUdpClientParams, SUdpConnectionParams};
use crate::net_pub::{EConnectionState, INetwork, ServiceType};
use crate::smart_pub::{StdArc, StdResult, ArcMutex, StdMutex, WatchReceiver, WatchSender};
use crate::udp_connection::CUdpConnection;
#[derive(Debug)]
pub struct CUdpClient {
exit_event_sender: WatchSender<bool>,
exit_event_receiver: WatchReceiver<bool>,
params: SUdpClientParams,
connection: ArcMutex<Option<CUdpConnection>>,
}
impl CUdpClient {
pub fn new(params: SUdpClientParams) -> Self {
let (exit_sender, exit_receiver) = tokio::sync::watch::channel(false);
CUdpClient {
exit_event_sender: exit_sender,
exit_event_receiver: exit_receiver,
params,
connection: StdArc::new(StdMutex::new(None)),
}
}
async fn connect_to(self: StdArc<Self>, stream: UdpSocket) -> StdResult<()> {
self.clone().execute_state(EConnectionState::Connecting);
tracing::info!("Start connect. ");
let connect_result = stream.connect(self.params.remote_address_inner()).await;
match connect_result {
Ok(_r) => {
let conn_params = SUdpConnectionParams {
call_back: self.params.call_back_inner(),
serv_type: ServiceType::Client,
extra_params: self.params.handle_inner(),
remote_addr: self.params.remote_address_inner(),
};
tracing::info!("Start new connection");
let conn =
CUdpConnection::new(stream, conn_params, self.exit_event_receiver.clone());
tracing::info!("Start connection OK");
let conn_arc = self.connection.clone();
let mut conn_mutex = conn_arc.lock().unwrap();
*conn_mutex = Some(conn);
}
Err(e) => {
tracing::warn!("{:?}", e);
}
}
Ok(())
}
fn execute_state(self: StdArc<Self>, state: EConnectionState<(String, Vec<u8>)>) {
let call_back_option = self.params.call_back_inner();
call_back_option.clone().state_callback(
self.params.handle_inner(),
state,
self.params.service_type(),
);
}
}
impl INetwork for CUdpClient {
#[tracing::instrument(name = "UDP Client Start!", skip(self), fields(
local_address = %self.params.local_address_ref(),
remote_address = %self.params.remote_address_ref()
))]
async fn start(self: StdArc<Self>) -> StdResult<()> {
tracing::info!("Start Udp Client");
tracing::info!("Start bind.");
let ret = tokio::net::UdpSocket::bind(self.params.local_address_inner()).await;
match ret {
Ok(stream) => {
let _ret = self.clone().connect_to(stream).await;
}
Err(e) => {
tracing::error!("Bind Failure {:?}", e);
}
}
tracing::info!("Stop Udp Client");
Ok(())
}
#[tracing::instrument(name = "UDP Client Stop!", skip(self), fields(
local_address = %self.params.local_address_inner(),
remote_address = %self.params.remote_address_inner()
))]
async fn stop(self: StdArc<Self>) -> StdResult<()> {
tracing::info!("UDP Client will stop");
let _ = self.exit_event_sender.send(true);
Ok(())
}
}