ugly_smart_lib 0.1.1

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

use tokio::net::{TcpListener, TcpStream};

use crate::{
    net_params::{ITcpConnectionParams, SConnectionParams, STcpServerParams},
    smart_pub::{ArcMutex, StdMutex},
};

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

/// TCP Server
#[derive(Debug)]
pub struct CTcpServer {
    exit_event_sender: WatchSender<bool>,
    exit_event_receiver: WatchReceiver<bool>,
    params: STcpServerParams,
}

impl CTcpServer {
    pub fn new(params: STcpServerParams) -> Self {
        let (event_sender, event_receiver) = tokio::sync::watch::channel(false);
        CTcpServer {
            exit_event_sender: event_sender,
            exit_event_receiver: event_receiver,
            params,
        }
    }

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

        let connection_state: ArcMutex<bool> = StdArc::new(StdMutex::new(true));

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

        Ok(())
    }
}

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

        // 绑定网络端口
        let listener = TcpListener::bind(self.params.address_inner()).await?;
        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));
                }
            };
        }

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

    #[tracing::instrument(name = "TCP 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)),
        }
    }
}