Skip to content

Responses流解析

从 HTTP 响应、SSE 解帧和事件分类追踪到 Core 消费,解释完成、断流、超时、错误保留与取消清理。

基于rust-v0.150.0
CodexRustModelSSE

Responses流解析 ​

本文承接 Codex API类型 和 ModelClient结构。前者已经解释 ResponseEvent 的类型边界,后者建立了 Session 级 ModelClient 与 Turn 级 ModelClientSession 的所有权;本文回答更具体的问题:HTTP Responses 请求成功建立后,字节流怎样变成 Core 能消费的事件,服务端失败、连接断开和用户取消又在哪一层 结束这条管道。

这里的 SSE 是 Server-Sent Events:HTTP response body 由若干 event:、data: 行和空行分隔的帧 组成。本文不分析请求字段如何构造,也不展开每一种 text、reasoning、tool delta 如何合并成最终 item;这些 属于 模型请求构造 和后续增量专题。读完本文,读者应能从 ResponsesClient::stream_request 追到 run_sampling_request,并判断一个异常究竟发生在 HTTP、SSE、 JSON 事件、业务错误映射还是消费者取消层。

1. 管道分层 ​

一条 HTTP Responses 流经过两级有界 channel。codex-api 拥有网络字节流和 SSE 解析任务,Core 的 map_response_events 拥有遥测、provider 错误转换和已完成 item 列表;Turn 只消费最终 core::ResponseStream。两级都使用容量 1600,但它们服务不同边界,不能视为同一个队列。

层次所有者输入输出终止依据
HTTP endpointResponsesClientrequest、headers、compressionStreamResponse请求失败或取得 response
SSE parserspawn_response_stream 派生任务headers、ByteStreamcodex_api::ResponseStreamCompleted、错误、EOF、idle timeout、接收者关闭
Core mappermap_response_events 派生任务API ResponseEventCore ResponseStreamconsumer drop 或 API stream 结束
Turn consumerrun_sampling_requestCore event streamruntime event、token 状态、结果Completed、错误或 Turn 取消

这层划分解释了一个容易混淆的现象:HTTP status 为成功只表示 body 可以开始读取,不表示模型 Turn 已完成; 反过来,收到 response.completed 后 parser 可以立即结束,不需要等待服务器关闭 HTTP 连接。

SSE 与 WebSocket 共享 process_responses_event,但错误终止策略不同:SSE parser 会暂存已分类的 response.failed,直到 EOF 再把它发送给下游;WebSocket reader 在 process_responses_event 返回错误时立即 结束当前读取循环。两者都要求 response.completed 才算成功,但“错误何时可观察”由传输实现决定。

这张时序图只描述 API 层,不把 Core mapper 或 Turn 的取消混进来;那两个消费者边界在后文单独展开。

2. HTTP入口 ​

Core 的 HTTP 分支先构造 API client,再调用 stream_request。只有建立流之前的 401 会进入认证恢复循环; 流已经返回后出现的 SSE 错误不会回到这个 loop 重新发送请求。

源码位置:codex-rs/core/src/client.rs :: ModelClientSession::stream_responses_api 返回与认证分支

rust
match stream_result {
    Ok(stream) => {
        let (stream, _) = map_response_stream(
            stream,
            request_session_telemetry,
            inference_trace_attempt,
            Arc::clone(&self.client.state.provider),
        );
        return Ok(stream);
    }
    Err(ApiError::Transport(
        unauthorized_transport @ TransportError::Http { status, .. },
    )) if status == StatusCode::UNAUTHORIZED => {
        let response_debug_context =
            extract_response_debug_context(&unauthorized_transport);
        inference_trace_attempt.record_failed(
            &unauthorized_transport,
            response_debug_context.request_id.as_deref(),
            /*output_items*/ &[],
        );
        pending_retry = PendingUnauthorizedRetry::from_recovery(
            handle_unauthorized(
                unauthorized_transport,
                &mut auth_recovery,
                session_telemetry,
                &self.client.state.provider,
            )
            .await?,
        );
        continue;
    }
    Err(err) => {
        let response_debug_context =
            extract_response_debug_context_from_api_error(&err);
        let err = self.client.state.provider.map_api_error(err);
        inference_trace_attempt.record_failed(
            &err,
            response_debug_context.request_id.as_deref(),
            /*output_items*/ &[],
        );
        return Err(err);
    }
}

