Skip to content

Codex API类型

解析 ResponsesApiRequest、WebSocket request、ResponseItem 与 ResponseEvent 之间的序列化和事件边界。

基于rust-v0.150.0
CodexRustModelResponses

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

rust
#[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

rust
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

rust
#[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

rust
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

rust
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 payloadWebSocket payload 与 HTTP request 共享字段,并保留专用控制项不证明 WebSocket 连接成功
parses_items_and_completedSSE output item + completeditem 转为 ResponseItem,完成事件保留 id/usage/end_turn不覆盖所有 event kind
error_when_missing_completedSSE 只发 output item 后关闭返回 stream closed before completed不证明 retry 策略

这些测试覆盖请求序列化、WebSocket 视图、正常 SSE 完成和缺少完成事件的错误路径。它们不覆盖全部事件类型、真实连接建立或 core 的后续重试策略。

9. 复现路径 ​

bash
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 消费者三者分开。