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 endpoint | ResponsesClient | request、headers、compression | StreamResponse | 请求失败或取得 response |
| SSE parser | spawn_response_stream 派生任务 | headers、ByteStream | codex_api::ResponseStream | Completed、错误、EOF、idle timeout、接收者关闭 |
| Core mapper | map_response_events 派生任务 | API ResponseEvent | Core ResponseStream | consumer drop 或 API stream 结束 |
| Turn consumer | run_sampling_request | Core event stream | runtime 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 返回与认证分支
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
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
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
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
#[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 增量分支
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 完成分支
"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 失败分支
"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));
}
}| 服务端 code | ApiError | 关键后果 |
|---|---|---|
context_length_exceeded | ContextWindowExceeded | 上层可进入上下文压缩策略 |
insufficient_quota | QuotaExceeded | 不是普通退避重试 |
usage_not_included | UsageNotIncluded | 表示账户能力不包含请求用途 |
cyber_policy | CyberPolicy | 保留消息,空消息使用固定 fallback |
invalid_prompt、bio_policy | InvalidRequest | 不应当作瞬时传输故障 |
server_is_overloaded、slow_down | ServerOverloaded | provider 可映射成服务过载策略 |
| 其他已解析错误 | 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
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
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 接收循环
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
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:
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再运行最小测试集:
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。