这里的重试边界是“取得 ResponseStream 之前”。这是正确的幂等性边界:流中可能已经产生文本或工具调用, 若 mapper 收到一半事件后自动重发完整请求,消费者可能看到重复副作用。流内错误由上层 Turn retry 策略结合 已发生状态处理,而不是在 endpoint 内静默重播。

codex-api endpoint 负责序列化请求、补 session headers,并明确声明接收 SSE。底层 stream_encoded_json_with 返回 headers 与异步 bytes,而不是先缓冲完整响应体。

源码位置:codex-rs/codex-api/src/endpoint/responses.rs :: ResponsesClient::stream_encoded

rust
async fn stream_encoded(
    &self,
    body: EncodedJsonBody,
    extra_headers: HeaderMap,
    compression: Compression,
    turn_state: Option<Arc<OnceLock<String>>>,
) -> Result<ResponseStream, ApiError> {
    let request_compression = match compression {
        Compression::None => RequestCompression::None,
        Compression::Zstd => RequestCompression::Zstd,
    };

    let stream_response = self
        .session
        .stream_encoded_json_with(
            Method::POST,
            Self::path(),
            extra_headers,
            Some(body),
            |req| {
                req.headers.insert(
                    http::header::ACCEPT,
                    HeaderValue::from_static("text/event-stream"),
                );
                req.compression = request_compression;
            },
        )
        .await?;

    Ok(spawn_response_stream(
        stream_response,
        self.session.provider().stream_idle_timeout,
        self.sse_telemetry.clone(),
        turn_state,
    ))
}

stream_idle_timeout 不是整次 Turn 的总时限。parser 每次等待下一帧时重新计时,所以持续到达事件的长响应 可以运行超过该时长;只有相邻事件之间长时间无数据才触发 idle timeout。

3. 响应头事件 ​

SSE body 不是所有响应事实的来源。spawn_response_stream 在解析首个帧之前读取 HTTP headers,保存 x-request-id,并把 model、rate limits、models etag 与 reasoning 标志转换成先行事件。sticky routing 的 x-codex-turn-state 则写入传入的 OnceLock,不会作为普通 ResponseEvent 发送。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: spawn_response_stream

rust
pub fn spawn_response_stream(
    stream_response: StreamResponse,
    idle_timeout: Duration,
    telemetry: Option<Arc<dyn SseTelemetry>>,
    turn_state: Option<Arc<OnceLock<String>>>,
) -> ResponseStream {
    let rate_limit_snapshots = parse_all_rate_limits(&stream_response.headers);
    let models_etag = stream_response
        .headers
        .get("X-Models-Etag")
        .and_then(|v| v.to_str().ok())
        .map(ToString::to_string);
    let server_model = stream_response
        .headers
        .get(OPENAI_MODEL_HEADER)
        .and_then(|v| v.to_str().ok())
        .map(ToString::to_string);
    let upstream_request_id = stream_response
        .headers
        .get(REQUEST_ID_HEADER)
        .and_then(|value| value.to_str().ok())
        .map(str::to_string);
    let safety_buffering_treatment =
        treatment_from_headers(&stream_response.headers).unwrap_or_default();

    if let Some(turn_state) = turn_state.as_ref()
        && let Some(header_value) = stream_response
            .headers
            .get(X_CODEX_TURN_STATE_HEADER)
            .and_then(|value| value.to_str().ok())
    {
        let _ = turn_state.set(header_value.to_string());
    }

    let (tx_event, rx_event) =
        mpsc::channel::<Result<ResponseEvent, ApiError>>(1600);
    tokio::spawn(async move {
        if let Some(model) = server_model {
            let _ = tx_event.send(Ok(ResponseEvent::ServerModel(model))).await;
        }
        for snapshot in rate_limit_snapshots {
            let _ = tx_event.send(Ok(ResponseEvent::RateLimits(snapshot))).await;
        }
        if let Some(etag) = models_etag {
            let _ = tx_event.send(Ok(ResponseEvent::ModelsEtag(etag))).await;
        }
        process_sse_with_treatment(
            stream_response.bytes,
            tx_event,
            idle_timeout,
            telemetry,
            safety_buffering_treatment,
        )
        .await;
    });

    ResponseStream { rx_event, upstream_request_id }
}

