如何在actix的add_stream中集成awc websocket客户端实现通信
actix-awc Websocket客户端实现方案
use actix::{Actor, AsyncContext, Context, Running, StreamHandler, WrapFuture, Message, Handler}; use actix_web_actors::ws::{Frame, ProtocolError}; use awc::ws::{Codec, Message as WsMessage}; use awc::BoxedSocket; use futures::{FutureExt, SinkExt, StreamExt, stream::SplitSink}; use log::info; use openssl::ssl::SslConnector; use serde_json::json; use tokio::net::TcpSocket; pub struct RosClient { // 存储websocket写端,用于后续发消息 pub ws_sender: Option<SplitSink<awc::ws::Framed<BoxedSocket, Codec>, WsMessage>>, } impl RosClient { pub fn new() -> Self { Self { ws_sender: None } } } impl Actor for RosClient { type Context = Context<Self>; fn started(&mut self, ctx: &mut Self::Context) { // 连接逻辑包装成actor future,可以访问self和ctx let connect_fut = async move { let ssl = { let mut ssl = SslConnector::builder(openssl::ssl::SslMethod::tls()).unwrap(); let _ = ssl.set_alpn_protos(b"\x08http/1.1"); ssl.build() }; let connector = awc::Connector::new().ssl(ssl).finish(); let ws = awc::ClientBuilder::new() .connector(connector) .finish() .ws("ws://127.0.0.1:9090") .set_header("Host", "0.0.0.0:9090"); let (resp, connection) = ws .connect() .await .map_err(|e| { println!("连接失败:{:?}", e); e }) .unwrap(); println!("连接响应:{:?}", resp); // 拆分Framed为读流和写端 let (sink, stream) = connection.split(); // 注册读流到上下文,自动触发StreamHandler的handle方法处理收到的消息 ctx.add_stream(stream); // 存储写端到结构体,供后续发消息使用 self.ws_sender = Some(sink); // 发送初始订阅消息 let message = json!({ "op": "subscribe", "topic": "/client_count" }) .to_string(); if let Some(sender) = self.ws_sender.as_mut() { sender.send(WsMessage::Text(message)).await.unwrap(); } } .into_actor(self); ctx.spawn(connect_fut); } fn stopping(&mut self, _ctx: &mut Self::Context) -> Running { Running::Stop } } impl StreamHandler<Result<Frame, ProtocolError>> for RosClient { fn handle(&mut self, item: Result<Frame, ProtocolError>, _ctx: &mut Self::Context) { match item.unwrap() { Frame::Text(text_bytes) => { println!("收到文本消息:{:?}", std::str::from_utf8(&text_bytes)); } Frame::Binary(_) => {} Frame::Continuation(_) => {} Frame::Ping(bin) => { println!("收到Ping:{:?}", std::str::from_utf8(&bin)) } Frame::Pong(_) => {} Frame::Close(_) => {} } } } // 示例:定义发送消息的actor消息,用于后续主动给服务端发消息 #[derive(Message)] #[rtype(result = "()")] pub struct SendWsMessage(pub String); impl Handler<SendWsMessage> for RosClient { type Result = (); fn handle(&mut self, msg: SendWsMessage, ctx: &mut Self::Context) -> Self::Result { if let Some(sender) = self.ws_sender.as_mut() { // 把发送逻辑包装成future提交到上下文执行,避免阻塞actor let send_fut = async move { if let Err(e) = sender.send(WsMessage::Text(msg.0)).await { println!("发送消息失败:{:?}", e); } } .into_actor(self); ctx.spawn(send_fut); } } }
问题1:add_stream正确实现逻辑
- 不能直接把连接用的future转成stream注册,需要等websocket握手完成拿到
Framed对象后,拆分出只读的stream端,直接传给ctx.add_stream()即可,注册后框架会自动消费stream中的消息,触发StreamHandler的handle方法处理收到的帧。
问题2:存储连接用于后续发送消息
Framed<BoxedSocket, Codec>本身实现了Sink和Stream双trait,使用futures提供的split()方法可以将其拆分为独立的SplitSink(写端)和SplitStream(读端),两者互不影响,不需要克隆。- 将拆分得到的
SplitSink存到actor的结构体字段中,后续需要主动发消息时,直接调用sink.send()方法即可,示例中额外提供了Handler实现,展示了如何通过actor消息触发websocket消息发送。
内容的提问来源于stack exchange,提问作者Fabrex
相关产品推荐
相关产品推荐

