group_interface 2025.8.14

A brief description of your crate
Documentation

//use group_controller::GroupController;

use std::sync::Arc;
//use std::{any, ffi::c_void};
//use async_trait::async_trait;
//use serde::Serialize;
use tokio::io::AsyncWriteExt;
use tokio::net::{TcpStream,tcp::{OwnedReadHalf,OwnedWriteHalf}};
//use dashmap::DashMap;
//use serde_json::Value;
use serde::{Serialize,Deserialize};
//use std::cell::RefCell;
//use anyhow::Result;
//use tokio::io::util::AsyncWriteExt;
//use std::{result::Result  as stdResult};
//use std::io::Error as stdError;
//use std::sync::Arc;
use std::clone::Clone;


/* 

                         //   **************以下这几项都不能动,关乎数据转换*********************   
//设备模块接口
//#[async_trait]

//#[repr(C)] 不能用,用了DashMap没法插入 有#[repr(C)]的结构体不能使用rust复杂类型,它会强制rust结构体进行c布局,
///破坏结构,公共头文件中定义不透明结构体(C 可见),可以跨C边界,隐藏数据。只有rust发送者接受者明白这个指针代表的类型,
/// 可以安全的类型转换。为什么要跨C边界,因为所有操作系统都是c语言开发的,支持c结构数据
#[repr(C)]//结构体中的元素必须符合FFI规则
pub struct GroupControllerHandle;

// pub type StartGroupFn = unsafe fn(&mut GroupController);
///定义动态库跨边界调用函数
pub type StartGroupFn = unsafe fn(*mut GroupControllerHandle);
//pub type CreateModuleFn = unsafe fn() -> *mut dyn DeviceModule;

///Group的共享容器,键值String是Group的名字,Arc<GroupController>是Group的控制及资源接口,便于Group使用
 pub type GroupMap = Arc<DashMap<String, Arc<GroupController>>>;//可能被多个组同时访问,因为使用了DashMap可以读写,Arc只是提供指针
//pub type GroupMap = Arc<DashMap<String, Arc<Mutex<GroupController>>>>;//可能被多个组同时访问,因为使用了DashMap可以读写

///Group的控制及资源接口结构体,因为Group拥有其共享指针,可以安全的访问。其元素都符合Send+Sync
#[derive(Clone)] 
pub struct GroupController
{
       name:String,
       group_map: GroupMap,//设备类的集合引用,智能指针,是GroupMap(Arc)的克隆 Arc 也是Send+Sync
       node_map: DashMap<String,Arc<(Arc<OwnedReadHalf>,Arc<OwnedWriteHalf>)>>,//同组节点的TcpStream的集合,
       //外代码无法访问,是线程安全的,Send +Sync,这时的TcpClient可以读写,与Arc无关
     
}   
 impl  GroupController {
        pub fn new() ->Self{ 
        Self {
          name:"group".to_string(),
          group_map: Arc::new(DashMap::new()),
          node_map: DashMap::new(),
    
        }
       }
        ///设定组名称
        pub fn set_group_name(&mut self,name:String)
        { 
        self.name = name;
      // println!("setname..{}",&self.name);
        }
       ///得到组的名称
        pub fn get_group_name(&self) ->&str
        {
      //println!("getname..{}",&self.name);
         &self.name
        }
       ///设定group集合指针Arc
         pub fn set_group_map(&mut self,group_map:GroupMap){ self.group_map=group_map;}  

    }
  */
                    //以上几项都不能修改,否则程序转换会失败