x-request-id 留在 ResponseStream 字段中,是因为它服务错误诊断、trace 和 feedback tag,不是模型内容。 而 ServerModel 必须进入事件流,因为安全路由可能让实际模型不同于请求模型,Turn 消费者需要更新对外状态。

4. 解帧循环 ​

核心循环每次只等待一个 SSE frame。eventsource_stream 负责跨网络 chunk 拼接 SSE 行;Codex 从 sse.data 反序列化 ResponsesStreamEvent,提取 metadata 类旁路事件,再把业务事件交给 process_responses_event。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: process_sse_with_treatment

rust
let mut stream = stream.eventsource();
let mut response_error: Option<ApiError> = None;
let mut last_server_model: Option<String> = None;

loop {
    let start = Instant::now();
    let response = timeout(idle_timeout, stream.next()).await;
    if let Some(t) = telemetry.as_ref() {
        t.on_sse_poll(&response, start.elapsed());
    }

    let sse = match response {
        Ok(Some(Ok(sse))) => sse,
        Ok(Some(Err(e))) => {
            let _ = tx_event.send(Err(ApiError::Stream(e.to_string()))).await;
            return;
        }
        Ok(None) => {
            let error = response_error.unwrap_or(ApiError::Stream(
                "stream closed before response.completed".into(),
            ));
            let _ = tx_event.send(Err(error)).await;
            return;
        }
        Err(_) => {
            let _ = tx_event
                .send(Err(ApiError::Stream("idle timeout waiting for SSE".into())))
                .await;
            return;
        }
    };

    let event: ResponsesStreamEvent = match serde_json::from_str(&sse.data) {
        Ok(event) => event,
        Err(e) => {
            debug!("Failed to parse SSE event: {e}, data: {}", &sse.data);
            continue;
        }
    };

    match process_responses_event(event) {
        Ok(Some(event)) => {
            let is_completed = matches!(event, ResponseEvent::Completed { .. });
            if tx_event.send(Ok(event)).await.is_err() {
                return;
            }
            if is_completed {
                return;
            }
        }
        Ok(None) => {}
        Err(error) => response_error = Some(error.into_api_error()),
    }
}

这段代码有三种看似相近、实际不同的失败策略:

输入异常当前处理原因
SSE framing/transport error立即发送 ApiError::Stream 并结束字节边界已不可信,无法继续解帧
单帧 data 不是合法 JSON记录 debug 后跳过后续 SSE frame 仍可能独立有效
合法 response.failed暂存分类后的 response_error等待 stream 关闭,以保留随后可能到达的旁路事件

第三点尤其重要:process_responses_event 返回错误时循环没有立即退出,而是把它存入 response_error。如果后续没有 response.completed,EOF 会优先返回这个具体错误,而不是用“提前断流”覆盖 服务端已经给出的错误原因。

说明: 非法 JSON 帧被跳过属于当前实现的 fail-open 选择,但它不代表整次请求一定成功。若被跳过的正是 response.completed,随后 EOF 仍会产生 stream closed before response.completed。

5. 事件分类 ​

ResponsesStreamEvent 是宽松的 wire DTO:除 type 外,大多数字段都是 Option。这让不同事件共用一个 反序列化结构,但也意味着每个 match 分支必须重新验证该事件所需字段。字段缺失通常返回 Ok(None),不会 伪造空 delta 或残缺 item。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: ResponsesStreamEvent

