【Rust投稿】从零实现消息中间件(4)-SERVER.CLIENT

2024-06-23 00:32

本文主要是介绍【Rust投稿】从零实现消息中间件(4)-SERVER.CLIENT,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

这部分主要说的是服务器端对于来自client连接的数据的处理. 主要功能包括

  1. 接收消息

  2. 收到sub消息,就记录到全局列表中

  3. 收到pub消息,就发送给相关订阅的client

  4. 出错,删除订阅,关闭连接

数据结构定义

Client中除了cid以外,其他两项都使用了Mutex进行保护,上一篇讲到过,凡是多线程读写的都需要Arc<Mutex>保护.

  • srv: 主要还是pub sub的时候都需要访问全局的sublist.

  • msg_sender: 之所以用Mutex保护是因为除了client自己要发送消息,当其他client pub消息的时候也要通过这个ClientMessageSender发送消息
    ClientMessageSender在我们这个版本中则非常简单,就是一个TcpStream的writer.
    rust #[derive(Debug)] pub struct Client<T: SubListTrait> { pub srv: Arc<Mutex<ServerState<T>>>, pub cid: u64, pub msg_sender: Arc<Mutex<ClientMessageSender>>, } #[derive(Debug)] pub struct ClientMessageSender { writer: WriteHalf<TcpStream>, }

代码实现

process_connection

  • 创建Client以及可以共享使用的ClientMessageSender

  • 启动client_task

    impl<T: SubListTrait + Send + 'static> Client<T> {pub fn process_connection(cid: u64,srv: Arc<Mutex<ServerState<T>>>,conn: TcpStream,
    ) -> Arc<Mutex<ClientMessageSender>> {let (reader, writer) = tokio::io::split(conn);let msg_sender = Arc::new(Mutex::new(ClientMessageSender::new(writer)));let c = Client {srv: srv,cid,msg_sender: msg_sender.clone(),};tokio::spawn(Client::client_task(c, reader));msg_sender
    }...
    }
    

client_task

主要功能:

  • 读取,解析消息

  • 分发消息给相应的处理函数

    • process_error

    • process_sub

    • process_pub

这个其实就是一个tcp连接的主循环,说到这里我想把tokio::spawn 和 go语言中的go关键字做一个类比.
在go中TcpServer接收到一个连接以后,紧接着就是单独起一个goroutine来处理.类似于go client.processConnection(),而到了Rust中基本上可以等价为

tokio::spawn(async move{Client::process_connection();
});

当然Rust重要复杂很多,涉及到所有权,生命周期等一系列问题.

 async fn client_task(self, mut reader: ReadHalf<TcpStream>) {let mut buf = [0; 1024];let mut parser = Parser::new();let mut subs = HashMap::new();loop {let r = reader.read(&mut buf[..]).await;if r.is_err() {let e = r.unwrap_err();self.process_error(e, subs).await;return;}let r = r.unwrap();let n = r;if n == 0 {self.process_error(NError::new(ERROR_CONNECTION_CLOSED), subs).await;return;}let mut buf = &buf[0..n];loop {let r = parser.parse(&buf[..]);if r.is_err() {self.process_error(r.unwrap_err(), subs).await;return;}let (result, left) = r.unwrap();match result {ParseResult::NoMsg => {break;}ParseResult::Sub(ref sub) => {if let Err(e) = self.process_sub(sub, &mut subs).await {self.process_error(e, subs).await;return;}}ParseResult::Pub(ref pub_arg) => {if let Err(e) = self.process_pub(pub_arg).await {self.process_error(e, subs).await;return;}}}if left == buf.len() {break;}buf = &buf[left..];}}}

从整个代码中也可以看出client_task的主要工作就是接受消息,并处理.

process_error

  1. 删除所有订阅

  2. 关闭连接
    rust async fn process_error<E: Error>(&self, err: E, subs: HashMap<String, ArcSubscription>) { println!("client {} process err {:?}", self.cid, err); { let mut sublist = &mut self.srv.lock().await.sublist; for (_, sub) in subs { sublist.remove(sub); } } let r = self.msg_sender.lock().await.writer.shutdown().await; if r.is_err() { println!("shutdown err {:?}", r.unwrap_err()); } }

process_sub

