Codex API类型
本文承接 模型请求构造,研究 codex-api 如何把请求和响应分成三层:ResponsesApiRequest 是内部统一请求,ResponseCreateWsRequest 是 WebSocket 序列化视图,ResponseEvent 是 SSE/WebSocket 事件转换后供 core 消费的事件。
问题边界很重要:本文不讲 build_responses_request 怎么决定字段,也不展开下一篇的 SSE 流程调度。读完后,读者应能根据源码判断:哪些 Option 字段会被省略,哪些字段只是 WebSocket 附加项,以及哪类服务端事件会被丢弃、转成事件或转成错误。
1. 三层边界
ResponsesApiRequest 是所有者:HTTP endpoint 直接序列它,WebSocket 通过 From<&ResponsesApiRequest> 借用它的字段。因此不应把 WebSocket JSON 当成另一套独立 request model。client_metadata 在 WebSocket 视图中会被复制,而 request input/tools/reasoning 等字段仍以借用视图复用。
2. Request序列化
源码位置:codex-rs/codex-api/src/common.rs :: ResponsesApiRequest
#[derive(Debug, Serialize, Clone, PartialEq)]
pub struct ResponsesApiRequest {
pub model: String,
#[serde(skip_serializing_if = "String::is_empty")]
pub instructions: String,
pub input: Vec<ResponseItem>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tools: Option<ResponsesApiTools>,
pub tool_choice: String,
pub parallel_tool_calls: bool,
pub reasoning: Option<Reasoning>,
pub store: bool,
pub stream: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub stream_options: Option<StreamOptions>,
pub include: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub service_tier: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub prompt_cache_key: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub text: Option<TextControls>,
#[serde(skip_serializing_if = "Option::is_none")]
pub client_metadata: Option<HashMap<String, String>>,
}这里有两个 wire 规则:空 instructions 不会产生 JSON 字段;所有标记 skip_serializing_if = Option::is_none 的字段在 None 时也不会产生字段。reasoning 没有 skip 标记,因此 builder 即使没有 effort,也可能发送 reasoning 对象;对象内部的可选字段再单独省略。
3. WebSocket视图
ResponseCreateWsRequest 不复制所有字符串,而是借用 ResponsesApiRequest 中的数据。previous_response_id 和 generate 是 WebSocket request 专用控制项,普通 Responses request 没有这两个字段。
源码位置:codex-rs/codex-api/src/common.rs :: From<&ResponsesApiRequest> for ResponseCreateWsRequest
impl<'a> From<&'a ResponsesApiRequest> for ResponseCreateWsRequest<'a> {
fn from(request: &'a ResponsesApiRequest) -> Self {
Self {
model: &request.model,
instructions: &request.instructions,
previous_response_id: None,
input: &request.input,
tools: request.tools.as_ref().map(ResponsesApiTools::as_raw_value),
tool_choice: &request.tool_choice,
parallel_tool_calls: request.parallel_tool_calls,
reasoning: request.reasoning.as_ref(),
store: request.store,
stream: request.stream,
stream_options: request.stream_options.as_ref(),
include: &request.include,
service_tier: request.service_tier.as_deref(),
prompt_cache_key: request.prompt_cache_key.as_deref(),
text: request.text.as_ref(),
generate: None,
client_metadata: request.client_metadata.clone(),
}
}
}tools 使用 RawValue 视图,这不是为了改写工具 JSON,而是避免在 WebSocket 包装时重新解析和重新构造。
4. ResponseItem类型
ResponseItem 是 input 和 output 共用的协议枚举,通过 serde 的 type tag 区分变体。这里的“类型”不仅是 Rust enum 名称,也决定了 wire JSON 的 type 字段。
源码位置:codex-rs/protocol/src/models.rs :: ResponseItem
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, JsonSchema, TS)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ResponseItem {
AdditionalTools {
#[serde(default, skip_serializing_if = "Option::is_none")]
id: Option<ResponseItemId>,
role: String,
tools: Vec<serde_json::Value>,
},
Message {
#[serde(default, skip_serializing_if = "Option::is_none")]
id: Option<ResponseItemId>,
role: String,
content: Vec<ContentItem>,
#[serde(default, skip_serializing_if = "Option::is_none")]
phase: Option<MessagePhase>,
#[serde(default, skip_serializing_if = "Option::is_none")]
internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>,
},
AgentMessage {
#[serde(default, skip_serializing_if = "Option::is_none")]
id: Option<ResponseItemId>,
author: String,
recipient: String,
content: Vec<AgentMessageInputContent>,
#[serde(default, skip_serializing_if = "Option::is_none")]
internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>,
},
Reasoning {
#[serde(default, skip_serializing_if = "Option::is_none")]
id: Option<ResponseItemId>,
summary: Vec<ReasoningItemReasoningSummary>,
#[serde(default, skip_serializing_if = "should_serialize_reasoning_content")]
content: Option<Vec<ReasoningItemContent>>,
encrypted_content: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>,
},
LocalShellCall {
#[serde(default, skip_serializing_if = "Option::is_none")]
id: Option<ResponseItemId>,
call_id: Option<String>,
status: LocalShellStatus,
action: LocalShellAction,
#[serde(default, skip_serializing_if = "Option::is_none")]
internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>,
},
FunctionCall {
#[serde(default, skip_serializing_if = "Option::is_none")]
id: Option<ResponseItemId>,
name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
namespace: Option<String>,
arguments: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
encrypted_function_args: Option<Vec<String>>,
call_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>,
},
}上面是与本篇主线直接相关的连续变体摘录;当前 ResponseItem 还包含 shell、tool search、compaction 等变体。不能把这段摘录当成完整枚举,也不能根据 FunctionCall.arguments: String 推断服务端发送的是 JSON object:源码注释明确它是“包含 JSON 的字符串”。
5. Text与Schema
create_text_param_for_request 把 verbosity 和 output schema 收敛成同一个 TextControls。两者都没有时返回 None;schema 的 wire name 固定为 codex_output_schema,strict 值由调用方传入。
源码位置:codex-rs/codex-api/src/common.rs :: create_text_param_for_request
pub fn create_text_param_for_request(
verbosity: Option<VerbosityConfig>,
output_schema: &Option<Value>,
output_schema_strict: bool,
) -> Option<TextControls> {
if verbosity.is_none() && output_schema.is_none() {
return None;
}
Some(TextControls {
verbosity: verbosity.map(std::convert::Into::into),
format: output_schema.as_ref().map(|schema| TextFormat {
r#type: TextFormatType::JsonSchema,
strict: output_schema_strict,
schema: schema.clone(),
name: "codex_output_schema".to_string(),
}),
})
}6. Event映射
ResponseEvent 是 core 消费的稳定事件集合,不等于服务端原始事件名集合。process_responses_event 对已知事件做字段完整性检查:缺少必要字段时可能返回 Ok(None),而不是构造半成品事件;response.failed 会按 context/quota/policy/server 等错误映射为不同 ApiError。
源码位置:codex-rs/codex-api/src/sse/responses.rs :: process_responses_event
pub fn process_responses_event(
event: ResponsesStreamEvent,
) -> std::result::Result<Option<ResponseEvent>, ResponsesEventError> {
match event.kind.as_str() {
"response.output_item.done" => {
if let Some(item_val) = event.item {
if let Ok(item) = serde_json::from_value::<ResponseItem>(item_val) {
return Ok(Some(ResponseEvent::OutputItemDone(item)));
}
debug!("failed to parse ResponseItem from output_item.done");
}
}
"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.done" => {
if let (Some(item_id), Some(text), Some(summary_index)) =
(event.item_id, event.text, event.summary_index)
{
return Ok(Some(ResponseEvent::ReasoningSummaryDone {
item_id,
text,
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,
}));
}
}
"response.created" => {
if event.response.is_some() {
return Ok(Some(ResponseEvent::Created {}));
}
}
"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) {
let message = cyber_policy_message(error.message);
response_error = ApiError::CyberPolicy { message };
} else if matches!(error.code.as_deref(), Some("invalid_prompt" | "bio_policy"))
{
let message = error
.message
.unwrap_or_else(|| "Invalid request.".to_string());
response_error = ApiError::InvalidRequest { message };
} else if is_server_overloaded_error(&error) {
response_error = ApiError::ServerOverloaded;
} else {
let delay = try_parse_retry_after(&error);
let message = error.message.unwrap_or_default();
response_error = ApiError::Retryable { message, delay };
}
}
return Err(ResponsesEventError::Api(response_error));
}
return Err(ResponsesEventError::Api(ApiError::Stream(
"response.failed event received".into(),
)));
}
"response.incomplete" => {
let reason = event.response.as_ref().and_then(|response| {
response
.get("incomplete_details")
.and_then(|details| details.get("reason"))
.and_then(Value::as_str)
});
let reason = reason.unwrap_or("unknown");
let message = format!("Incomplete response returned, reason: {reason}");
return Err(ResponsesEventError::Api(ApiError::Stream(message)));
}
"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}");
debug!("{error}");
return Err(ResponsesEventError::Api(ApiError::Stream(error)));
}
}
}
}
"response.output_item.added" => {
if let Some(item_val) = event.item {
if let Ok(item) = serde_json::from_value::<ResponseItem>(item_val) {
return Ok(Some(ResponseEvent::OutputItemAdded(item)));
}
debug!("failed to parse ResponseItem from output_item.added");
}
}
"response.reasoning_summary_part.added" => {
if let Some(summary_index) = event.summary_index {
return Ok(Some(ResponseEvent::ReasoningSummaryPartAdded {
summary_index,
}));
}
}
_ => {
trace!("unhandled responses event: {}", event.kind);
}
}
Ok(None)
}这个处理函数返回 Ok(None) 时,上层 SSE loop 不会收到事件;但 response.failed 、response.incomplete 和完成事件解析失败会返回 Err。这个返回值是 event mapper 与流程调度之间的边界,不是服务端原始事件的一一镜像。
7. 错误边界
response.failed 不是统一的普通消息:它会根据 error code 分类为 context window、quota、usage、cyber policy、invalid request、server overloaded 或 retryable。response.incomplete 则转成带 reason 的 stream error。如果 SSE 在 response.completed 前关闭,也会被视为 stream error,不会被当作正常完成。
8. 协议测试
| 测试 | 输入与动作 | 断言 | 覆盖边界 |
|---|---|---|---|
responses_client_stream_request_preserves_item_ids | 构造带 item id 的 request 并通过 recording transport 发送 | JSON body 与请求对象一致,content type 为 JSON | 不证明服务端业务执行 |
direct_serialization_preserves_websocket_request_payload | 从 API request 生成 response.create payload | WebSocket payload 与 HTTP request 共享字段,并保留专用控制项 | 不证明 WebSocket 连接成功 |
parses_items_and_completed | SSE output item + completed | item 转为 ResponseItem,完成事件保留 id/usage/end_turn | 不覆盖所有 event kind |
error_when_missing_completed | SSE 只发 output item 后关闭 | 返回 stream closed before completed | 不证明 retry 策略 |
这些测试覆盖请求序列化、WebSocket 视图、正常 SSE 完成和缺少完成事件的错误路径。它们不覆盖全部事件类型、真实连接建立或 core 的后续重试策略。
9. 复现路径
RUST_MIN_STACK=16777216 cargo test -p codex-api responses_client_stream_request_preserves_item_ids
RUST_MIN_STACK=16777216 cargo test -p codex-api direct_serialization_preserves_websocket_request_payload
RUST_MIN_STACK=16777216 cargo test -p codex-api parses_items_and_completed
RUST_MIN_STACK=16777216 cargo test -p codex-api error_when_missing_completed源码练习:在 ResponsesApiRequest 中将 instructions 设为空字符串,观察 JSON 中它是否还存在;再在 SSE fixture 中写入未知 type,验证它是被忽略还是被视为错误。
10. 源码导航
先读 codex-api/src/common.rs 中的 ResponsesApiRequest 和 WebSocket 视图,再对照 protocol/src/models.rs 的 ResponseItem serde tag;最后读 sse/responses.rs 的 event mapping 和错误分类。下一篇进入流式事件解析时,应将“原始 event”、ResponseEvent 和 core 消费者三者分开。