rust
#[derive(Deserialize, Debug)]
pub struct ResponsesStreamEvent {
    #[serde(rename = "type")]
    pub(crate) kind: String,
    pub(crate) headers: Option<Value>,
    metadata: Option<Value>,
    response: Option<Value>,
    item: Option<Value>,
    item_id: Option<String>,
    call_id: Option<String>,
    delta: Option<String>,
    text: Option<String>,
    summary_index: Option<i64>,
    content_index: Option<i64>,
    safety_buffering: Option<Value>,
}

正文增量、工具输入和 reasoning 分别保留自己的关联键。工具 delta 优先使用 item_id,缺失时才用 call_id 作为 item key,同时仍保留原始 call_id;这使消费者能兼容不同服务端事件形态。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: process_responses_event 增量分支

rust
match event.kind.as_str() {
    "response.output_text.delta" => {
        if let Some(delta) = event.delta {
            return Ok(Some(ResponseEvent::OutputTextDelta(delta)));
        }
    }
    "response.custom_tool_call_input.delta" => {
        if let (Some(delta), Some(item_id)) =
            (event.delta, event.item_id.clone().or(event.call_id.clone()))
        {
            return Ok(Some(ResponseEvent::ToolCallInputDelta {
                item_id,
                call_id: event.call_id,
                delta,
            }));
        }
    }
    "response.reasoning_summary_text.delta" => {
        if let (Some(delta), Some(summary_index)) =
            (event.delta, event.summary_index)
        {
            return Ok(Some(ResponseEvent::ReasoningSummaryDelta {
                delta,
                summary_index,
            }));
        }
    }
    "response.reasoning_text.delta" => {
        if let (Some(delta), Some(content_index)) =
            (event.delta, event.content_index)
        {
            return Ok(Some(ResponseEvent::ReasoningContentDelta {
                delta,
                content_index,
            }));
        }
    }
    _ => trace!("unhandled responses event: {}", event.kind),
}

Ok(None)

未知 type 被忽略是协议前向兼容策略。它只保证旧客户端不会因新增的非关键事件立刻失败,并不保证新事件 的语义被实现;如果服务端把完成语义迁移到旧客户端未知的事件,最终仍会因缺少 response.completed 失败。

6. 协议终点 ​

response.completed 同时承担三件事:提供 response id、把 usage 转成 Codex TokenUsage,并作为 parser 的成功终止屏障。它一旦成功发送到 channel,SSE 任务立即 return,不等待底层 byte stream EOF。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: process_responses_event 完成分支

rust
"response.completed" => {
    if let Some(resp_val) = event.response {
        match serde_json::from_value::<ResponseCompleted>(resp_val) {
            Ok(resp) => {
                return Ok(Some(ResponseEvent::Completed {
                    response_id: resp.id,
                    token_usage: resp.usage.map(Into::into),
                    end_turn: resp.end_turn,
                }));
            }
            Err(err) => {
                let error = format!("failed to parse ResponseCompleted: {err}");
                return Err(ResponsesEventError::Api(ApiError::Stream(error)));
            }
        }
    }
}

状态图中的 ErrorSaved → Completed 是代码允许的路径:若流在 response.failed 后又异常给出合法 response.completed,完成事件会成为终点。正常协议不应同时产生二者,但 parser 选择以显式 completed 屏障为准,而不是自行维护更严格的服务端状态机。

7. 错误保留 ​

response.failed 不是统一转换成字符串。错误码决定后续 Turn 是否可能重试、是否应触发 compact,或者应 直接把配额、内容策略和服务过载呈现为不同类别。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: process_responses_event 失败分支

