Skip to content

WebSocket Client生命周期

从代理策略到握手完成,追踪 Codex WebSocket Client 的连接、TLS、收发适配、地址竞争与关闭边界。

基于rust-v0.150.0
CodexRustModelWebSocket

WebSocket Client生命周期 ​

本文承接 HTTP Client路由与中间件 和 Responses重试策略。前文分别解释了 HTTP 路由策略和上层重试预算;本文继续向下追踪 WebSocket 一次连接从构造到关闭的真实生命周期。

范围固定在 codex-websocket-client:它接收一个 Tungstenite Request,从 HttpClientFactory 获取代理路线,建立直连或代理 TCP,按需要完成代理 TLS 和目标 TLS,执行 WebSocket 握手,最后把底层连接包装成同时实现 Stream 与 Sink 的 WebSocketConnection。它不保存会话消息,也不实现重连循环;连接断开后是否重新创建 connector、是否重放请求,由上层 ModelClient 决定。

读者读完后应能解释四个容易混淆的现象:为什么 wss 的 TLS 可能发生两次,为什么 NO_PROXY 会影响 WebSocket 但不改变目标 URL,为什么 IPv6 首选地址卡住时连接仍能快速尝试 IPv4,以及为什么 WebSocketConnection 的 poll_close 只是把关闭动作转交给当前底层流。

1. 生命周期总览 ​

一次连接经过四个 owner:factory 决定路线,connector 保存 TLS 与 TCP 选项,dialer 消费路线并建立 socket,connection 对上层暴露协议流。握手成功后,connector 不再拥有消息级状态。

这里的“生命周期”不是一个后台 task 状态机。连接建立前,异步 future 由调用方持有;握手成功后,WebSocketConnection 由调用方轮询;future 被取消或 sink 被关闭时,底层流的 drop/close 语义负责释放网络资源。

2. Connector配置 ​

WebSocketConnector 是跨连接复用的配置 owner。它保存 factory 的 clone、可选的显式 rustls 配置和 TCP_NODELAY 选择。new 会立即读取 native roots 与 Codex custom CA;TungsteniteDefault 则推迟 TLS 配置交给 Tungstenite。

源码位置:codex-rs/websocket-client/src/lib.rs :: WebSocketConnector、WebSocketTlsMode

rust
#[derive(Clone)]
pub struct WebSocketConnector {
    http_client_factory: HttpClientFactory,
    tls_config: Option<Arc<ClientConfig>>,
    tcp_nodelay: TcpNodelay,
}

#[derive(Clone, Copy, PartialEq, Eq)]
pub enum WebSocketTlsMode {
    ExplicitCodexTls,
    TungsteniteDefault,
}

impl WebSocketConnector {
    pub fn new(
        http_client_factory: &HttpClientFactory,
    ) -> Result<Self, BuildCustomCaTransportError> {
        Self::new_with_tls_mode(http_client_factory, WebSocketTlsMode::ExplicitCodexTls)
    }

    pub fn new_with_tls_mode(
        http_client_factory: &HttpClientFactory,
        tls_mode: WebSocketTlsMode,
    ) -> Result<Self, BuildCustomCaTransportError> {
        let tls_config = match tls_mode {
            WebSocketTlsMode::ExplicitCodexTls => {
                Some(build_rustls_client_config_with_custom_ca()?)
            }
            WebSocketTlsMode::TungsteniteDefault => None,
        };
        Ok(Self {
            http_client_factory: http_client_factory.clone(),
            tls_config,
            tcp_nodelay: TcpNodelay::Default,
        })
    }

    pub fn with_tcp_nodelay(mut self) -> Self {
        self.tcp_nodelay = TcpNodelay::Enabled;
        self
    }
}

这段代码说明了两个生效时机。custom CA 在 new 时就可能失败,属于 connector 构造错误;TCP_NODELAY 只是保存选项,在真正 connect_tcp 成功后才调用 set_nodelay。因此“能创建 connector”不等于“目标已经可达”。

3. 连接入口 ​

connect 是公开连接入口。它不会自行猜测代理,也不会从环境变量直接读取代理;目标 URI 先交给 factory 的异步解析,再把解析结果和 request/config/TLS 选项交给内部 dialer。