对于收到的sub则是

  1. 全局订阅列表中保存一份

  2. 本地连接保存一份,这样连接断开的时候好删除
    为了避免内存分配,我们的SubArg里面使用的还是Parer缓冲区中的内存,当我们需要在连接之外访问这些信息的时候,我们就必须单独保存一份了,这里我们用的是sub.subject.to_string()来分配一个新的内存.
    rust async fn process_sub( &self, sub: &SubArg<'_>, subs: &mut HashMap<String, ArcSubscription>, ) -> crate::error::Result<()> { let sub = Subscription { subject: sub.subject.to_string(), queue: sub.queue.map(|q| q.to_string()), sid: sub.sid.to_string(), msg_sender: self.msg_sender.clone(), }; let sub = Arc::new(sub); subs.insert(sub.subject.clone(), sub.clone()); let sublist = &mut self.srv.lock().await.sublist; sublist.insert(sub); Ok(()) }

process_pub

收到pub消息,

  1. 查找所有的订阅

  2. 将消息逐一转发给他们
    转发的过程中要稍微麻烦一点,因为考虑到设计中的负载均衡问题,qsubs则是从同一个queue中随机选择一个来推送消息.
    rust async fn process_pub(&self, pub_arg: &PubArg<'_>) -> crate::error::Result<()> { let sub_result = { let sub_list = &mut self.srv.lock().await.sublist; sub_list.match_subject(pub_arg.subject)? }; for sub in sub_result.subs.iter() { self.send_message(sub.as_ref(), pub_arg) .await .map_err(|e| NError::new(ERROR_CONNECTION_CLOSED))?; } //qsubs 要考虑负载均衡问题 let mut rng = rand::rngs::StdRng::from_entropy(); for qsubs in sub_result.qsubs.iter() { let n = rng.next_u32(); let n = n as usize % qsubs.len(); let sub = qsubs.get(n).unwrap(); self.send_message(sub.as_ref(), pub_arg) .await .map_err(|e| NError::new(ERROR_CONNECTION_CLOSED))?; } Ok(()) }

send_message

就是拼装消息格式
因为是第一个版本,也是展示关键api的使用,里面用到了大量的await,实际上没有必要.
实际项目中,肯定会使用缓冲区来做.

///消息格式
///```
/// MSG <subject> <sid> <size>\r\n
/// <message>\r\n
/// ```
async fn send_message(&self, sub: &Subscription, pub_arg: &PubArg<'_>) -> std::io::Result<()> {let writer = &mut sub.msg_sender.lock().await.writer;writer.write("MSG ".as_bytes()).await?;writer.write(sub.subject.as_bytes()).await?;writer.write(" ".as_bytes()).await?;writer.write(sub.sid.as_bytes()).await?;writer.write(" ".as_bytes()).await?;writer.write(pub_arg.size_buf.as_bytes()).await?;writer.write("\r\n".as_bytes()).await?;writer.write(pub_arg.msg).await?;writer.write("\r\n".as_bytes()).await?;Ok(())
}

其他

相关代码都在我的github rnats 欢迎围观

https://github.com/nkbai/learnrustbynats

这篇关于【Rust投稿】从零实现消息中间件(4)-SERVER.CLIENT的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1085827

相关文章

SQL server数据库如何下载和安装

《SQLserver数据库如何下载和安装》本文指导如何下载安装SQLServer2022评估版及SSMS工具,涵盖安装配置、连接字符串设置、C#连接数据库方法和安全注意事项,如混合验证、参数化查... 目录第一步:打开官网下载对应文件第二步:程序安装配置第三部:安装工具SQL Server Manageme

C#连接SQL server数据库命令的基本步骤

《C#连接SQLserver数据库命令的基本步骤》文章讲解了连接SQLServer数据库的步骤,包括引入命名空间、构建连接字符串、使用SqlConnection和SqlCommand执行SQL操作,... 目录建议配合使用:如何下载和安装SQL server数据库-CSDN博客1. 引入必要的命名空间2.

Linux下删除乱码文件和目录的实现方式

《Linux下删除乱码文件和目录的实现方式》:本文主要介绍Linux下删除乱码文件和目录的实现方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录linux下删除乱码文件和目录方法1方法2总结Linux下删除乱码文件和目录方法1使用ls -i命令找到文件或目录

SpringBoot+EasyExcel实现自定义复杂样式导入导出

《SpringBoot+EasyExcel实现自定义复杂样式导入导出》这篇文章主要为大家详细介绍了SpringBoot如何结果EasyExcel实现自定义复杂样式导入导出功能,文中的示例代码讲解详细,... 目录安装处理自定义导出复杂场景1、列不固定,动态列2、动态下拉3、自定义锁定行/列,添加密码4、合并

mybatis执行insert返回id实现详解

《mybatis执行insert返回id实现详解》MyBatis插入操作默认返回受影响行数,需通过useGeneratedKeys+keyProperty或selectKey获取主键ID,确保主键为自... 目录 两种方式获取自增 ID:1. ​​useGeneratedKeys+keyProperty(推

Spring Boot集成Druid实现数据源管理与监控的详细步骤

《SpringBoot集成Druid实现数据源管理与监控的详细步骤》本文介绍如何在SpringBoot项目中集成Druid数据库连接池,包括环境搭建、Maven依赖配置、SpringBoot配置文件... 目录1. 引言1.1 环境准备1.2 Druid介绍2. 配置Druid连接池3. 查看Druid监控

Linux在线解压jar包的实现方式

《Linux在线解压jar包的实现方式》:本文主要介绍Linux在线解压jar包的实现方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录linux在线解压jar包解压 jar包的步骤总结Linux在线解压jar包在 Centos 中解压 jar 包可以使用 u

c++ 类成员变量默认初始值的实现

《c++类成员变量默认初始值的实现》本文主要介绍了c++类成员变量默认初始值,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧... 目录C++类成员变量初始化c++类的变量的初始化在C++中,如果使用类成员变量时未给定其初始值,那么它将被

Qt使用QSqlDatabase连接MySQL实现增删改查功能

《Qt使用QSqlDatabase连接MySQL实现增删改查功能》这篇文章主要为大家详细介绍了Qt如何使用QSqlDatabase连接MySQL实现增删改查功能,文中的示例代码讲解详细,感兴趣的小伙伴... 目录一、创建数据表二、连接mysql数据库三、封装成一个完整的轻量级 ORM 风格类3.1 表结构

基于Python实现一个图片拆分工具

《基于Python实现一个图片拆分工具》这篇文章主要为大家详细介绍了如何基于Python实现一个图片拆分工具,可以根据需要的行数和列数进行拆分,感兴趣的小伙伴可以跟随小编一起学习一下... 简单介绍先自己选择输入的图片,默认是输出到项目文件夹中,可以自己选择其他的文件夹,选择需要拆分的行数和列数,可以通过