Rust下Actix WebSocket与RabbitMQ(lapin)双向消息代理集成咨询
我已经使用Rust的Actix框架编写了一个WebSocket服务,现在需要将RabbitMQ集成到该项目中,我找到了lapin crate用于RabbitMQ操作,但在将其与Actix框架集成时遇到了问题,我希望基于现有的WebSocket实现,将RabbitMQ的消息代理转发到客户端。
以下是我初步的实现方案,目前还处于非常早期的阶段,如果我的实现路径存在问题请告知我,方便我及时调整方向。
现有项目上下文
WebSocket Actor会和另一个存储Socket与房间信息的Actor通信,核心代码如下:
pub struct WSConn { pub id: Uuid, room_id: Uuid, hb: Instant, lobby_address: Addr<Lobby>, } impl WSConn { pub fn new(lobby: Addr<Lobby>, rabbit: Addr<MyRabbit>) -> Self { Self { id: Uuid::new_v4(), room_id: Uuid::nil(), hb: Instant::now(), lobby_address: lobby, } } fn hb(&self, context: &mut WebsocketContext<Self>) { /// HEARTBEAT CODE GOES HERE /// } }
当Actor启动时,我会向Lobby发送一条连接消息,Lobby会处理将连接添加到自身结构体的全部逻辑。
impl Actor for WSConn { type Context = WebsocketContext<Self>; fn started(&mut self, ctx: &mut Self::Context) { info!("Starting hearbeat"); self.hb(ctx); let wsserver_address = ctx.address(); info!("A new client has connected with id {}", self.id); self.lobby_address.send(Connect { address: wsserver_address.recipient(), room_id: self.room_id, id: self.id }).into_actor(self).then(|res, _, ctx| { match res { Ok(_) => (), _ => ctx.stop() } fut::ready(()) }).wait(ctx); } fn stopping(&mut self, _ctx: &mut Self::Context) -> Running { info!("Stopping actor"); self.lobby_address.do_send(Disconnect { room_id: self.room_id, id: self.id }); Running::Stop } }
Lobby的定义如下:
type Socket = Recipient<WSServerMessage>; pub struct Lobby { pub sessions: DashMap<Uuid, Socket>,//self id to self pub rooms: DashMap<Uuid, DashSet<Uuid>>,//room id to list of users id }
Lobby的完整代码量较大,此处不贴出,如果需要查看可以提供相关代码。客户端连接后会被分配到默认房间,客户端发送的消息会由StreamHandler处理。
impl StreamHandler<Result<ws::Message, ws::ProtocolError>> for WSConn { fn handle(&mut self, item: Result<Message, ProtocolError>, ctx: &mut Self::Context) { match item.unwrap() { ws::Message::Binary(bin) => ctx.binary(bin), ws::Message::Ping(bin) => { self.hb = Instant::now(); ctx.pong(&bin); } ws::Message::Pong(_) => self.hb = Instant::now(), ws::Message::Close(reason) => { ctx.close(reason); ctx.stop(); } ws::Message::Text(text) => { let command = serde_json::from_str::<Command>(&text) .expect(&format!("Can't parse message {}", &text)); info!("{:?}", command); if command.command.starts_with("/") { info!("This is a {} request", command.command); match command.command.as_ref() { "/join" => { info!("Join Room {}", command.payload); let uid = Uuid::from_str(command.payload.as_str().unwrap()).expect("Can't parse message {} to uuid"); self.lobby_address.send(Join { current_room: self.room_id, room_id: uid, id: self.id }).into_actor(self).then(|res, _, ctx| { match res { Ok(_) => (), _ => ctx.stop() } fut::ready(()) }).wait(ctx); self.room_id = uid; }, _ => () } } else { info!("Text is {}", text); self.lobby_address.do_send(ClientActorMessage { id: self.id, msg: command, room_id: self.room_id, }); } } _ => { info!("Something weird happened. Closing"); self.lobby_address.do_send(Disconnect { room_id: self.room_id, id: self.id }); ctx.stop();} } } }
如上述代码所示,客户端发送携带/join命令和有效uuidv4 payload的消息即可加入对应房间,我限制了客户端同一时间只能加入一个房间,加入新房间时会自动退出上一个房间。
目前RabbitMQ集成实现
我计划使用连接池维护RabbitMQ连接,再基于连接创建Channel,首先我定义了存储连接池的结构体:
use actix::{Actor, Context, Handler, StreamHandler}; use deadpool_lapin::{Config, Pool, Runtime}; use deadpool_lapin::lapin::Error; use lapin::message::Delivery; use crate::lapin_server::messages::CreateChannel; pub struct MyRabbit { pub pool: Pool, } impl MyRabbit { pub fn new() -> Self { let mut cfg = Config::default(); cfg.url = Some("amqps://ghqcmhat:KbhPAA309QRg7TjdgFEV14pQRheoh44P@codfish.rmq.cloudamqp.com/ghqcmhat".into()); let new_pool = cfg.create_pool(Some(Runtime::Tokio1)).expect("Can't create pool"); MyRabbit { pool: new_pool } } } impl Actor for MyRabbit { type Context = Context<Self>; }
将其封装为Actor后,我可以在启动WebSocket服务的同时启动该Actor,这部分逻辑在main函数中实现:
#[actix_web::main] async fn main() -> std::io::Result<()>{ std::env::set_var("RUST_LOG", "actix_web=info,info"); env_logger::init(); // Start Lobby actor and get his address let websocket_lobby = Lobby::default().start(); let rabbit = MyRabbit::new().start(); let application_data = web::Data::new(Appdata::new()); info!("Starting server on 127.0.0.1:8080"); let server = HttpServer::new(move || { App::new() .wrap(Logger::default()) .route("/ws/",web::get().to(websocket_handler)) .app_data(Data::new(websocket_lobby.clone())) .app_data(Data::new(rabbit.clone())) .app_data(application_data.clone()) }); server.bind("127.0.0.1:8080")?.run().await }
为了适配新的Actor,我在请求处理器中新增了对应参数:
pub async fn websocket_handler(request: HttpRequest, stream: web::Payload, srv: Data<Addr<Lobby>>, rab: Data<Addr<MyRabbit>>, data: Data<Appdata>) -> Result<HttpResponse, Error> { let mut counter = data.counter.lock().unwrap(); counter.add_assign(1); info!("This is request # {}", counter); let ws = WSConn::new(srv.get_ref().clone(), rab.get_ref().clone()); let response = ws::start(ws, &request, stream); debug!("Response: {:?}", &response); response }
调整后的WSConn结构体如下:
pub struct WSConn { pub id: Uuid, room_id: Uuid, hb: Instant, lobby_address: Addr<Lobby>, rabbit_address: Addr<MyRabbit> }
我知道要从RabbitMQ的主题消费消息,需要exchange名称、类型、routing key和队列名称,因此我也定义了对应的结构体:
pub struct Channel { pub queue_name: String, pub exchange_name: String, pub exchange_type: ExchangeKind, pub routing_key: String } impl Default for Channel { fn default() -> Self { Self { queue_name: "".to_string(), exchange_name: "".to_string(), exchange_type: Default::default(), routing_key: "".to_string() } } }
待解决问题与最终目标
目前我卡在了后续实现环节:
- 不确定deadpool_lapin是否是适用于该场景的crate
- 不清楚如何适配lapin官方文档中的示例,官方示例使用
async_global_executor::block_on,且通过async_global_executor::spawn创建新线程来消费消息,不知道如何和Actix的Actor模型结合
我最终要实现的效果是WebSocket与RabbitMQ双向消息代理:客户端连接WebSocket后发送如下格式的消息:
{ "command": "SUBSCRIBE", "payload": "topic_name" }
就可以收到对应RabbitMQ主题下发布的所有消息,发送UNSUBSCRIBE命令即可取消订阅。恳请各位提供相关实现建议,如果需要更多信息可以告知我。
内容的提问来源于stack exchange,提问作者Fabrex