/* 
 ///Group异步接口,支持多线程共享和传递
 #[async_trait]
 pub trait IController: Send + Sync {  //为了指针接口是同步的,不能Send+Sync
    
        fn test(&self);
     
        async fn log(&self,msg_class:&str,msg:&str) -> io::Result<()>;
       

         async  fn who(&self);
         
   

     /* 
     ///接收 TcpStream,并将其劈成读写两部分,以元组的方式存入DashMap中,便于多线程全双工读写
       pub  async   fn handle_socket(&self, node_id: String,  mut socket: TcpStream);
       
    */
      

    /// 从节点容器中找到目标节点,并获得读写句柄元组,便于多线程快速读写
    /// group是组名,同一类或者统一标准的类的名称,对应不同动态库;node_id是节点的名字或者id
    /// 支持寻址,可以找到其它Group的节点,并收发信息
    /// 下一步支持到其他程序或者其他网络上的节点的通讯
    // fn get_node_tcpclient(&self, group: &str, node_id: &str) -> Option<Arc<(Arc<OwnedReadHalf>,Arc<OwnedWriteHalf>)>>;
    

       //ReadHalf, WriteHalf是原始流的引用,而OwnedReadHalf,OwnedWriteHalf拥有Tcp流的完整所有权,分离后原始TcpStream不再可用
     ///共享读写,用Arc实现共享读写,将TcpStream劈成读写两部分,便于异步线程同时读写,
     /// 减少TcpStream读写的竞争,因为TcpStream是全双工的
      fn sharing_stream(stream: TcpStream) -> (Arc<OwnedReadHalf>,Arc<OwnedWriteHalf>);
    
}
*/

 // #[async_trait]
 // impl  IController for GroupController  {
    
      
     /// 程序测试用,会在终端打印出"controller test......ok",说明调用GroupController的功能正常
     pub      fn test(){ println!("group_interface test......ok"); }
     
     pub     async fn log(name:&str,msg_class:&str,msg:&str) -> io::Result<()> 
          {

            logs(name, msg_class, msg).await?;
          // 测试日志写入
           // log("test_group", "error", "系统启动失败").await?;
           // log("test_group", "error", "连接超时").await?;
          
          // 添加跨小时测试
          //  log("test_group", "debug", "调试信息").await?;
          
          // 确保缓冲区内容刷新到文件
            if let Some((_, mut writer)) = NOW_BUFWRITER.write().await.take() {
              writer.flush().await?;
             }
          
              Ok(())
           }

        pub   async  fn who()
          {
            println!("I am liurunzi from rizhao 柿树圆")
           }
     //  }
     

     /* 
     ///接收 TcpStream,并将其劈成读写两部分,以元组的方式存入DashMap中,便于多线程全双工读写
       pub  async   fn handle_socket(&self, node_id: String,  mut socket: TcpStream)
       {
          //  println!("rrrrr");
        // 使用 Arc 包装 TcpStream 以便安全共享
         let message = format!("{}上线!可以收发信息了。\n", node_id);
         let gg=  socket.write_all(message.as_bytes()).await;
         match gg {
             Ok(()) => println!("发送字符成功"),
             Err(e) => println!("发送失败{}",e),
            }
       // println!("成功发送{}字节",3);
          let node_id_str=node_id.clone();
        // 使用 DashMap 无需显式获取锁
  
         // 使用entry API避免竞争条件
         if let Some(old_socket) = 
                  self.node_map.insert(node_id, Arc::new(GroupController::sharing_stream(socket))) {
          // 处理被替换的旧socket(如果需要)
          drop(old_socket);
          eprintln!("Warning: Replaced existing socket for node {}",node_id_str);
          }
      
         //println!("ddddddddd");
          println!("[Node {} socket stored. Total sockets into_split: {}", self.name,self.node_map.len());
         // Ok(())
    }
    */
      
/* 
    /// 从节点容器中找到目标节点,并获得读写句柄元组,便于多线程快速读写
    /// group是组名,同一类或者统一标准的类的名称,对应不同动态库;node_id是节点的名字或者id
    /// 支持寻址,可以找到其它Group的节点,并收发信息
    /// 下一步支持到其他程序或者其他网络上的节点的通讯
     fn get_node_tcpclient(&self, group: &str, node_id: &str) -> Option<Arc<(Arc<OwnedReadHalf>,Arc<OwnedWriteHalf>)>>
    {
       
      if group == &self.name {

              
        self.node_map.get(node_id).map(|entry| {
          // 克隆Arc,增加引用计数(不复制流本身)
          Arc::clone(entry.value())
           //  entry
         //  Box::new(entry.value())
          // Box::new(**entry.value())
                      })
                    }
     
                      else {
                        self.group_map.get(group)
                            // 从子组中获取节点映射
                      .and_then(|group_entry: dashmap::mapref::one::Ref<'_, String, Arc<GroupController>>| {
                           group_entry.value().node_map.get(node_id)
                           .map(|node_entry|Arc::clone(node_entry.value()))}
                           )
                    }
      }
      */
        // }

       //ReadHalf, WriteHalf是原始流的引用,而OwnedReadHalf,OwnedWriteHalf拥有Tcp流的完整所有权,分离后原始TcpStream不再可用
     ///共享读写,用Arc实现共享读写,将TcpStream劈成读写两部分,便于异步线程同时读写,
     /// 减少TcpStream读写的竞争,因为TcpStream是全双工的
    pub   fn sharing_stream(stream: TcpStream) -> (Arc<OwnedReadHalf>,Arc<OwnedWriteHalf>)
        {
       
      // 分离为读半部分和写半部分
        let (read_half, write_half) = stream.into_split();
      
      // 用 Arc 共享读句柄(读操作不需要可变引用)
        let read_arc = Arc::new(read_half);
      // 用 Arc 共享写句柄(写操作需要 &mut,但 WriteHalf 实现了 Unpin)
        let write_arc = Arc::new(write_half);
      
      
        let r: Arc<tokio::net::tcp::OwnedReadHalf> = Arc::clone(&read_arc);
        let w: Arc<tokio::net::tcp::OwnedWriteHalf> = Arc::clone(&write_arc);

      // println!("[sharing {}_stream iss ok", self.name );
        (r,w) 
       }