源码位置:codex-rs/websocket-client/src/lib.rs :: WebSocketConnector::connect

rust
pub async fn connect(
    &self,
    request: Request,
    config: WebSocketConfig,
) -> Result<(WebSocketConnection, Response), WebSocketError> {
    let proxy_route = self
        .http_client_factory
        .resolve_proxy_route_async(request.uri().to_string())
        .await
        .map_err(WebSocketError::Io)?;
    let (inner, response) = dialer::connect(
        request,
        config,
        self.tls_config.clone(),
        proxy_route,
        self.tcp_nodelay,
    )
    .await?;
    Ok((WebSocketConnection { inner }, response))
}

如果路线解析失败,错误在握手前以 WebSocketError::Io 返回;如果路线成功但 TCP、代理 CONNECT、TLS 或 WebSocket HTTP upgrade 失败,则由 dialer 返回 Tungstenite 错误。这里没有自动重试,所以调用方可以区分“尚未建立连接”和“握手后流读取失败”。

4. 三种路线 ​

dialer 首先按 OutboundProxyRoute 分支。TransportDefault 完全使用 Tungstenite 的默认代理解析,并直接返回 ConnectionInner::TransportDefault;Direct 直接连接目标 host/port,但握手完成后包装为 ConnectionInner::Routed;Proxy 通常先连接代理并发送 CONNECT。Proxy 携带 no_proxy 时会先交给 Tungstenite 判断:成功时也可能返回 TransportDefault,只有 HTTPS proxy scheme 不被其环境代理解析支持时,才回到 Codex 的显式代理路径并包装为 Routed。

源码位置:codex-rs/websocket-client/src/dialer.rs :: connect

rust
pub(crate) async fn connect(
    request: Request,
    config: WebSocketConfig,
    tls_config: Option<Arc<ClientConfig>>,
    proxy_route: OutboundProxyRoute,
    tcp_nodelay: TcpNodelay,
) -> Result<(ConnectionInner, Response), WebSocketError> {
    let disable_nagle = tcp_nodelay == TcpNodelay::Enabled;
    let proxy_url = match proxy_route {
        OutboundProxyRoute::TransportDefault => {
            let (stream, response) = connect_async_tls_with_config(
                request,
                Some(config),
                disable_nagle,
                tls_config.map(Connector::Rustls),
            )
            .await?;
            return Ok((ConnectionInner::TransportDefault(stream), response));
        }
        OutboundProxyRoute::Direct => None,
        OutboundProxyRoute::Proxy { url, no_proxy: None } => Some(url),
        OutboundProxyRoute::Proxy { url, no_proxy: Some(_) } => {
            match connect_async_tls_with_config(
                request.clone(),
                Some(config),
                disable_nagle,
                tls_config.clone().map(Connector::Rustls),
            )
            .await
            {
                Ok((stream, response)) => {
                    return Ok((ConnectionInner::TransportDefault(stream), response));
                }
                Err(WebSocketError::Url(UrlError::UnsupportedProxyScheme)) => Some(url),
                Err(error) => return Err(error),
            }
        }
    };

    let stream: Box<dyn AsyncIo> = match proxy_url {
        None => {
            let host = websocket_host(&request)?;
            let port = websocket_port(&request)?;
            Box::new(
                connect_tcp(host_port(host, port), tcp_nodelay)
                    .await
                    .map_err(WebSocketError::Io)?,
            )
        }
        Some(url) => {
            let proxy = ProxyEndpoint::parse(&url)?;
            let host = websocket_host(&request)?;
            let port = websocket_port(&request)?;
            let stream = connect_tcp(proxy.config.authority(), tcp_nodelay)
                .await
                .map_err(WebSocketError::Io)?;
            connect_via_proxy(stream, &proxy.config, host, port).await?
        }
    };

    let (stream, response) = client_async_tls_with_config(
        request,
        stream,
        Some(config),
        tls_config.map(Connector::Rustls),
    )
    .await?;
    Ok((ConnectionInner::Routed(stream), response))
}

代码中的 no_proxy: Some(_) 分支很有教学价值:它先尝试 Tungstenite 自己的完整 NO_PROXY 语义;如果环境代理 URL 使用 Tungstenite 不支持的 HTTPS proxy scheme,才回退到 Codex 显式 TLS-to-proxy 路径。不能把这个分支简化成“有 no_proxy 就直连”。