rust
"response.failed" => {
    if let Some(resp_val) = event.response {
        let mut response_error =
            ApiError::Stream("response.failed event received".into());
        if let Some(error) = resp_val.get("error")
            && let Ok(error) = serde_json::from_value::<Error>(error.clone())
        {
            if is_context_window_error(&error) {
                response_error = ApiError::ContextWindowExceeded;
            } else if is_quota_exceeded_error(&error) {
                response_error = ApiError::QuotaExceeded;
            } else if is_usage_not_included(&error) {
                response_error = ApiError::UsageNotIncluded;
            } else if is_cyber_policy_error(&error) {
                response_error = ApiError::CyberPolicy {
                    message: cyber_policy_message(error.message),
                };
            } else if matches!(
                error.code.as_deref(),
                Some("invalid_prompt" | "bio_policy")
            ) {
                response_error = ApiError::InvalidRequest {
                    message: error
                        .message
                        .unwrap_or_else(|| "Invalid request.".to_string()),
                };
            } else if is_server_overloaded_error(&error) {
                response_error = ApiError::ServerOverloaded;
            } else {
                response_error = ApiError::Retryable {
                    delay: try_parse_retry_after(&error),
                    message: error.message.unwrap_or_default(),
                };
            }
        }
        return Err(ResponsesEventError::Api(response_error));
    }
}
服务端 codeApiError关键后果
context_length_exceededContextWindowExceeded上层可进入上下文压缩策略
insufficient_quotaQuotaExceeded不是普通退避重试
usage_not_includedUsageNotIncluded表示账户能力不包含请求用途
cyber_policyCyberPolicy保留消息,空消息使用固定 fallback
invalid_prompt、bio_policyInvalidRequest不应当作瞬时传输故障
server_is_overloaded、slow_downServerOverloadedprovider 可映射成服务过载策略
其他已解析错误Retryable可携带从消息提取的 retry delay

try_parse_retry_after 只在 code 为 rate_limit_exceeded 时解析 “try again in 11.054s” 或毫秒形式。 这是对服务端 message 的窄解析,不是通用自然语言时间解析器;格式不匹配时 delay 为 None。

图中所有分支都汇入同一个保存点,强调分类和终止不是同一步:分类决定 ApiError,EOF 才决定何时把它 交给下游。response.completed 若随后到达,则由完成屏障结束 parser。

8. Mapper取消 ​

API parser 解决 wire 语义,Core mapper 解决运行时语义。它记录已完成 item、TTFT、inference trace 和 provider 错误映射,并创建 CancellationToken 监听最终消费者是否 drop stream。

源码位置:codex-rs/core/src/client.rs :: map_response_events

rust
let (tx_event, rx_event) =
    mpsc::channel::<Result<ResponseEvent>>(RESPONSE_STREAM_CHANNEL_CAPACITY);
let (tx_last_response, rx_last_response) = oneshot::channel::<LastResponse>();
let consumer_dropped = CancellationToken::new();
let consumer_dropped_for_stream = consumer_dropped.clone();

tokio::spawn(async move {
    let mut items_added: Vec<ResponseItem> = Vec::new();
    let mut api_stream = api_stream;

    loop {
        let event = tokio::select! {
            _ = consumer_dropped.cancelled() => {
                inference_trace_attempt.record_cancelled(
                    STREAM_DROPPED_REASON,
                    upstream_request_id,
                    &items_added,
                );
                return;
            }
            event = api_stream.next() => event,
        };

        let Some(event) = event else { break };
        match event {
            Ok(ResponseEvent::OutputItemDone(item)) => {
                items_added.push(item.clone());
                if tx_event
                    .send(Ok(ResponseEvent::OutputItemDone(item)))
                    .await
                    .is_err()
                {
                    inference_trace_attempt.record_cancelled(
                        STREAM_DROPPED_REASON,
                        upstream_request_id,
                        &items_added,
                    );
                    return;
                }
            }
            Ok(ResponseEvent::Completed {
                response_id,
                token_usage,
                end_turn,
            }) => {
                inference_trace_attempt.record_completed(
                    &response_id,
                    upstream_request_id,
                    &token_usage,
                    &items_added,
                );
                if let Some(sender) = tx_last_response.take() {
                    let _ = sender.send(LastResponse {
                        response_id: response_id.clone(),
                        items_added: std::mem::take(&mut items_added),
                    });
                }
                if tx_event
                    .send(Ok(ResponseEvent::Completed {
                        response_id,
                        token_usage,
                        end_turn,
                    }))
                    .await
                    .is_err()
                {
                    return;
                }
            }
            Ok(event) => {
                if tx_event.send(Ok(event)).await.is_err() {
                    return;
                }
            }
            Err(err) => {
                let mapped = provider.map_api_error(err);
                inference_trace_attempt.record_failed(
                    &mapped,
                    upstream_request_id,
                    &items_added,
                );
                if tx_event.send(Err(mapped)).await.is_err() {
                    return;
                }
            }
        }
    }
});

