ugly_smart_lib 0.1.1

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

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;

/// TCP 连接
#[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::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()),
            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 ", String::from_utf8_lossy(&buf));
                            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 {
            // buf.clear();
            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 ", String::from_utf8_lossy(&buf));
                            // tracing::debug!("From {} {} {:?}", addr, " Recved ", buf);
                            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(())
    }
}