5. TLS分层 ​

wss 并不只意味着“调用一次 TLS”。直连时只有目标 TLS;HTTPS proxy 场景先对 proxy authority 建 TLS,再在 CONNECT 建立的字节隧道中对目标 host 建第二层 TLS。tls_config 是目标和代理 TLS 可共享的 rustls 配置;如果 connector 选择 Tungstenite 默认模式,只有显式 HTTPS proxy 需要临时构造 Codex TLS 配置。

源码位置:codex-rs/websocket-client/src/dialer.rs :: HTTPS proxy 分支

rust
let stream: Box<dyn AsyncIo> = if proxy.tls {
    let proxy_tls_config = match &tls_config {
        Some(tls_config) => Arc::clone(tls_config),
        None => build_rustls_client_config_with_custom_ca()
            .map_err(|error| WebSocketError::Io(error.into()))?,
    };
    let server_name = ServerName::try_from(proxy.config.host.clone())
        .map_err(|_| WebSocketError::Tls(TlsError::InvalidDnsName))?;
    let stream = TlsConnector::from(proxy_tls_config)
        .connect(server_name, stream)
        .await
        .map_err(WebSocketError::Io)?;
    Box::new(stream)
} else {
    Box::new(stream)
};
connect_via_proxy(stream, &proxy.config, host, port).await?

这里的 server name 是代理主机,不是最终 WebSocket 主机;目标 server name 则由 client_async_tls_with_config 根据 request 再处理。把两个名字混用会导致企业 HTTPS proxy 或目标证书验证失败。

6. 地址竞争 ​

TCP 连接不是简单的 TcpStream::connect。connect_tcp 先解析所有地址,再由 connect_happy_eyeballs 交错排列首选地址族和备用地址族:第一次立即启动,后续尝试每隔 250ms 推进;任一连接成功就返回,全部失败才返回最后一个错误。

源码位置:codex-rs/websocket-client/src/dialer.rs :: connect_tcp、connect_happy_eyeballs

rust
const HAPPY_EYEBALLS_DELAY: Duration = Duration::from_millis(250);

async fn connect_tcp(address: String, tcp_nodelay: TcpNodelay) -> io::Result<TcpStream> {
    let addresses = tokio::net::lookup_host(address).await?.collect::<Vec<_>>();
    let stream = connect_happy_eyeballs(addresses, TcpStream::connect).await?;
    if tcp_nodelay == TcpNodelay::Enabled {
        stream.set_nodelay(/*nodelay*/ true)?;
    }
    Ok(stream)
}

async fn connect_happy_eyeballs<T, F, Fut>(
    addresses: Vec<SocketAddr>,
    mut connect: F,
) -> io::Result<T>
where
    F: FnMut(SocketAddr) -> Fut,
    Fut: Future<Output = io::Result<T>>,
{
    let mut addresses = addresses.into_iter();
    let Some(first_address) = addresses.next() else {
        return Err(io::Error::new(
            io::ErrorKind::InvalidInput,
            "could not resolve to any address",
        ));
    };

    let first_is_ipv4 = first_address.is_ipv4();
    let mut preferred = VecDeque::new();
    let mut alternate = VecDeque::new();
    for address in addresses {
        if address.is_ipv4() == first_is_ipv4 {
            preferred.push_back(address);
        } else {
            alternate.push_back(address);
        }
    }

    let mut attempts = FuturesUnordered::new();
    attempts.push(connect(first_address));
    let mut next_attempt_at = Instant::now() + HAPPY_EYEBALLS_DELAY;
    let mut last_error = None;

    loop {
        if addresses.is_empty() {
            match attempts.next().await {
                Some(Ok(stream)) => return Ok(stream),
                Some(Err(error)) => {
                    if attempts.is_empty() {
                        return Err(error);
                    }
                    last_error = Some(error);
                }
                None => {
                    return Err(last_error.unwrap_or_else(|| {
                        io::Error::other("connection attempts ended without an error")
                    }));
                }
            }
            continue;
        }

        tokio::select! {
            result = attempts.next() => {
                match result {
                    Some(Ok(stream)) => return Ok(stream),
                    Some(Err(error)) => {
                        last_error = Some(error);
                        let address = take_next_address(&mut addresses)?;
                        attempts.push(connect(address));
                        next_attempt_at = Instant::now() + HAPPY_EYEBALLS_DELAY;
                    }
                    None => {
                        let address = take_next_address(&mut addresses)?;
                        attempts.push(connect(address));
                        next_attempt_at = Instant::now() + HAPPY_EYEBALLS_DELAY;
                    }
                }
            }
            _ = sleep_until(next_attempt_at) => {
                let address = take_next_address(&mut addresses)?;
                attempts.push(connect(address));
                next_attempt_at = Instant::now() + HAPPY_EYEBALLS_DELAY;
            }
        }
    }
}

