ugly_smart_lib 0.1.1

jonny's ugly smart lib for Rust
Documentation
//! src/client.rs
//!
//! `client` 是一个会话客户端文件
//!
//! 它包含以下结构。
//! - Client 会话客户端,有重连机制
//!

use crate::{
    net_params::{ITcpConnectionParams, STcpClientParams, SConnectionParams},
    net_pub::{EConnectionState, INetwork, ServiceType},
    smart_pub::{ArcMutex, StdArc, StdMutex, StdResult, WatchReceiver, WatchSender},
    tcp_connection::CTcpConnection,
};
// use std::sync::atomic::{AtomicBool, Ordering};

/// 使用 Lazy_static 来创建一个线程安全的单例  
// use lazy_static::lazy_static;

// lazy_static! {
//     static ref CLIENT_STATE_FLAG: AtomicBool = AtomicBool::new(false);
// }

#[derive(Debug)]
pub struct CTcpClient {
    exit_event_sender: WatchSender<bool>,
    exit_event_receiver: WatchReceiver<bool>,
    params: STcpClientParams,
}

impl CTcpClient {
    #[tracing::instrument(name = "TCP Client New")]
    pub fn new(params: STcpClientParams) -> Self {
        let (exit_sender, exit_receiver) = tokio::sync::watch::channel(false);
        CTcpClient {
            exit_event_sender: exit_sender,
            exit_event_receiver: exit_receiver,
            params,
        }
    }

    async fn connect_to(self: StdArc<Self>) -> StdResult<()> {
        let duration = tokio::time::Duration::from_secs(1);
        self.clone().execute_state(EConnectionState::Connecting);
        let client_connection_state: ArcMutex<bool> = StdArc::new(StdMutex::new(false));// AtomicBool::new(false);
        loop {
            match self.clone().exit_event_receiver.has_changed() {
                Ok(state) if state => {
                    tracing::debug!("Exit Signal Recved");
                    break;
                }
                Err(e) => {
                    tracing::error!("{:?}", e);
                }
                _ => {}
            }

            // let cnn_state = CLIENT_STATE_FLAG.load(Ordering::Relaxed);
            let cnn_state = {
                let state_arc = client_connection_state.clone();
                let state_mutex = state_arc.lock().unwrap();
                *state_mutex
            };

            if cnn_state {
                tokio::time::sleep(duration).await;
                continue;
            }

            let ret = tokio::net::TcpStream::connect(self.params.address_inner()).await;
            match ret {
                Ok(stream) => {
                    let conn_params = SConnectionParams {
                        call_back: self.params.call_back_inner(),
                        serv_type: ServiceType::Client,
                        extra_params: self.params.handle_inner(),
                    };

                    tracing::info!("Start new connection");
                    let _conn =
                        CTcpConnection::new(stream, conn_params, self.exit_event_receiver.clone(), client_connection_state.clone());
                    tracing::info!("Start connection OK");

                    // CLIENT_STATE_FLAG.store(true, Ordering::Relaxed);
                    {
                        let state_arc = client_connection_state.clone();
                        let mut state_mutex = state_arc.lock().unwrap();
                        *state_mutex = true;
                    };
                }
                Err(e) => {
                    tracing::error!(
                        "Connection Failure [{:?}], Try reconnecting after 1 second",
                        e
                    );
                    tokio::time::sleep(duration).await;
                }
            }
        }
        self.clone().execute_state(EConnectionState::Stoped);
        Ok(())
    }

    fn execute_state(self: StdArc<Self>, state: EConnectionState<Vec<u8>>) {
        let call_back_option = self.params.call_back_inner();
        call_back_option.clone().state_callback(
            self.params.handle_inner(),
            state,
            ServiceType::Client,
        );
    }
}

impl INetwork for CTcpClient {
    #[tracing::instrument(name = "TCP Client Start!", skip(self), fields(
        address = %self.params.address_ref()
    ))]
    async fn start(self: StdArc<Self>) -> StdResult<()> {
        tracing::info!("Start Tcp Client");
        loop {
            match self.clone().exit_event_receiver.has_changed() {
                Ok(state) if state => {
                    tracing::debug!("Exited State Recved");
                    break;
                }
                Err(e) => {
                    tracing::error!("{:?}", e);
                }
                _ => {}
            }
            let ret = self.clone().connect_to().await;
            match ret {
                _ => {}
            }
        }
        tracing::info!("Stop Tcp Client");
        Ok(())
    }

    #[tracing::instrument(name = "TCP Client Stop!", skip(self), fields(
        address = %self.params.address_ref()
    ))]
    async fn stop(self: StdArc<Self>) -> StdResult<()> {
        tracing::info!("TCP Client will stop");
        let _ = self.exit_event_sender.send(true);
        Ok(())
    }
}

// impl Drop for Client {
//     fn drop(&mut self) {
//         let _ = self.exit_event_sender.send(true);
//     }
// }