Core ResponseStream 的 Drop 会 cancel 这个 token。取消不是向模型事件流注入一个伪造错误,而是让 mapper 停止 poll API stream并记录本次 attempt 为 cancelled。mapper 被 drop 后,内部 API ResponseStream 和 channel receiver 随 task 一起释放,SSE producer 随后在 send 失败处分支 return,完成两级资源清理。

源码位置:codex-rs/core/src/client_common.rs :: ResponseStream

rust
pub struct ResponseStream {
    pub(crate) rx_event: mpsc::Receiver<Result<ResponseEvent>>,
    pub(crate) consumer_dropped: CancellationToken,
}

impl Stream for ResponseStream {
    type Item = Result<ResponseEvent>;

    fn poll_next(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
    ) -> Poll<Option<Self::Item>> {
        self.rx_event.poll_recv(cx)
    }
}

impl Drop for ResponseStream {
    fn drop(&mut self) {
        self.consumer_dropped.cancel();
    }
}

图中的清理是所有权链自然传播的结果:取消 token 停 mapper,receiver drop 让上游 send 失败,上游 task return 再释放网络流。代码没有单独调用“关闭 SSE socket”的全局函数。

9. Turn消费 ​

Turn 使用自己的 cancellation token 包装 stream.next()。用户中断首先得到 TurnAborted;普通 stream error 保留 mapper 映射后的 CodexErr;没有 error 却 channel 关闭则再生成提前断流错误。

源码位置:codex-rs/core/src/session/turn.rs :: run_sampling_request 接收循环

rust
let event = match stream
    .next()
    .instrument(trace_span!(parent: &handle_responses, "receiving"))
    .or_cancel(&cancellation_token)
    .await
{
    Ok(event) => event,
    Err(codex_async_utils::CancelErr::Cancelled) => {
        break Err(CodexErr::TurnAborted);
    }
};

let event = match event {
    Some(Ok(event)) => event,
    Some(Err(err)) => break Err(err),
    None => {
        break Err(CodexErr::Stream(
            "stream closed before response.completed".into(),
        ));
    }
};

match event {
    ResponseEvent::Completed {
        response_id,
        token_usage,
        end_turn,
    } => {
        sess.send_event(
            &turn_context,
            EventMsg::RawResponseCompleted(RawResponseCompletedEvent {
                response_id,
                token_usage: token_usage.clone(),
            }),
        )
        .await;
        let budget_result = sess
            .record_token_usage_info(&turn_context, token_usage.as_ref())
            .await;
        should_emit_token_count = true;
        should_emit_turn_diff = true;
        if let Err(err) = budget_result {
            break Err(err);
        }
        if let Some(false) = end_turn {
            needs_follow_up = true;
        }
        break Ok(SamplingRequestResult {
            needs_follow_up,
            last_agent_message,
        });
    }
    ResponseEvent::OutputTextDelta(delta) => {
        if let Some(active) = active_item.as_ref() {
            let item_id = active.id();
            let parsed = assistant_message_stream_parsers
                .parse_delta(&item_id, &delta);
            emit_streamed_assistant_text_delta(
                &sess,
                &turn_context,
                plan_mode_state.as_mut(),
                &item_id,
                parsed,
            )
            .await;
        }
    }
    _ => {}
}

