ugly_smart_lib 0.1.1

jonny's ugly smart lib for Rust
Documentation
//! src/udp_server.rs
//!
//! `server` 是一个会话服务器文件
//!
//! 它包含以下结构。
//! - Server TCP服务器

use tokio::net::UdpSocket;

use crate::{
    net_params::{IUdpConnectionParams, SUdpConnectionParams, SUdpServerParams},
    net_pub::spawn_and_log_error,
    smart_pub::{ArcMutex, StdMutex},
    udp_connection::CUdpConnection,
};

use super::{
    net_pub::{INetwork, ServiceType},
    smart_pub::{StdArc, StdResult, WatchReceiver, WatchSender},
};

/// TCP Server
#[derive(Debug)]
pub struct CUdpServer {
    exit_event_sender: WatchSender<bool>,
    exit_event_receiver: WatchReceiver<bool>,
    params: SUdpServerParams,
    connection: ArcMutex<Option<CUdpConnection>>,
}

impl CUdpServer {
    pub fn new(params: SUdpServerParams) -> Self {
        let (event_sender, event_receiver) = tokio::sync::watch::channel(false);
        CUdpServer {
            exit_event_sender: event_sender,
            exit_event_receiver: event_receiver,
            params,
            connection: StdArc::new(StdMutex::new(None)),
        }
    }

    // 建立新的连接
    #[tracing::instrument(name = "TCP Server Start", skip(self)/* , fields(cleint_address=tracing::field::Empty) */)]
    async fn connect_new(self: StdArc<Self>, stream: UdpSocket) -> StdResult<()> {
        // let remote_addr = stream.peer_addr().unwrap();
        let connection_params = SUdpConnectionParams {
            call_back: self.params.call_back_inner(),
            serv_type: ServiceType::Server,
            extra_params: self.params.local_address_inner(),
            remote_addr: "".into(),
        };

        tracing::info!("New Client");
        // tracing::info!("New Client {:?}", remote_addr);
        let conn = CUdpConnection::new(
            stream,
            connection_params,
            self.exit_event_receiver.clone(),
        );
        tracing::info!("New Client OK");

        let conn_arc = self.connection.clone();
        let mut conn_mutex = conn_arc.lock().unwrap();
        *conn_mutex = Some(conn);

        Ok(())
    }
}

/// 实现 server 端的特质
impl INetwork for CUdpServer {
    #[tracing::instrument(name = "UDP Server Start", skip(self)/* , fields(address=%self.params.handle_ref()) */)]
    async fn start(self: StdArc<Self>) -> StdResult<()> {
        tracing::info!("Start Server! {}", self.params.local_address_ref());

        // 绑定网络端口
        // let stream = UdpSocket::bind(self.params.local_address_inner()).await;
        match UdpSocket::bind(self.params.local_address_inner()).await {
            Ok(stream) => {
                spawn_and_log_error(self.clone().connect_new(stream));
                // stream.set_nonblocking(true);
            }
            Err(e) => {
                tracing::error!("{:?}", e);
            }
        }
        // let mut exit_receiver = self.clone().exit_event_receiver.clone();
        // loop {
        //     tokio::select! {
        //         flag = exit_receiver.changed() =>
        //         {
        //             match flag {
        //                 Ok(_exit_flag) => { break;}
        //                 Err(_) => {}
        //             }
        //         }
        //         ret = listener.accept() =>
        //         {
        //             let (stream, _addr) = ret?;
        //             // 建立新的连接
        //             tokio::spawn(self.clone().connect_new(stream));
        //         }
        //     };
        // }

        // spawn_and_log_error(self.clone().connect_new(stream));

        // tracing::info!("Stop Server!");
        Ok(())
    }

    #[tracing::instrument(name = "UDP Server Stop", skip(self), fields(address=%self.params.handle_ref()))]
    async fn stop(self: StdArc<Self>) -> StdResult<()> {
        let r = self.exit_event_sender.send(true);
        match r {
            Ok(()) => Ok(()),
            Err(e) => Err(Box::new(e)),
        }
    }
}