拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Rig 的 WebSocket 传输后端:rig-tungstenite 源码与实战指南

Rig 的 WebSocket 传输后端:rig-tungstenite 源码与实战指南
  • AI Agent
  • Agent 框架
  • RAG
  • 后端

【免费下载链接】rig

⚙️🦀 Build modular and scalable LLM Applications in Rust

项目地址:https://gitcode.com/GitHub_Trending/rig2/rig
点击查看免费下载

本篇指南聚焦 Rig 仓库中crates/rig-tungstenite这一官方捆绑的 WebSocket 传输后端:它基于tokio-tungstenite实现rig_http::ws_client::WebSocketClientExt契约,让 OpenAI Responses 会话可以无缝运行在 WebSocket 之上。读完本文,你将掌握如何在 Rig 中开启 WebSocket 会话、如何显式注入自定义后端、该后端在有无 Tokio 运行时两种场景下的工作机制,以及升级拒绝、超时与帧序等边界行为是如何被设计与验证的。

分层架构:协议归 rig-core,契约归 rig-http,Socket 归 rig-tungstenite

Rig 对 WebSocket 支持做了清晰的三层切分,rig-tungstenite只是其中负责"物理 Socket"的最底层:

  • rig-core拥有 WebSocket 协议:OpenAI Responses 会话、事件信封(event envelope)以及整个 turn 状态机都位于rig_core::providers::openai::responses_api::websocket(见 crates/rig-core/src/providers/openai/responses_api/websocket.rs)。这一层面向传输无关的rig_http::ws_client契约编写,并在rig_core::ws_client处重新导出。
  • rig-http持有传输契约:包括WebSocketClientExt、WebSocketConnection、Frame、ConnectOptions等抽象,见 crates/rig-http/src/ws_client.rs。
  • rig-tungstenite只拥有 Socket:TungsteniteClient是该契约的一个具体实现,内部驱动tokio-tungstenite完成握手、收发帧与关闭握手。正如rig-reqwest只拥有 HTTP 传输一样,这个 crate 不关心任何会话协议细节。

这种"协议与传输分离"的设计带来一个直接好处:WebSocket 后端是可插拔的。只要实现WebSocketClientExt,任何自定义后端(例如基于web_sys::WebSocket的浏览器实现)都可以接入 Rig 的 Responses WebSocket 会话。

从依赖关系看,rig-tungstenite只依赖rig-http(开启websocketfeature),外加futures与thiserror(见 crates/rig-tungstenite/Cargo.toml)。tokio与tokio-tungstenite则被限定在cfg(not(target_family = "wasm"))条件下引入,因为这本质是一个面向原生环境的实现。

快速上手:一行代码打开 Responses WebSocket 会话

方式一:启用tungstenitefeature,免指定后端

在rig-core中启用tungstenitefeature 后,ResponsesWebSocketSessionBuilder::connect()会自动使用内置的TungsteniteClient,无需命名后端:

let session = model.responses_websocket().connect().await?;

对应的实现位于 crates/rig-core/src/providers/openai/responses_api/websocket.rs#L251-L258:

#[cfg(all(feature = "tungstenite", not(target_family = "wasm")))] pub async fn connect(self) -> Result<ResponsesWebSocketSession, ProviderError> { self.connect_with(&rig_tungstenite::TungsteniteClient::new()) .await }