fn take_next_address(addresses: &mut VecDeque<SocketAddr>) -> io::Result<SocketAddr> {
    addresses
        .pop_front()
        .ok_or_else(|| io::Error::other("connection address queue unexpectedly empty"))
}

地址队列耗尽后,函数不会取消已经发起的 attempts,而是继续等待它们完成;某个 attempt 提前失败时又会立即补发下一地址,不必等待 250ms timer。关键结论仍需由测试验证,不能仅凭延迟常量推断所有网络都能改善。

7. Stream与Sink ​

握手完成后,WebSocketConnection 不再暴露 ConnectionInner 的具体 transport 类型。它对两种 inner 分支分别转发 Stream::poll_next 和 Sink 的四个 poll 方法,因此 API client 可以用 Message 读写,而无需知道当前是直连、默认代理还是显式代理。

源码位置:codex-rs/websocket-client/src/lib.rs :: WebSocketConnection 的 Stream/Sink 实现

rust
pub struct WebSocketConnection {
    inner: ConnectionInner,
}

impl Stream for WebSocketConnection {
    type Item = Result<Message, WebSocketError>;

    fn poll_next(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
        match &mut self.get_mut().inner {
            ConnectionInner::TransportDefault(stream) => Pin::new(stream).poll_next(context),
            ConnectionInner::Routed(stream) => Pin::new(stream).poll_next(context),
        }
    }
}

impl Sink<Message> for WebSocketConnection {
    type Error = WebSocketError;

    fn poll_ready(
        self: Pin<&mut Self>,
        context: &mut Context<'_>,
    ) -> Poll<Result<(), Self::Error>> {
        match &mut self.get_mut().inner {
            ConnectionInner::TransportDefault(stream) => Pin::new(stream).poll_ready(context),
            ConnectionInner::Routed(stream) => Pin::new(stream).poll_ready(context),
        }
    }

    fn poll_close(
        self: Pin<&mut Self>,
        context: &mut Context<'_>,
    ) -> Poll<Result<(), Self::Error>> {
        match &mut self.get_mut().inner {
            ConnectionInner::TransportDefault(stream) => Pin::new(stream).poll_close(context),
            ConnectionInner::Routed(stream) => Pin::new(stream).poll_close(context),
        }
    }
}

poll_ready、start_send、poll_flush 和 poll_close 都继续委托给 Tungstenite;connector 不插入心跳 task,也不把 Close frame 转成额外的应用事件。上层若要实现 ping、响应超时或重连,必须在这个 stream/sink 之上另建 owner。

这张关系图强调的是所有权而不是协议字段:connector 在握手前拥有路线和 TLS 配置,connection 在握手后拥有具体 inner;factory 只被 clone 用来解析路线,不会成为消息流的 owner。

8. 失败与关闭 ​

连接失败不是一个单一状态:空 DNS 结果、无效 proxy URL、代理 TLS 失败、目标 TLS 失败、握手拒绝和消息流错误分别在不同阶段产生。连接建立后,Stream 返回 None 或错误表示读取侧结束;主动关闭则通过 Sink::poll_close 进入 Tungstenite 的关闭协议。

当前 crate 没有“断线后自动回到 Configured 并重试”的边。它只返回错误或结束流;恢复路径必须由调用方重新创建 request/connector 或交给上层 retry policy。这个边界也意味着一个已经使用过的 request body、previous response id 或模型 turn 状态不能由本 crate 自动判断是否安全重放。

9. 连接测试 ​

