use crate::net_params::{ITcpConnectionParams, SConnectionParams};
use crate::net_pub::{spawn_and_log_error, EConnectionState};
use crate::smart_pub::{MpscAsyncReceiver, StdResult, WatchReceiver, ArcMutex};
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
net::{tcp::OwnedReadHalf, tcp::OwnedWriteHalf, TcpStream},
};
#[derive(Debug)]
pub struct CTcpConnection {
}
impl CTcpConnection {
#[tracing::instrument(name="TCP Connection New", skip(socket, params, exit_receiver, connection_state), fields(handle=%params.handle_inner(), state=tracing::field::Empty))]
pub fn new(
socket: TcpStream,
params: SConnectionParams,
exit_receiver: WatchReceiver<bool>,
connection_state: ArcMutex<bool>,
) -> Self {
let (socket_reader, socket_writer) = socket.into_split();
let (data_reader, data_writer) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
let exit_receiver1 = exit_receiver.clone();
tracing::info!("Start spawn_write_loop.");
spawn_and_log_error(Self::spawn_write_loop(
exit_receiver1,
socket_writer,
data_writer,
));
tracing::info!("spawn_write_loop ok.");
let service_type = params.service_type();
let params2: SConnectionParams = params.clone();
tracing::info!("Start spawn_read_loop.");
spawn_and_log_error(Self::spawn_read_loop(
params.clone(),
exit_receiver.clone(),
socket_reader,
connection_state,
));
tracing::info!("spawn_read_loop ok");
tracing::info!("Notify client state Running.");
params2.call_back_inner().state_callback(
params2.handle_inner(),
EConnectionState::Running(data_reader.clone()),
service_type,
);
CTcpConnection {}
}
async fn spawn_read_loop(
conn_params: SConnectionParams,
mut exit_receiver: WatchReceiver<bool>,
mut socket_reader: OwnedReadHalf,
connection_state: ArcMutex<bool>,
) -> StdResult<()> {
let mut buf = Vec::new();
buf.resize(1024 * 10240, 0);
loop {
buf.clear();
tokio::select! {
flag = exit_receiver.changed() => {
match flag {
Ok(_exit_flag) => { break;}
Err(_) => {}
}
}
read_size = socket_reader.read_buf(&mut buf) => {
match read_size
{
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!("{} {} {:?}", conn_params.handle_ref(), " Recved ", String::from_utf8_lossy(&buf));
conn_params.call_back_inner().data_callback(buf.clone(), conn_params.handle_inner(), conn_params.serv_type);
}
Ok(size) if size == 0 => {
tracing::warn!("{} {}", conn_params.handle_ref(), " Recved Size = 0, Will exit.");
conn_params.call_back_inner().state_callback(conn_params.handle_inner(), EConnectionState::Interrupt, conn_params.serv_type);
let state_arc = connection_state.clone();
let mut state_mutex = state_arc.lock().unwrap();
*state_mutex = false;
break;
}
_ => unreachable!()
}
}
}
}
Ok(())
}
async fn spawn_write_loop(
mut exit_receiver: WatchReceiver<bool>,
mut socket_writer: OwnedWriteHalf,
mut data_writer: MpscAsyncReceiver<Vec<u8>>,
) -> StdResult<()> {
loop {
tokio::select! {
flag = exit_receiver.changed() => {
match flag {
Ok(_exit_flag) => { break;}
Err(_) => {}
}
}
msg = data_writer.recv() => {
match msg {
Some(msg_data) => {
let _ = socket_writer.write_all(&msg_data).await;
}
None => {}
}
}
}
}
Ok(())
}
}