ARTICLE · 1130023
Codex 源码-LLM Client 与 Responses API 流式通信
Codex 源码解析系列
第 12 讲:LLM Client 与 Responses API 流式通信
基于 OpenAI Codex 源码 · 2026-10-06
💡 本讲一句话:Codex 每次向模型要答案,背后都跑着一条"WebSocket 优先、HTTP/SSE 兜底"的流式管道——本讲逐行拆这条管道的四个关键件:传输分发(client.rs)、SSE 解析与错误分类(codex-api/sse/responses.rs)、事件桥接与遥测(map_response_events)、以及三级重试状态机(responses_retry.rs),读完你能说清"断线重连、降级 HTTP、限流退避"分别是谁在什么时机做的。
一、stream():一个入口,两种传输
第 10 讲我们看过 Agent Loop 怎么驱动一轮对话,但"把 prompt 发给模型、把 token 流式收回来"这件事发生在更底层。入口就是 ModelClientSession::stream()——它先问 provider 支持哪种 wire API,再决定走 WebSocket 还是 HTTP/SSE。
📄 codex-rs/core/src/client.rs (第 1961-2030 行)
pub async fn stream( // 对外唯一入口:一次模型请求,返回事件流&mut self, // 会话级客户端(持有 WebSocket 连接状态),可变借用因为要改传输状态prompt: &Prompt, // 本轮完整提示词(含历史 + 工具结果)model_info: &ModelInfo, // 模型元信息:slug、是否支持 WS、重试上限等session_telemetry: &SessionTelemetry, // 会话级遥测句柄,后面每个分支都要打点effort: Option<ReasoningEffortConfig>, // 推理强度(可选),透传给模型侧summary: ReasoningSummaryConfig, // 是否要推理摘要service_tier: Option<String>, // 服务等级(如 priority),影响路由与计费responses_metadata: &CodexResponsesMetadata, // turn_id/session_id/thread_id 等元数据inference_trace: &InferenceTraceContext, // rollout trace 上下文,显式传入避免传输分支各写一套) -> Result<ResponseStream> { // 统一返回事件流,调用方不关心底层是 WS 还是 HTTP let wire_api = self.client.state.provider.info().wire_api; // 从 provider 配置读协议类型(Responses/Chat) match wire_api { // 目前只有 Responses 一种走这条路径 WireApi::Responses => { // 进入 Responses API 分支 if self.client.responses_websocket_enabled() { // WS 开关:provider 支持 && 本会话没被降级过 let request_trace = current_span_w3c_trace_context(); // 抓当前 W3C trace 上下文,跨传输传递链路 ID match self // 先尝试 WebSocket 传输 .stream_responses_websocket( // WS 路径:连接复用 + 增量 delta(第五节细讲) prompt, model_info, session_telemetry, effort.clone(), summary, // 参数原样透传 service_tier.clone(), responses_metadata, /*warmup*/ false, request_trace, inference_trace, // warmup=false 表示真实推理请求,要进 trace ) .await? { // 等待结果;? 把不可恢复错误直接抛给调用方 WebsocketStreamOutcome::Stream(stream) => return Ok(stream), // WS 成功:直接返回事件流,HTTP 路径不再执行 WebsocketStreamOutcome::FallbackToHttp => { // WS 被服务端拒绝(426)或重试耗尽 self.try_switch_fallback_transport(session_telemetry, model_info); // 永久降级本会话到 HTTP,并清空 WS 状态 } } // 到这里说明:WS 没开、或刚被降级——统一落到 HTTP/SSE self.stream_responses_api( // HTTP/SSE 路径(第二节细讲) prompt, model_info, session_telemetry, effort, summary, service_tier, responses_metadata, inference_trace, ) .await // 返回同样的 ResponseStream 类型,上层完全无感 }为什么这样设计:把"选哪种传输"收敛在一个函数里,上层(turn loop)只面对一个 ResponseStream——传输是纯实现细节。注意降级不是"这次失败下次再试 WS":force_http_fallback()(第 562-581 行)用 disable_websockets.swap(true) 原子翻转一个全局 AtomicBool,整个会话剩余生命周期都走 HTTP。这是刻意的:WS 失败通常意味着网络环境不支持长连接(代理、防火墙),反复重试只会浪费延迟;一次降级 + 打点 codex.transport.fallback_to_http,比无限抖动更稳。
🔑 关键决策:降级是"会话级一次性"的,不是"请求级重试"。WS 开关 = provider 支持 && !disable_websockets(第 926-934 行),两个条件任一不满足就永远走 HTTP。
二、SSE 解析:超时、断连、错误分类
HTTP 路径的字节流由 codex-api/src/sse/responses.rs 负责解析。先看入口 spawn_response_stream()(第 36-102 行):它把 HTTP 响应头里的"元信息事件"先塞进通道,再启动 SSE 解析任务。
📄 codex-rs/codex-api/src/sse/responses.rs (第 36-102 行)
pub fn spawn_response_stream( // 把"HTTP 响应头 + SSE 字节流"统一包装成事件流 stream_response: StreamResponse, // 已收到的 HTTP 响应(headers + body 字节流) idle_timeout: Duration, // 空闲超时:多久没收到任何 SSE 事件就判死 telemetry: Option<Arc<dyn SseTelemetry>>, // SSE 级遥测(每 poll 一次打点) turn_state: Option<Arc<OnceLock<String>>>, // sticky routing token:服务端回传后锁存,后续请求带上) -> ResponseStream { let rate_limit_snapshots = parse_all_rate_limits(&stream_response.headers); // 从响应头解析限流快照(X-RateLimit-*) let models_etag = stream_response.headers.get("X-Models-Etag")...map(ToString::to_string); // 模型列表 etag,供缓存失效判断 let server_model = stream_response.headers.get(OPENAI_MODEL_HEADER)...map(ToString::to_string); // 服务端实际用的模型(可能与请求不同) let reasoning_included = stream_response.headers.get(X_REASONING_INCLUDED_HEADER).is_some(); // 是否包含推理内容 let upstream_request_id = stream_response.headers.get(REQUEST_ID_HEADER)...map(str::to_string); // x-request-id:排障时对上服务端日志的关键 ID let safety_buffering_treatment = treatment_from_headers(&stream_response.headers).unwrap_or_default(); // 安全缓冲策略(按头解析,缺省默认) if let Some(turn_state) = turn_state.as_ref() // 若调用方提供了 sticky routing 槽位 && let Some(header_value) = stream_response.headers.get(X_CODEX_TURN_STATE_HEADER)... { // 且服务端回传了路由状态头 let _ = turn_state.set(header_value.to_string()); // OnceLock 一次性写入:同一 turn 后续请求都带这个 token,粘住后端实例 } let (tx_event, rx_event) = mpsc::channel::<Result<ResponseEvent, ApiError>>(1600); // 通道容量 1600:够缓冲突发 delta,又不至于内存失控 tokio::spawn(async move { // 独立任务消费字节流,调用方只拿 rx_event——生产者/消费者解耦 if let Some(model) = server_model { tx_event.send(Ok(ResponseEvent::ServerModel(model))).await; } // 先广播"实际模型"事件(可能发生了路由切换) for snapshot in rate_limit_snapshots { tx_event.send(Ok(ResponseEvent::RateLimits(snapshot))).await; } // 再广播限流快照,UI 可提前显示额度 if let Some(etag) = models_etag { tx_event.send(Ok(ResponseEvent::ModelsEtag(etag))).await; } // 模型列表 etag 事件 if reasoning_included { tx_event.send(Ok(ResponseEvent::ServerReasoningIncluded(true))).await; } // 推理内容开关事件 process_sse_with_treatment(stream_response.bytes, tx_event, idle_timeout, telemetry, safety_buffering_treatment).await; // 最后进入真正的 SSE 循环 }); ResponseStream { rx_event, upstream_request_id } // 对外只暴露接收端 + 上游请求 ID}为什么这样设计:响应头里的信息(实际模型、限流额度、路由 token)和 SSE body 是两种来源,但都被归一成 ResponseEvent 从同一个通道流出——下游不需要区分"这是头还是事件"。turn_state 用 OnceLock 锁存是典型的 sticky session 手法:第一轮拿到路由 token,后续请求都带上,让同一 turn 尽量落在同一后端实例上(对 prompt cache 命中率很关键)。
真正的解析循环 process_sse_with_treatment()(第 568-683 行)是这段代码里防御性最强的部分:
📄 codex-rs/codex-api/src/sse/responses.rs (第 568-683 行,节选)
async fn process_sse_with_treatment( // SSE 主循环:字节流 → ResponseEvent stream: ByteStream, tx_event: mpsc::Sender<Result<ResponseEvent, ApiError>>, idle_timeout: Duration, telemetry: Option<Arc<dyn SseTelemetry>>, safety_buffering_treatment: SafetyBufferingTreatment,) { let mut stream = stream.eventsource(); // 字节流升级为 SSE eventsource(按 "data:" 行切事件) let mut response_error: Option<ApiError> = None; // 暂存"服务端声明的错误",等流结束再上报 loop { let start = Instant::now(); // 记录本次 poll 起点,用于遥测耗时 let response = tokio::select! { // select 竞争两个未来:消费者断开 vs 超时读事件 biased; // biased:优先检查"是否还有人收"——没人收就立刻退出,不白等网络 _ = tx_event.closed() => return, // 下游已 drop(用户取消/turn 结束)→ 直接停,省资源 response = timeout(idle_timeout, stream.next()) => response, // idle_timeout 内没事件 → Err(Elapsed):判"流卡死" }; if let Some(t) = telemetry.as_ref() { t.on_sse_poll(&response, start.elapsed()); } // 每次 poll 都打点(含耗时),SSE 健康度可观测 let sse = match response { Ok(Some(Ok(sse))) => sse, // 正常拿到一个 SSE 事件 Ok(Some(Err(e))) => { tx_event.send(Err(ApiError::Stream(e.to_string()))).await; return; } // 传输层错误:上报后结束 Ok(None) => { // 服务端主动关流(EOF) let error = response_error.unwrap_or(ApiError::Stream("stream closed before response.completed".into())); // 若之前收到过 failed 用那个错,否则报"未完成就断了" tx_event.send(Err(error)).await; return; } Err(_) => { tx_event.send(Err(ApiError::Stream("idle timeout waiting for SSE".into()))).await; return; } // 超时:明确区分"卡死"和"正常结束" }; let event: ResponsesStreamEvent = match serde_json::from_str(&sse.data) { // 解析事件 JSON Ok(event) => event, Err(e) => { debug!(...); continue; } // 单条解析失败只记日志跳过——不让一条脏数据杀死整条流 }; let model_verifications = event.model_verifications(); // 提取模型校验信息(若有) let turn_moderation_metadata = event.turn_moderation_metadata(); // 提取本轮审核元数据(若有) let safety_buffering = event.safety_buffering(&safety_buffering_treatment); // 按策略决定是否缓冲安全事件 if let Some(model) = event.response_model() && last_server_model.as_deref() != Some(model.as_str()) { // 模型中途切换(路由变更) tx_event.send(Ok(ResponseEvent::ServerModel(model.clone()))).await...; // 广播新模型事件,UI/日志能感知"换模型了" } if let Some(verifications) = model_verifications { tx_event.send(Ok(ResponseEvent::ModelVerifications(verifications))).await...; } // 校验信息透传 if let Some(metadata) = turn_moderation_metadata { tx_event.send(Ok(ResponseEvent::TurnModerationMetadata(metadata))).await...; } // 审核元数据透传 if let Some(buffering) = safety_buffering { tx_event.send(Ok(ResponseEvent::SafetyBuffering(buffering))).await...; } // 安全缓冲事件透传 match process_responses_event(event) { // 核心分发:按 event.kind 路由(见下方错误分类) Ok(Some(event)) => { let is_completed = matches!(event, ResponseEvent::Completed { .. }); // 判断是否终态事件 if tx_event.send(Ok(event)).await.is_err() { return; } // 发送失败=下游已走,退出循环 if is_completed { return; } // response.completed 之后不再处理任何事件——协议上它就是终点 } Ok(None) => {} // 未识别/忽略的事件(如 in_progress),继续等下一条 Err(error) => { response_error = Some(error.into_api_error()); } // 服务端声明失败:先记下,等流自然结束再上报(保留完整上下文) } }为什么这样设计:三个防御点值得注意。① biased select 把"消费者是否还在"放在最高优先级——用户按 Esc 取消 turn 的瞬间,解析任务立刻退出,而不是傻等到超时。② 单条 JSON 解析失败只 continue:SSE 流里混入一条畸形事件不该杀死整轮对话,这是"尽力而为"的容错哲学。③ response.failed 不立即报错而是暂存到 response_error,等流关闭时统一上报——因为服务端发完 failed 事件后可能还有收尾数据,提前抛错会丢上下文。
错误分类是重试策略的输入。第 711-730 行一组谓词把服务端 error code 映射成语义化错误:
📄 codex-rs/codex-api/src/sse/responses.rs (第 711-730 行)
fn is_context_window_error(error: &Error) -> bool { error.code.as_deref() == Some("context_length_exceeded") } // 上下文超长:不可重试,应触发压缩(第14讲主题)fn is_quota_exceeded_error(error: &Error) -> bool { error.code.as_deref() == Some("insufficient_quota") } // 额度用尽:重试无意义,直接告知用户fn is_usage_not_included(error: &Error) -> bool { error.code.as_deref() == Some("usage_not_included") } // 套餐不含该用量:同上,属于"别试了"类错误fn is_cyber_policy_error(error: &Error) -> bool { error.code.as_deref() == Some("cyber_policy") } // 网络安全策略拦截:需要特殊文案提示fn is_server_overloaded_error(error: &Error) -> bool { // 服务端过载(两种 code) error.code.as_deref() == Some("server_is_overloaded") || error.code.as_deref() == Some("slow_down") // slow_down 是限流变体,同样可退避重试}为什么这样设计:错误被分成"可重试(overloaded/rate_limit)"和"不可重试(quota/context/cyber_policy)"两大类,这个分类直接决定第四节的重试状态机走哪条分支。限流场景还有更细的处理——第 685-709 行 try_parse_retry_after() 用正则从错误消息里抠出 "try again in Ns",把服务端建议的等待时间变成 Duration:尊重服务端的退避建议,比客户端自己猜更准。
context_length_exceeded
映射为:ContextWindowExceeded 可重试:否
客户端动作:标记 token 满,触发压缩流程(第 14 讲主题)
insufficient_quota / usage_not_included
映射为:QuotaExceeded / UsageNotIncluded 可重试:否
客户端动作:更新限流快照,提示用户额度问题
cyber_policy / misalignment_policy_violation
映射为:CyberPolicy / MisalignmentPolicyViolation 可重试:否
客户端动作:展示专用安全提示文案
server_is_overloaded / slow_down
映射为:ServerOverloaded 可重试:是
客户端动作:进入重试状态机,退避后重发
rate_limit_exceeded(含 "try again in Ns")
映射为:RateLimitExceeded { delay } 可重试:是
客户端动作:按服务端建议的 Ns 等待后重试
其他未知 code
映射为:Retryable { message, delay } 可重试:是(保守)
客户端动作:按客户端 backoff 重试,耗尽后降级/报错
三、事件桥接:map_response_events()
SSE/WS 解析出来的 ResponseEvent 还不能直接给 turn loop 用——中间要过一层"桥接任务",负责三件事:遥测打点、TTFT(首 token 延迟)测量、以及把"本次请求新增的所有 item"攒起来供增量复用。这就是 map_response_events()(client.rs 第 2109-2257 行):
📄 codex-rs/core/src/client.rs (第 2109-2257 行,节选)
fn map_response_events<S>( // 泛型 S:任何"产出 ResponseEvent 的流"都能桥接(HTTP/WS 共用) upstream_request_id: Option<String>, api_stream: S, session_telemetry: SessionTelemetry, inference_trace_attempt: InferenceTraceAttempt, provider: SharedModelProvider,) -> (ResponseStream, oneshot::Receiver<LastResponse>) // 返回:事件流 + "最终响应"一次性通道(供增量复用)where S: futures::Stream<Item = std::result::Result<ResponseEvent, ApiError>> + Unpin + Send + 'static { let (tx_event, rx_event) = mpsc::channel::<Result<ResponseEvent>>(RESPONSE_STREAM_CHANNEL_CAPACITY); // 容量 1600(第2083行常量):缓冲突发 delta let (tx_last_response, rx_last_response) = oneshot::channel::<LastResponse>(); // oneshot:只发一次"本次请求的全部新增 item + response_id" let consumer_dropped = CancellationToken::new(); // 消费者(turn loop)drop 时触发,桥接任务立刻停 tokio::spawn(async move { // 独立任务:上游流 → 下游通道 let mut items_added: Vec<ResponseItem> = Vec::new(); // 攒本次请求新增的 item(assistant 消息、工具调用等) let (request_start, mut ttft_ms) = (Instant::now(), None); // TTFT 计时起点 + 首 token 延迟(未测得为 None) loop { let event = tokio::select! { // 又是"消费者优先"的 select _ = consumer_dropped.cancelled() => { inference_trace_attempt.record_cancelled(STREAM_DROPPED_REASON, upstream_request_id, &items_added); return; } // 下游走了:记"被取消"trace(带已收到的 item),退出 event = api_stream.next() => event, // 否则取下一个上游事件 }; let Some(event) = event else { break; }; // 上游流结束 match event { Ok(ResponseEvent::OutputItemDone(item)) => { // 一个完整 item(消息/工具调用)完成 items_added.push(item.clone()); // 先入"本次新增"账本——增量复用靠它 if tx_event.send(Ok(ResponseEvent::OutputItemDone(item))).await.is_err() { ...record_cancelled...; return; } // 再转发下游;发送失败=下游已走,记取消并退出 } Ok(ResponseEvent::Completed { response_id, token_usage, usage_metadata, end_turn }) => { // 终态事件:请求完成 feedback_tags!(last_model_response_id = &response_id); // 打反馈标签(用户点踩时可关联到具体响应) if let Some(usage) = &token_usage { session_telemetry.sse_event_completed(usage, ttft_ms); } // 上报 token 用量 + TTFT:核心性能指标 inference_trace_attempt.record_completed(&response_id, upstream_request_id, &token_usage, &items_added); // rollout trace 记完成(含全部 item,可回放) if let Some(sender) = tx_last_response.take() { sender.send(LastResponse { response_id: response_id.clone(), items_added: std::mem::take(&mut items_added) }); } // 把"response_id + 新增 item"交给会话层——下一轮增量请求的基线 if tx_event.send(Ok(ResponseEvent::Completed { ... })).await.is_err() { return; } // 最后才转发 Completed(账本先落,事件后发) } Ok(event) => { // 其他所有事件(delta、added 等) if matches!(&event, ResponseEvent::OutputItemAdded(_)) && ttft_ms.is_none() { // 第一个 item added = "模型开始产出",测 TTFT ttft_ms = Some(i64::try_from(request_start.elapsed().as_millis()).unwrap_or(i64::MAX)); // 毫秒级精度;溢出兜底 i64::MAX } if tx_event.send(Ok(event)).await.is_err() { ...record_cancelled...; return; } // 透传下游,失败即退出 } Err(err) => { // 上游错误 let mapped = provider.map_api_error(err); // provider 级映射(如 Bedrock 特殊错误码) inference_trace_attempt.record_failed(&mapped, upstream_request_id, &items_added); // trace 记失败(带已收 item,排障用) if !logged_error { session_telemetry.see_event_completed_failed(&mapped); logged_error = true; } // 遥测只报一次,避免重复计数 if tx_event.send(Err(mapped)).await.is_err() { return; } // 错误转发下游(turn loop 的重试逻辑会接手) } }为什么这样设计:这层桥接是"一次请求的完整生命周期账本"。注意 LastResponse(response_id + items_added)走的是独立的 oneshot 通道,而不是混在事件流里——因为它的消费者是下一轮请求的增量逻辑(第五节),和 UI 消费事件流的节奏完全不同。TTFT 用"第一个 OutputItemAdded"而非"第一个 delta"来测:item added 代表模型产出了结构化内容,比纯文本 delta 更能代表"用户看到东西了"的时刻。
💡 设计模式:整条链路是"生产者-消费者 + CancellationToken"的组合:SSE 任务生产 → 桥接任务中转(记账)→ turn loop 消费。任何一环 drop,上游通过 tx_event.closed() / token 取消在下一个 select 点立刻感知——没有忙等、没有泄漏。
四、三级重试状态机:responses_retry.rs
流断了怎么办?答案在 core/src/responses_retry.rs(全文仅 178 行,但把重试策略收得极干净)。状态只有三个字段:
📄 codex-rs/core/src/responses_retry.rs (第 17-47 行)
const INITIAL_CONNECTION_RETRY_DELAY: Duration = Duration::from_secs(5); // 连接级重试初始等待:5 秒const MAX_CONNECTION_RETRY_DELAY: Duration = Duration::from_secs(60); // 指数退避上限:封顶 60 秒,防止无限翻倍#[derive(Debug, Clone, Copy)]pub(crate) enum ResponsesStreamRequest { // 区分两种请求——重试策略不同 Sampling, // 正常采样(用户对话) RemoteCompactionV2, // 远程压缩(第14讲主题),更保守}pub(crate) struct ResponsesStreamRetryState { // 一个 turn 内跨多次尝试共享的状态 retries: u64, // "请求级"重试计数(受 max_retries 约束) connection_retries: u64, // "连接级"重试计数(可无限,见下) connection_retry_delay: Duration, // 当前连接退避时长(5s → 10s → ... → 60s 封顶)}/// Server retry advice retained after stream retries are exhausted. The turn ID/// prevents a reused Guardian session from applying advice from an earlier review. // 注释:耗尽后保留的服务端建议;turn_id 防止跨轮误用pub(crate) struct ExhaustedResponseRetry { // 重试彻底耗尽时的"遗言",存进 thread_extension_data pub(crate) turn_id: String, // 绑定到具体 turn——Guardian 会话复用时不会拿旧 review 的建议当新建议 pub(crate) retry_at: Option<tokio::time::Instant>, // 服务端建议的下次可重试时刻(若有)}核心函数 handle_retryable_response_stream_error()(第 51-144 行)按优先级走三条分支:
📄 codex-rs/core/src/responses_retry.rs (第 92-143 行)
if retry_state.retries >= max_retries // 分支②前置:请求级重试已耗尽 && client_session.try_switch_fallback_transport( // 且降级动作"首次生效"(返回 true) &turn_context.session_telemetry, turn_context.model_info(), ) { // ——WS→HTTP 传输降级,给一次"换条路再试"的机会 sess.send_event(turn_context, EventMsg::Warning(WarningEvent { message: format!("Falling back from WebSockets to HTTPS transport. {err:#}") })).await; // UI 弹 Warning:让用户知道发生了什么(而不是干等) retry_state.retries = 0; // 重置计数——HTTP 路径有自己全新的重试预算 return Ok(()); // Ok(())=调用方继续 loop 重发;Err=彻底放弃 } if retry_state.retries < max_retries { // 分支③:正常请求级重试(预算内) retry_state.retries += 1; // 计数 +1 let delay = err.retry_delay().unwrap_or_else(|| backoff(retry_count)); // 优先用服务端建议的延迟,没有就客户端指数退避 log_retry(request, turn_context, &err, retry_count, max_retries, delay); // warn! 日志:turn_id + 第几次/共几次 + 等多久(排障第一现场) let report_error = retry_count > 1 || cfg!(debug_assertions) || !sess.services.model_client.responses_websocket_enabled(); // release 下隐藏"第 1 次 WS 重试"提示——减少瞬时重连的噪音 if report_error { sess.notify_stream_error(turn_context, format!("Reconnecting... {retry_count}/{max_retries}"), err).await; } // UI 显示进度:用户不再面对"假死"界面 codex_client::record_retry!(retry_count, delay, operation); // 遥测打点(按 Sampling/Compaction 分开统计) tokio::time::sleep(delay).await; // 睡够再返回——退避是硬等待,不是乐观重试 return Ok(()); } sess.services.thread_extension_data.insert(ExhaustedResponseRetry { // 分支④:全部耗尽——留下"遗言"再报错 turn_id: turn_context.sub_id.clone(), retry_at: err.retry_delay().and_then(|delay| tokio::time::Instant::now().checked_add(delay)), }); Err(err) // 返回错误给 turn loop,本轮失败为什么这样设计:三级结构对应三种故障语义:连接级(网络瞬断,可无限等——第 65-90 行在 feature flag UnboundedConnectionRetries 下对 Sampling 请求生效,5s→60s 指数退避);传输级(WS 这条路走不通了,换 HTTP 再给一轮预算);请求级(同一条路上重试 N 次)。把"换传输"放在"耗尽后、报错前"的中间位置很讲究:它既不是第一次失败就降级(太激进),也不是彻底放弃才想起还有 HTTP(太晚)。另外 ExhaustedResponseRetry 带 turn_id 存进 extension data,是防"Guardian 会话复用旧建议"的卫生措施——重试建议必须和产生它的那一轮绑定。
调用方在 session/turn.rs(第 1561-1647 行)的 turn loop 里,重试语义被进一步收紧:
📄 codex-rs/core/src/session/turn.rs (第 1561-1647 行,节选)
let max_retries = turn_context.provider.info().stream_max_retries(); // 重试上限来自 provider 配置(不同后端容忍度不同) let mut retry_state = ResponsesStreamRetryState::default(); // 每个 turn 一份独立状态,跨轮不串味 loop { turn_context.extension_data.remove::<codex_api::ResponseId>(); // 关键卫生:清掉上一轮的 response_id——重试不能把新工具调用归因到旧响应上 let prompt_input = if let Some(input) = initial_input.take() { input } else { sess.clone_history().await.for_prompt(...) }; // 首轮用原始输入;重试轮从历史重建(已执行的工具结果已在历史里) ...build_prompt(...); // 组装完整 prompt let err = match try_run_sampling_request(...).await { // 发起一次采样请求(内部走本讲的 stream()) Ok(output) => return Ok((output, original_input.unwrap_or(prompt.input))), // 成功:直接返回,loop 结束 Err(err) => match err.details() { // 失败:先做"不可重试错误"的快速通道 CodexErrorDetails::ContextWindowExceeded => { sess.set_total_tokens_full(&turn_context).await; return Err(err); } // 上下文满:标记后直接上抛(触发压缩,别浪费重试) CodexErrorDetails::UsageLimitReached(e) => { ...update_rate_limits...; return Err(err); } // 额度尽:更新限流快照后上抛 _ => err, // 其他错误继续走下面的可重试判断 }, }; if !err.is_retryable() { return Err(err); } // 不可重试(quota/cyber_policy 等):立即失败,不进状态机 handle_retryable_response_stream_error(&mut retry_state, max_retries, err, client_session, &sess, &turn_context, ResponsesStreamRequest::Sampling).await?; // 可重试:交给三级状态机(Ok=继续 loop;Err=上抛) turn_context.turn_timing_state.record_sampling_retry(); // 计时器记一次重试(TTFT/总耗时统计要扣除等待时间) }为什么这样设计:turn loop 是"重试的守门人":它先拦截两类重试无意义的错误(上下文满、额度尽),避免状态机白跑;每轮开头 remove::<ResponseId>() 是防归因错乱的保险丝——上一轮的 response_id 如果残留,重试后新产生的工具调用会被错误地挂到旧响应上。重试输入用 clone_history() 重建而非复用旧 prompt:因为上一轮可能已经执行了部分工具,这些结果必须包含在新请求里。
五、WebSocket 连接复用与增量 delta
WS 路径比 HTTP 多一个杀手锏:长会话里只发新增的 item,不重发全量历史。连接复用判断在 websocket_connection()(client.rs 第 1382-1448 行):
📄 codex-rs/core/src/client.rs (第 1396-1440 行,节选)
let needs_new = match self.websocket_session.connection.as_ref() { // 已有连接? Some(conn) => self.websocket_session.endpoint != Some(endpoint) || conn.is_closed().await, // endpoint 变了(换模型/换路由)或连接已断 → 需要新连接 None => true, // 没有连接 → 必须新建 }; let owner_changed = self.websocket_session.auth_owner_generation != auth_owner_generation // 认证属主代际不一致(账号切换中) || self.client.auth_owner_generation() != auth_owner_generation; // ——旧连接属于"上一个主人",不能复用 if owner_changed { self.turn_state = Arc::new(OnceLock::new()); } // 换主人 → sticky routing token 作废重建(新会话不该粘旧路由) if needs_new || owner_changed { // 需要新连接 self.reset_websocket_session(); // 清掉旧状态(last_request/last_response 等增量基线全部作废) let new_conn = match self.client.connect_websocket(...).await { // 发起 WS 握手 Ok(new_conn) => new_conn, Err(err) => { if matches!(err, ApiError::Transport(TransportError::Timeout)) { self.reset_websocket_session(); } return Err(err); } // 超时失败也要清状态,防止半开连接残留 }; self.websocket_session.connection = Some(new_conn); // 记录新连接 self.websocket_session.endpoint = Some(endpoint); // 记住 endpoint——下次复用判断用 self.websocket_session.auth_owner_generation = auth_owner_generation; // 记住属主代际 self.websocket_session.set_connection_reused(/*connection_reused*/ false); // 遥测标记:这是新建(统计连接命中率) } else { self.websocket_session.set_connection_reused(/*connection_reused*/ true); // 复用成功——省掉一次完整握手 + TLS }为什么这样设计:复用的判据是"endpoint 相同 && 连接活着 && 属主没变",三个条件缺一不可。特别是 auth_owner_generation:用户切换账号的瞬间,旧 WS 连接上还挂着上一个账号的凭证和路由状态——复用它会跨账号泄漏上下文,所以代际一变就强制重建 + 重置 turn_state。
增量 delta 的核心校验在 get_incremental_items()(第 1248-1286 行)——注释写得很直白:"只有当非 input 字段没变、且 input 是上一次的严格扩展时才复用增量":
📄 codex-rs/core/src/client.rs (第 1248-1286 行)
fn get_incremental_items( // 判断本次请求能否只发 delta &self, request: &ResponsesApiRequest, last_response: Option<&LastResponse>, allow_empty_delta: bool, ) -> Option<Vec<ResponseItem>> { // Some(delta)=可增量;None=必须全量重发 let previous_request = self.websocket_session.last_request.as_ref()?; // 没有上一轮请求记录 → 无从比较,直接 None if !responses_request_properties_match(previous_request, request) { // model/tools/reasoning 等"非 input 字段"必须完全一致 trace!("incremental request failed, websocket reuse properties didn't match"); return None; // 不一致(如换了工具集)→ 服务端缓存基线失效,全量重发 } let response_items = last_response.map_or(&[][..], |response| response.items_added.as_slice()); // 上一轮服务端返回的 item 也算"基线的一部分"——不能重发给它 let previous_items_len = previous_request.input.len().checked_add(response_items.len()?); // 基线长度 = 上次 input + 上次输出(防溢出用 checked) let Some((request_items_to_compare, incremental_items)) = request.input.split_at_checked(previous_items_len) else { // 本次 input 必须"前缀 == 基线、后缀 = delta" trace!("incremental request failed, incompatible request length"); return None; // 长度不够(历史被压缩/裁剪过)→ 无法对齐,全量重发 }; let previous_items = previous_request.input.iter().chain(response_items); // 拼出完整基线序列 if !previous_items.zip(request_items_to_compare).all(|(previous, current)| response_items_equal_ignoring_internal_metadata(previous, current)) { // 逐条比对前缀(忽略内部元数据差异) trace!("incremental request failed, items didn't match"); return None; // 任何一条不一致 → 基线被破坏,全量重发 } if !allow_empty_delta && incremental_items.is_empty() { return None; } // delta 为空且不允许空增量(无新内容)→ 不发 Some(incremental_items.to_vec()) // 校验通过:返回纯新增部分为什么这样设计:增量协议的正确性完全靠"前缀匹配"保证——服务端按 (previous_response_id + delta) 重建完整上下文,客户端必须证明"我的 input = 你的基线 + 新增"。任何一环对不上(换模型、历史被 compaction 裁剪、工具集变化)都保守地退回全量:宁可多传几 KB,也不能让服务端拼出错误的上下文。这是典型的"优化路径必须可验证,否则走安全路径"。
六、数据流全景:一次模型请求的完整旅程
一次采样请求的完整链路
① 发起:turn loop → stream()
清 ResponseId、重建 prompt,按 WS 开关选路。
▼
② 传输:WS(复用+delta)或 HTTP/SSE
🔹 WS:endpoint/属主校验 → 前缀匹配发增量🔹 HTTP:全量请求 + sticky routing header
▼
③ 解析:SSE/WS → ResponseEvent
idle timeout + biased select;错误分类(可重试/不可重试)。
▼
④ 桥接:map_response_events()
TTFT 计时、items_added 账本、遥测/trace 打点。
▼
⑤ 消费:turn loop 驱动工具执行
事件 → UI/工具调用;LastResponse 留给下一轮增量。
▼
⑥ 失败兜底:三级重试状态机
🔹 连接级无限退避(5s→60s)🔹 WS→HTTP 降级换预算 🔹 请求级 N 次 → ExhaustedResponseRetry
七、小结
这一讲拆的是一条"看起来透明、实际处处设防"的管道:传输选择对上层不可见,但降级是会话级一次性的;SSE 解析容忍单条脏数据但不容忍流卡死(idle timeout);事件桥接把遥测、TTFT、增量账本三件事收进一个任务;重试状态机用三个字段区分"网络抖动/传输失效/请求失败"三种语义。下一讲进入 Context Manager——历史怎么存、怎么裁剪,正好和本讲的"前缀匹配增量协议"接上:compaction 一旦裁剪历史,增量基线立刻作废,全量重发。
📚 系列导航
← 第 11 讲:CodexThread 与 ThreadManager
→ 第 13 讲:Context Manager 与历史管理
关注公众号「AI技术推荐官」获取更多源码解析内容