9.1 公开连接与收发 ​

public_connector_uses_factory_and_exposes_stream_and_sink 启动本地 echo WebSocket server,使用 ReqwestDefault factory 创建 connector,发送 Message::Text("hello"),再断言收到相同消息。它证明公开 connector 能把 factory、握手和 Stream/Sink 接口接起来,不证明真实代理或远端证书链。

源码位置:codex-rs/websocket-client/src/dialer_tests.rs :: public_connector_uses_factory_and_exposes_stream_and_sink

rust
#[tokio::test]
async fn public_connector_uses_factory_and_exposes_stream_and_sink() {
    let (target_addr, target_task) = start_echo_websocket_server(/*acceptor*/ None).await;
    let request = format!("ws://localhost:{}/v1/responses", target_addr.port())
        .into_client_request()
        .expect("websocket request should build");
    let factory = HttpClientFactory::new(OutboundProxyPolicy::ReqwestDefault);
    let connector = WebSocketConnector::new(&factory).expect("connector should build");

    let (mut websocket, _) = connector
        .connect(request, WebSocketConfig::default())
        .await
        .expect("websocket handshake should succeed");
    assert!(!websocket_tcp_nodelay(&websocket));
    let expected = Message::Text("hello".into());
    websocket
        .send(expected.clone())
        .await
        .expect("websocket should send");
    let actual = websocket
        .next()
        .await
        .expect("websocket should receive a message")
        .expect("websocket message should be valid");
    assert_eq!(actual, expected);

    target_task.await.expect("target task should finish");
}

9.2 代理与TLS ​

http_proxy_tunnels_secure_websocket_before_handshake 和 https_proxy_tunnels_secure_websocket_before_handshake 都调用 assert_proxy_tunnels_secure_websocket,测试代理先收到目标 authority 的 CONNECT,再让目标 TLS/WebSocket 握手完成。它证明代理 tunnel 顺序和两类 proxy scheme 的本地 fixture 行为,不证明 PAC/WPAD 或真实企业代理认证。

9.3 NO_PROXY边界 ​

environment_proxy_route_honors_no_proxy_in_a_subprocess 分别用匹配的 127.0.0.1 与不匹配的 unrelated.example 运行子进程。测试通过 proxy listener 是否收到 CONNECT 来断言 bypass 与 proxy 两条路径;子进程用于隔离环境变量,避免污染当前测试进程。

9.4 地址竞争 ​

happy_eyeballs_does_not_wait_for_stalled_preferred_family 使用 paused Tokio time,让首选 IPv6 future 永久 pending,同时让 IPv4 地址可成功;断言在一秒预算内返回 reachable 地址。这证明备用地址会按延迟启动,不证明真实 DNS 返回顺序或所有内核网络栈的时序。

10. 实践验证 ​

可以运行以下测试,观察配置、路线与连接行为之间的关系:

bash
cargo test -p codex-websocket-client public_connector_uses_factory_and_exposes_stream_and_sink
cargo test -p codex-websocket-client http_proxy_tunnels_secure_websocket_before_handshake
cargo test -p codex-websocket-client https_proxy_tunnels_secure_websocket_before_handshake
cargo test -p codex-websocket-client happy_eyeballs_does_not_wait_for_stalled_preferred_family

只读定位时,先从 WebSocketConnector::connect 进入,再沿 dialer::connect、connect_tcp 和 WebSocketConnection::poll_close 搜索。若现象是“无法连接”,先区分 route/TCP/proxy/TLS/upgrade 阶段;若现象是“连接成功但没有消息”,应继续检查上层是否轮询 Stream,而不是在本 crate 中寻找不存在的后台接收 task。

11. 技术边界 ​

这些测试覆盖 WebSocket transport 的一次连接生命周期、代理和 TLS 分层、地址竞争、Stream/Sink 转发及关闭边界。它们没有覆盖上层 Responses WebSocket 如何编码事件、如何发送 ping、如何重连或如何把断线映射为 sampling retry;这些行为属于 codex-api、core 或更高层 client。下一步阅读 Realtime 专题时,还需要重新核对 feature gate、产品入口和音频/会话状态 owner,不能把本 crate 的通用 WebSocket 连接器等同于 Realtime 会话实现。