You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.25 02:36:03