- AI Agent
- Agent 框架
- RAG
- 后端
【免费下载链接】rig
⚙️🦀 Build modular and scalable LLM Applications in Rust
本篇指南聚焦 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,而是:
- 通过
LazyLock懒启动一个共享的单 worker Tokio 多线程运行时,线程名为"rig-tungstenite",enable_all()启用全部 driver;启动失败会被缓存,后续连接直接复用该失败结果(见 crates/rig-tungstenite/src/runtime.rs#L13-L31)。 - 握手与 socket 生命周期全部留在该 fallback runtime 上,调用方与连接之间只通过一对
futureschannel 通信(命令通道容量为 1,oneshot用于应答),调用方完全不需要亲自轮询 socket I/O(见 crates/rig-tungstenite/src/connection.rs#L179-L207)。 - 连接持有
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
相关推荐
rig-http 传输契约详解:为 Rig LLM 应用构建可替换的 HTTP 与 WebSocket 传输层
rig http 传输契约详解:为 Rig LLM 应用构建可替换的 HTTP 与 WebSocket 传输层 rig http 是 Rig 项目中所有模型提供
AI AgentAgent 框架RAG后端rig-surrealdb 实战指南:在 Rust 中为 Rig 框架构建 SurrealDB 向量检索(RAG)后端
rig surrealdb 实战指南:在 Rust 中为 Rig 框架构建 SurrealDB 向量检索(RAG)后端 导读 rig surrealdb 是 R
AI AgentAgent 框架RAG后端rig-reqwest 深度解析:Rig 框架的 reqwest HTTP 传输层与运行时语义
rig reqwest 深度解析:Rig 框架的 reqwest HTTP 传输层与运行时语义 导读 本文围绕 rig reqwest https://link
AI AgentAgent 框架RAG后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考