ugly_smart_lib 0.1.1

jonny's ugly smart lib for Rust
Documentation
//! connection crate
//!
//! `connection` 是一个tcp网络连接文件
//!
//! 它包含以下结构。
//! - SConnectionParams 参数
//! - Connection 连接
//!

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},
};
// use std::sync::atomic::{AtomicBool, Ordering};

/// TCP 连接
#[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 {
        // tracing::Span::current().record("state", &tracing::field::display(connection_state.clone()));
        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::Span::current().record("state", &tracing::field::display(&credentials.username));
        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(())
    }
}