【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

相关文章

C#控制台程序同步调用WebApi实现方式

《C#控制台程序同步调用WebApi实现方式》控制台程序作为Job时,需同步调用WebApi以确保获取返回结果后执行后续操作,否则会引发TaskCanceledException异常,同步处理可避免异... 目录同步调用WebApi方法Cls001类里面的写法总结控制台程序一般当作Job使用,有时候需要控制

SpringBoot集成P6Spy的实现示例

《SpringBoot集成P6Spy的实现示例》本文主要介绍了SpringBoot集成P6Spy的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面... 目录本节目标P6Spy简介抛出问题集成P6Spy1. SpringBoot三板斧之加入依赖2. 修改

Python实现数据可视化图表生成(适合新手入门)

《Python实现数据可视化图表生成(适合新手入门)》在数据科学和数据分析的新时代,高效、直观的数据可视化工具显得尤为重要,下面:本文主要介绍Python实现数据可视化图表生成的相关资料,文中通过... 目录前言为什么需要数据可视化准备工作基本图表绘制折线图柱状图散点图使用Seaborn创建高级图表箱线图热

Redis分布式锁中Redission底层实现方式

《Redis分布式锁中Redission底层实现方式》Redission基于Redis原子操作和Lua脚本实现分布式锁,通过SETNX命令、看门狗续期、可重入机制及异常处理,确保锁的可靠性和一致性,是... 目录Redis分布式锁中Redission底层实现一、Redission分布式锁的基本使用二、Red

基于Python实现数字限制在指定范围内的五种方式

《基于Python实现数字限制在指定范围内的五种方式》在编程中,数字范围限制是常见需求,无论是游戏开发中的角色属性值、金融计算中的利率调整,还是传感器数据处理中的异常值过滤,都需要将数字控制在合理范围... 目录引言一、基础条件判断法二、数学运算巧解法三、装饰器模式法四、自定义类封装法五、NumPy数组处理

Python中经纬度距离计算的实现方式

《Python中经纬度距离计算的实现方式》文章介绍Python中计算经纬度距离的方法及中国加密坐标系转换工具,主要方法包括geopy(Vincenty/Karney)、Haversine、pyproj... 目录一、基本方法1. 使用geopy库(推荐)2. 手动实现 Haversine 公式3. 使用py

MySQL进行分片合并的实现步骤

《MySQL进行分片合并的实现步骤》分片合并是指在分布式数据库系统中,将不同分片上的查询结果进行整合,以获得完整的查询结果,下面就来具体介绍一下,感兴趣的可以了解一下... 目录环境准备项目依赖数据源配置分片上下文分片查询和合并代码实现1. 查询单条记录2. 跨分片查询和合并测试结论分片合并(Shardin

Spring Security重写AuthenticationManager实现账号密码登录或者手机号码登录

《SpringSecurity重写AuthenticationManager实现账号密码登录或者手机号码登录》本文主要介绍了SpringSecurity重写AuthenticationManage... 目录一、创建自定义认证提供者CustomAuthenticationProvider二、创建认证业务Us

MySQL配置多主复制的实现步骤

《MySQL配置多主复制的实现步骤》多主复制是一种允许多个MySQL服务器同时接受写操作的复制方式,本文就来介绍一下MySQL配置多主复制的实现步骤,具有一定的参考价值,感兴趣的可以了解一下... 目录1. 环境准备2. 配置每台服务器2.1 修改每台服务器的配置文件3. 安装和配置插件4. 启动组复制4.

MySQL数据脱敏的实现方法

《MySQL数据脱敏的实现方法》本文主要介绍了MySQL数据脱敏的实现方法,包括字符替换、加密等方法,通过工具类和数据库服务整合,确保敏感信息在查询结果中被掩码处理,感兴趣的可以了解一下... 目录一. 数据脱敏的方法二. 字符替换脱敏1. 创建数据脱敏工具类三. 整合到数据库操作1. 创建服务类进行数据库