因此 Completed 的消费者副作用远多于“停止循环”:它刷新剩余文本 segment,发送 raw completion,写 token usage,并根据 end_turn=false 决定是否继续下一次 sampling。SSE parser 只证明一次 provider response 完成; 它不等于整个 Codex Turn 必然结束。

10. 流解析测试 ​

测试没有伪造 ResponseEvent 直接塞进 channel,而是用 tokio_test::io::Builder 和 ReaderStream 构造真实 字节 chunk,再经过 eventsource_stream。这能同时覆盖 chunk 读取、SSE framing、JSON 解码和事件发送。

源码位置:codex-rs/codex-api/src/sse/responses.rs :: collect_events

rust
async fn collect_events(
    chunks: &[&[u8]],
) -> Vec<Result<ResponseEvent, ApiError>> {
    let mut builder = IoBuilder::new();
    for chunk in chunks {
        builder.read(chunk);
    }

    let reader = builder.build();
    let stream = ReaderStream::new(reader)
        .map_err(|err| TransportError::Network(err.to_string()));
    let (tx, mut rx) =
        mpsc::channel::<Result<ResponseEvent, ApiError>>(16);
    tokio::spawn(process_sse(
        Box::pin(stream),
        tx,
        idle_timeout(),
        /* telemetry */ None,
    ));

    let mut events = Vec::new();
    while let Some(event) = rx.recv().await {
        events.push(event);
    }
    events
}

parses_items_and_completed 输入两个 response.output_item.done 和一个 response.completed,断言输出严格为 两个 message item 加一个 completion,并检查 commentary phase、response id、usage 和 end_turn。它证明 多 frame 顺序与基本转换,不证明网络 idle timeout 或 Core mapper 的取消。

error_when_missing_completed 只输入一个合法 item 后关闭 reader,断言先收到 item,再收到精确错误 stream closed before response.completed。它证明已经产出的 item 不会因尾部断流被撤回,也证明 EOF 不能冒充 成功完成;它不证明上层是否会重试整次 sampling。

error_when_error_event 输入 rate-limit response.failed,随后 reader EOF,断言最终得到 ApiError::Retryable,message 保持不变且 delay 精确为 11.054s。它证明已保存的业务错误优先于通用 EOF 错误,不证明所有 provider 都使用同一 message 格式。

emits_completed_without_stream_end 在 completion frame 后拼接永不结束的 pending stream,并用外层 timeout 收集结果。断言 parser 仍立即只返回一个 Completed,证明协议终点是 event,而不是 TCP EOF。

11. 只读复核 ​

先用下面的搜索区分四个终止点,不要只搜索 ResponseEvent::Completed:

bash
rg -n "spawn_response_stream|process_sse_with_treatment|map_response_events|run_sampling_request" \
  codex-rs/codex-api/src codex-rs/core/src

rg -n "stream closed before response.completed|idle timeout waiting for SSE|consumer_dropped" \
  codex-rs/codex-api/src codex-rs/core/src

再运行最小测试集:

bash
cd codex-rs
cargo test -p codex-api sse::responses::tests::parses_items_and_completed -- --exact
cargo test -p codex-api sse::responses::tests::error_when_missing_completed -- --exact
cargo test -p codex-api sse::responses::tests::error_when_error_event -- --exact
cargo test -p codex-api sse::responses::tests::emits_completed_without_stream_end -- --exact

复核时可以用两个故障问题检验是否真正走通源码:若 UI 已显示部分文本后报“stream closed before response.completed”,已显示 item 为什么没有回滚,错误最早在哪个 parser 分支产生;若用户中断 Turn,为什么 trace 应记录 cancelled 而不是 failed,哪一个 receiver 的释放最终让网络读取任务停止。答案都需要同时经过 SSE parser、Core mapper 和 Turn consumer,不能只引用其中一层。

本文结论只覆盖 HTTP SSE Responses 路径。WebSocket 使用相同 ResponseEvent 目标类型,但它有独立的帧、 连接复用和完成处理代码;不能把本文关于 eventsource_stream、HTTP headers 和 SSE idle poll 的细节外推到 WebSocket transport。