也就是说connect()本质上就是connect_with(&TungsteniteClient::new())的便捷封装。TungsteniteClient是一个无状态后端(#[derive(Clone, Copy, Debug, Default)],见 crates/rig-tungstenite/src/lib.rs),握手配置(URI、认证头、超时)全部来自每次请求本身。

方式二:显式传入后端

未启用tungstenitefeature(例如在 wasm 目标上),或需要注入自定义后端时,使用connect_with:

let session = model.responses_websocket().connect_with(&TungsteniteClient::new()).await?;

connect_with接受任意W: WebSocketClientExt,见 crates/rig-core/src/providers/openai/responses_api/websocket.rs#L262-L276。responses_websocket()入口本身不依赖任何特定后端,见 crates/rig-core/src/providers/openai/responses_api/websocket.rs#L870-L878。

会话时间配置

ResponsesWebSocketSessionBuilder提供两类超时配置(见 crates/rig-core/src/providers/openai/responses_api/websocket.rs#L202-L248):

配置项默认值说明
connect_timeout(Duration)30 秒(DEFAULT_CONNECT_TIMEOUT,定义于同文件第 34 行)建立 WebSocket 连接的握手超时
without_connect_timeout()—禁用连接超时
event_timeout(Duration)禁用(None)等待下一条 WebSocket 事件的最大时长
without_event_timeout()—禁用事件超时

注意:连接超时是由后端在握手阶段强制执行的(见ConnectOptions.timeout),与会话层的事件超时相互独立。

Cargo feature 与 TLS 选择

rig-tungstenite自身的 feature(见 crates/rig-tungstenite/Cargo.toml):

Feature默认作用
default = ["rustls"]是默认启用 rustls TLS
rustls是启用tokio-tungstenite/rustls-tls-webpki-roots,基于 webpki-roots 的 rustls 证书栈
native-tls否切换为tokio-tungstenite/native-tls

在rig-core侧,feature 的联动关系(见 crates/rig-core/Cargo.toml):

  • tungstenite = ["websocket", "dep:rig-tungstenite"]:同时拉起传输无关的websocket契约与捆绑后端;
  • rustls = ["rig-reqwest?/rustls", "rig-tungstenite?/rustls"]、native-tls = [...]:为 HTTP 与 WebSocket 两个传输统一选择 TLS 栈,恰好启用其中一个(默认rustls)。

wasm 目标的明确限制

rig-tungstenite是原生后端,tokio-tungstenite需要 Tokio reactor。在 wasm 目标上,该 crate 会在编译期直接报错(见 crates/rig-tungstenite/src/lib.rs#L22-L29):

#[cfg(target_family = "wasm")] compile_error!( "rig-tungstenite is a native websocket backend (tokio-tungstenite). On wasm, implement \ `rig_http::ws_client::WebSocketClientExt` over `web_sys::WebSocket` and open sessions with \ `connect_with(..)`." );

即:浏览器场景应基于web_sys::WebSocket自行实现WebSocketClientExt,再通过connect_with(..)打开会话。rig-http的ws_client契约正是为这种可插拔场景准备的接缝(seam)。这也解释了为什么rig-core将rig-tungstenite依赖放在[target.'cfg(not(target_family = "wasm"))'.dependencies]下,且 feature 层面做了降级处理(见 crates/rig-core/Cargo.toml#L59-L62)。

传输契约详解:Frame、CloseFrame 与 ConnectOptions

WebSocketClientExt与WebSocketConnection两个 trait 是后端实现的唯一入口(定义见 crates/rig-http/src/ws_client.rs):

  • WebSocketClientExt::connect(&self, request: Request<NoBody>, options: ConnectOptions):打开一条连接;升级被拒绝时必须保留 HTTP 状态码、响应头和响应体。
  • WebSocketConnection:提供send(Frame)、recv()(返回Ok(None)表示对端结束流)、close(Option<CloseFrame>)三个方法,以WasmBoxedFuture返回,保证会话的线程安全契约在 wasm 下也可编译。

Frame枚举覆盖 WebSocket 的全部数据与控制帧:

变体含义
Text(String)UTF-8 文本帧
Binary(Bytes)二进制帧
Ping(Bytes)Ping 帧,携带应用载荷
Pong(Bytes)Pong 帧,携带应用载荷
Close(Option<CloseFrame>)关闭帧;Some时携带对端的 RFC 6455 状态码与原因

CloseFrame { code: u16, reason: String }直接对应 RFC 6455 的关闭码与关闭原因。ConnectOptions { timeout: Option<Duration> }只负责握手超时,可通过ConnectOptions::new().with_timeout(...)构造。

值得注意的细节:在 crates/rig-tungstenite/src/connection.rs 的from_message中,Message::Frame(_)(原始帧,不承载会话级协议载荷)会被跳过并返回None,而不是把一段会被会话当作 JSON 解析的裸字节交上去;同时DirectConnection::recv会用循环持续读取,直到拿到一个有效帧(见 crates/rig-tungstenite/src/connection.rs#L40-L55)。connection/tests.rs中的测试逐一验证了 Text/Binary/Ping/Pong/Close 六种帧在 tungstenite 表示间的往返一致性,并确认原始帧被正确丢弃(见 crates/rig-tungstenite/src/connection/tests.rs)。

rig-http还提供了websocket_url(base_url, path)辅助函数,将https/http基地址转换为wss/ws地址并拼接路径(例如https://example.com/v1+responses→wss://example.com/v1/responses),不支持的 scheme 会返回错误(见 crates/rig-http/src/ws_client.rs#L117-L144)。

运行时策略:Direct 与 Forwarded 两条路径

TungsteniteClient::connect的核心逻辑(见 crates/rig-tungstenite/src/lib.rs#L66-L87)会根据调用方是否处于 Tokio 运行时,选择两种截然不同的连接实现:

connect() ├─ 调用方在 Tokio 运行时内 → 握手后返回 DirectConnection(由调用方运行时轮询 socket) └─ 调用方无 Tokio 运行时 → 握手移入 fallback runtime, 返回 ForwardedConnection(channel 转发的 actor 连接)

DirectConnection:调用方自己的 Tokio 运行时

当runtime::in_tokio()为真(即Handle::try_current()成功,见 crates/rig-tungstenite/src/runtime.rs#L37-L39),握手直接在当前运行时上执行,返回的DirectConnection直接包装WebSocketStream<MaybeTlsStream<TcpStream>>,由调用方运行时轮询。注意in_tokio()只能检测是否存在运行时句柄,无法检测 I/O 与 timer driver 是否启用——缺失 driver 时 socket I/O 可能 panic,宿主需要自行保证。

ForwardedConnection:懒启动的 fallback runtime

这是该 crate 最值得一提的设计。当调用方没有 Tokio 运行时(典型场景是 Bevy task pools、smol、futures::executor),连接不会要求调用方提供 reactor,而是:

  1. 通过LazyLock懒启动一个共享的单 worker Tokio 多线程运行时,线程名为"rig-tungstenite",enable_all()启用全部 driver;启动失败会被缓存,后续连接直接复用该失败结果(见 crates/rig-tungstenite/src/runtime.rs#L13-L31)。
  2. 握手与 socket 生命周期全部留在该 fallback runtime 上,调用方与连接之间只通过一对futureschannel 通信(命令通道容量为 1,oneshot用于应答),调用方完全不需要亲自轮询 socket I/O(见 crates/rig-tungstenite/src/connection.rs#L179-L207)。
  3. 连接持有OwnedTask句柄,drop 即 abort:被取消的连接会立刻释放传输资源,即使 actor 正阻塞在对端不读的写操作上(见 crates/rig-tungstenite/src/runtime.rs#L42-L59)。off_runtime_ownership集成测试专门用真实 loopback socket 验证了这一点(见 crates/rig-tungstenite/tests/off_runtime_ownership.rs)。

连接 actor:select 驱动的帧与命令仲裁

run_actor是 channel 转发连接的核心(见 crates/rig-tungstenite/src/connection.rs#L97-L177),它把 socketsplit()成 sink 与 stream 后进入事件循环,并通过futures::select!同时监听命令与入站帧。这一结构带来三个关键行为:

  • 空闲 receive 不会饿死 close:即使没有待处理的命令,actor 也会持续轮询入站帧,因此一个挂起的 receive 不会阻塞后续的 close 命令;
  • 取消的 receive 保留结果:若应答通道的接收端已取消(reply.send失败),未送达的帧或错误会被重新压回队首(inbound.push_front),下一次recv仍能观察到原始结果(见 crates/rig-tungstenite/src/connection.rs#L106-L124);
  • 有界 read-ahead 施加背压:入站队列上限READ_AHEAD = 256(见 crates/rig-tungstenite/src/connection.rs#L94-L95)。队列满或流结束时,actor 暂停读 socket 只处理命令,避免慢消费者导致无界缓冲;同时轮询 socket 本身会驱动 tungstenite 的自动 Pong 回复,维持心跳。

错误处理:升级拒绝与超时的保真

WebSocket 握手被服务端拒绝(如 API key 无效)时,tungstenite 会给出Error::Http(response)。from_tungstenite会把它还原为携带完整 status、headers、body 的Error::non_success_with_details,而不是退化为一条没有上下文的字符串错误(见 crates/rig-tungstenite/src/lib.rs#L133-L148)。这是有实测依据的:src/tests.rs中的测试使用真实端点录制的拒绝体(401+x-request-id+ OpenAI 错误信封 JSON)断言三者完整保留,并验证429时retry-after、x-ratelimit-remaining等限流头原样存活(见 crates/rig-tungstenite/src/tests.rs#L24-L88)。而握手超时则产生专门的ConnectTimeout错误,消息形如"timed out connecting the websocket after 30s",同样被测试固定下来防止意外变更(见 crates/rig-tungstenite/src/tests.rs#L138-L150)。

另一条保真边界是请求头的透传:client_request会把调用方请求(包括认证头)覆盖到 tungstenite 生成的握手请求上,而Sec-WebSocket-Key等握手专用头由后端自己补齐,调用方不应提供。the_client_request_keeps_the_callers_headers测试用wss://api.openai.com/v1/responses的Bearer头验证了这一点(见 crates/rig-tungstenite/src/tests.rs#L154-L174)。

非Http的传输错误(ConnectionClosed、Io、Protocol、Url等)则一律映射为Error::Instance,不携带任何"不存在的 provider 响应"——every_transport_failure_is_left_alone测试枚举了全部此类变体以钉死这条边界(见 crates/rig-tungstenite/src/tests.rs#L100-L132)。会话层随后通过websocket_provider_error从保留的响应头中提取request_id并合入ProviderError(见 crates/rig-core/src/providers/openai/responses_api/websocket.rs#L861-L868)。

无 Tokio 运行时的完整会话验证

rig-core 的集成测试websocket_off_runtime是整个后端能力最直接的证据:它用futures::executor::block_on在没有 Tokio 运行时的调用线程上驱动"连接 → 发送 → 接收 → 关闭"的完整会话(服务端则在独立线程的 Tokio 运行时上接受连接并回应response.create),并验证事件超时在服务端静默时依然生效(见 crates/rig-core/tests/websocket_off_runtime.rs)。此外,streaming_conformance_websocket、websocket_handshake_rejection等测试均以required-features = ["tungstenite"]的方式 gate 在 feature 之后(见 crates/rig-core/Cargo.toml#L120-L134)。

从使用角度总结,接入路径就是两条:默认环境启用rig-core/tungstenite后直接connect();特殊环境(wasm 浏览器、自定义传输、无 Tokio 运行时)则实现WebSocketClientExt后走connect_with()。协议层你完全不必关心——会话协议在 rig-core,Socket 在 rig-tungstenite,两者之间只隔着一层由rig-http定义的薄契约。

  • AI Agent
  • Agent 框架
  • RAG
  • 后端

【免费下载链接】rig

⚙️🦀 Build modular and scalable LLM Applications in Rust

项目地址:https://gitcode.com/GitHub_Trending/rig2/rig
点击查看免费下载
上一篇:Zotero插件市场:在文献管理软件中打造你的专属插件生态系统
下一篇:Zotero插件市场终极指南:一键安装插件,提升文献管理效率

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

返回列表