//  }
  
  

/* 

 ///Group异步接口,支持多线程共享和传递
 #[async_trait]
 pub trait GroupTrait: Send + Sync {  //为了指针接口是同步的,不能Send+Sync
    
     fn set_controller_to_group(&mut self, controller: Arc<GroupController>);
     fn get_controller_clone(&self) ->Arc<GroupController>;
    /// 处理新连接的socket
     fn start(controller: Arc<GroupController>);
    
   ///async fn send_to_socket(&self)->stdResult<(), stdError>;
      async  fn handle(&self) ->stdResult<(), stdError>;
      async  fn test(&self);
    //  async  fn write_data(device_id:String,data_buf:&[u8]) ->stdResult<(), std::io::Error>;
    //async  fn read_data(device_id: String) -> stdResult<Vec<u8>, std::io::Error>;
}
*/



       //全局运行时



// 全局运行时(程序生命周期内只初始化一次),所有的group程序都在这里运行
 use std::sync::OnceLock;
 use tokio::runtime::{Runtime, Handle};
///全局运行时,所有Group共同使用的
 static GLOBAL_RUNTIME: OnceLock<Runtime> = OnceLock::new();

/// 获取全局运行时句柄
pub fn global_runtime_handle() -> Handle {
    GLOBAL_RUNTIME
        .get_or_init(|| Runtime::new().expect("无法创建全局运行时"))
        .handle()
        .clone()
      
  }




//use std::net::IpAddr; // 使用标准库的IpAddr类型替代String

/// 发送目标结构
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct Target {
    pub app: String,            // 目标应用名称
    pub class: String,          // 功能类或设备类名称 (避免使用cls这样的缩写)
    pub node: String,         // 设备名称或编号
    pub message: String,        // 消息内容
}

/* 
/// 注册信息包结构
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct Reg {
    pub class: String,          // 设备类型
    pub node: String,         // 设备名称或账号
    pub version: Option<String>,// 可选的版本信息
    pub metadata: Option<String>,// 可选的元数据
}
 */

                      //*****登记信息*****

use tokio::fs::{self, File, OpenOptions};
use tokio::io::{self, BufWriter};
use tokio::sync::RwLock;
use std::path::PathBuf;
use chrono::Utc;
use once_cell::sync::Lazy;

/// 全局缓冲区:存储当前使用的文件名和对应的BufWriter<File>
static NOW_BUFWRITER: Lazy<RwLock<Option<(PathBuf, BufWriter<File>)>>> = Lazy::new(|| {
    RwLock::new(None)
});

/// 辅助函数:获取格式化后的时间
fn current_formatted_time() -> String {
    Utc::now().format("%Y-%m-%d %H:%M:%S").to_string()
   // Local::now().format("%Y-%m-%d %H:%M:%S").to_string()
}

/// 异步日志记录函数
/// group_name:组名,功能类型;msg_class:信息类型;msg:登记信息
async fn logs(group_name: &str, msg_class: &str, msg: &str) -> io::Result<()> {
    // 1. 准备日志文件名
    let now = Utc::now();
   //let now = Local::now();
  // println!("当前时间{}",now.to_string());
    let file_name = PathBuf::from(format!(
        "./info/{}/{}/{}/{}/{}_{}.log",
        group_name,
        msg_class,
        now.format("%m"),
        now.format("%d"),
        now.format("%H"),
        msg_class
    ));
   // println!("文件名{}",file_name.display());
    // 2. 确保日志目录存在
    if let Some(parent) = file_name.parent() {
        fs::create_dir_all(parent).await?;
    }

    // 3. 获取全局缓冲区的写锁
    let mut writer_guard = NOW_BUFWRITER.write().await;

    // 4. 检查是否需要切换文件
    let should_switch = match &*writer_guard {
        Some((cached_path, _)) => cached_path != &file_name,
        None => true,
    };

    // 5. 处理文件切换逻辑
    if should_switch {
        // 刷新并关闭之前的文件
        if let Some((_, mut writer)) = writer_guard.take() {
            writer.flush().await?;
        }

        // 打开新文件(创建+追加模式)
        let file = OpenOptions::new()
            .create(true)
            .append(true)
            .open(&file_name)
            .await?;

        // 创建带缓冲区的写入器
        *writer_guard = Some((file_name.clone(), BufWriter::with_capacity(8192 * 4, file)));
    }

    // 6. 写入日志内容
    if let Some((_, writer)) = &mut *writer_guard {
        let timestamp = current_formatted_time();
        writer
            .write_all(format!("[Utc {}] {}\n", timestamp, msg).as_bytes())
            .await?;
    }

    Ok(())
}