use crate::net_params::{IUdpConnectionParams, SUdpConnectionParams};
use crate::net_pub::{spawn_and_log_error, EConnectionState, ServiceType};
use crate::smart_pub::{MpscAsyncReceiver, StdResult, WatchReceiver};
use tokio::net::UdpSocket;
#[derive(Debug)]
pub struct CUdpConnection {}
impl CUdpConnection {
#[tracing::instrument(
name="Udp Connection new",
skip(socket, params, exit_receiver),
fields(handle=%params.handle_inner(), remote_address=%params.remote_address_ref(), state=tracing::field::Empty)
)]
pub fn new(
socket: UdpSocket,
params: SUdpConnectionParams,
exit_receiver: WatchReceiver<bool>,
) -> Self {
let (data_reader, data_writer) = tokio::sync::mpsc::unbounded_channel::<(String, Vec<u8>)>();
let params2: SUdpConnectionParams = params.clone();
tracing::info!("Start spawn_read_write_loop.");
if params.service_type() == ServiceType::Client {
spawn_and_log_error(Self::spawn_read_write_loop_for_client(
params.clone(),
exit_receiver.clone(),
socket,
data_writer,
));
} else {
spawn_and_log_error(Self::spawn_read_write_loop_for_server(
params.clone(),
exit_receiver.clone(),
socket,
data_writer,
));
}
tracing::info!("spawn_read_write_loop ok");
tracing::info!("notify client state Running.");
params2.call_back_inner().state_callback(
params2.handle_inner(),
EConnectionState::Running(data_reader.clone()),
params2.serv_type,
);
CUdpConnection {}
}
async fn spawn_read_write_loop_for_client(
conn_params: SUdpConnectionParams,
mut exit_receiver: WatchReceiver<bool>,
socket: UdpSocket,
mut data_writer: MpscAsyncReceiver<(String, Vec<u8>)>,
) -> StdResult<()> {
let mut buf = [0u8; 1024];
loop {
tokio::select! {
flag = exit_receiver.changed() => {
match flag {
Ok(_exit_flag) => { break;}
Err(_) => {}
}
}
read_result = socket.recv(&mut buf) => {
match read_result
{
Err(e) => {
tracing::warn!("{} {} {:?}", conn_params.handle_ref(), " Exited ", e);
conn_params.call_back_inner().state_callback(conn_params.handle_inner(), EConnectionState::Interrupt, conn_params.serv_type);
break;
},
Ok(size) if size > 0 => {
tracing::debug!("Recved Data");
conn_params.call_back_inner().data_callback(buf.to_vec(), conn_params.handle_inner(), conn_params.serv_type);
}
Ok(size) if size == 0 => {
tracing::warn!("{}", " Recved Size = 0, Will exit.");
conn_params.call_back_inner().state_callback(conn_params.handle_inner(), EConnectionState::Interrupt, conn_params.serv_type);
break;
}
_ => unreachable!()
}
}
msg = data_writer.recv() => {
match msg {
Some((_addr, msg_data)) => {
let _ = socket.send(&msg_data).await;
}
None => {}
}
}
}
}
Ok(())
}
async fn spawn_read_write_loop_for_server(
conn_params: SUdpConnectionParams,
mut exit_receiver: WatchReceiver<bool>,
socket: UdpSocket,
mut data_writer: MpscAsyncReceiver<(String, Vec<u8>)>,
) -> StdResult<()> {
let mut buf = [0u8; 1024];
tracing::debug!("spawn_read_write_loop_for_server");
loop {
tokio::select! {
flag = exit_receiver.changed() => {
match flag {
Ok(_exit_flag) => { break;}
Err(_) => {}
}
}
read_result = socket.recv_from(&mut buf) => {
match read_result
{
Err(e) => {
tracing::warn!("{} {} {:?}", conn_params.handle_ref(), " Exited ", e);
conn_params.call_back_inner().state_callback(conn_params.handle_inner(), EConnectionState::Interrupt, conn_params.serv_type);
break;
},
Ok((size, addr)) if size > 0 => {
tracing::debug!("From {} {}", addr, " Recved Data");
conn_params.call_back_inner().data_callback(buf.to_vec(), addr.to_string(), conn_params.serv_type);
}
Ok((size, addr)) if size == 0 => {
tracing::warn!("From {} {}", addr, " Recved Size = 0, Will exit.");
conn_params.call_back_inner().state_callback(addr.to_string(), EConnectionState::Interrupt, conn_params.serv_type);
break;
}
_ => unreachable!()
}
}
msg = data_writer.recv() => {
match msg {
Some((addr, msg_data)) => {
let _ = socket.send_to(&msg_data, addr).await;
}
None => {}
}
}
}
}
Ok(())
}
}