解包 OpenAI Codex:SQ/EQ 协议与 Agent 会话模型

OpenAI Codex 是一个终端 Agent,用 Rust 写核心,TypeScript 写前端。它不是一个简单的"调用 API 然后执行命令"的脚本,而是一套完整的 Agent 运行时。理解它的架构,从协议层开始最清晰。

本文聚焦三个核心:SQ/EQ 通信协议、CodexThread 会话模型、上下文压缩机制。代码引用来自 codex-rs 子目录。

一、SQ/EQ 协议:为什么要分开

协议文件开头第一行注释就说清楚了设计意图:

//! Defines the protocol for a Codex session between a client and an agent.
//!
//! Uses a SQ (Submission Queue) / EQ (Event Queue) pattern to asynchronously communicate
//! between user and agent.

SQ(Submission Queue)是用户向 Agent 发指令的通道,EQ(Event Queue)是 Agent 向用户推事件的通道。两个方向完全独立。

这个模式解决的核心问题是:Agent 执行是异步的、长时间的,中间会产生大量中间状态(命令开始执行、输出增量、需要审批、执行结束……),这些都要实时推给客户端。如果用请求-响应模型,要么客户端一直轮询,要么连接一直挂着等最终结果,都不合适。SQ/EQ 把"发指令"和"接事件"解耦,每个方向各自独立,互不阻塞。

SQ 的入口是 Submission 结构体:

pub struct Submission {
    /// Unique id for this Submission to correlate with Events
    pub id: String,
    /// Payload
    pub op: Op,
    pub client_user_message_id: Option<String>,
    pub trace: Option<W3cTraceContext>,
}

每个 Submission 有一个唯一 ID,EQ 里的 Event 会带同一个 ID,这样客户端就能把"我发了这个操作"和"Agent 产生的这些事件"关联起来。trace 字段支持 W3C 标准的分布式追踪上下文,说明这套协议本来就考虑了多层跳转的场景。

EQ 的入口是 Event

pub struct Event {
    /// Submission `id` that this event is correlated with.
    pub id: String,
    /// Payload
    pub msg: EventMsg,
}

二、Op 枚举:用户能做什么

Op 枚举定义了所有从客户端到 Agent 的操作类型。它标注了 #[non_exhaustive](外部 crate 不能对其做穷举匹配,协议未来可以继续扩展而不破坏兼容性),以及 #[serde(tag = "type", rename_all = "snake_case")](在 JSON wire format 里每个变体用 type 字段区分,值为 snake_case,比如 "type": "user_input")。

按功能分组,Op 分为六类:

会话控制

/// 中断当前任务,但不终止后台终端进程
/// 服务端响应 EventMsg::TurnAborted
Interrupt,

/// 终止该线程下所有正在运行的后台终端进程
/// 用于主动停止长期运行的后台 shell(有别于 Interrupt)
CleanBackgroundTerminals,

/// 关闭整个 Codex 实例
Shutdown,

InterruptCleanBackgroundTerminals 的区别值得注意:Interrupt 只是打断当前 Agent Turn,后台 shell 进程(比如长时间运行的 dev server)还活着;CleanBackgroundTerminals 则主动终止所有后台进程,用于用户确实想停掉一切的场景。

用户输入

/// 用户发送消息,可附带本次 Turn 的环境选择、输出 Schema、
/// Responses API 元数据、附加上下文,以及永久生效的线程配置覆盖
UserInput {
    items: Vec<UserInput>,
    environments: Option<Vec<TurnEnvironmentSelection>>,
    final_output_json_schema: Option<Value>,
    responsesapi_client_metadata: Option<HashMap<String, String>>,
    additional_context: BTreeMap<String, AdditionalContextEntry>,
    // flatten:直接内嵌线程配置字段,无需嵌套对象
    #[serde(default, flatten)]
    thread_settings: ThreadSettingsOverrides,
},

/// 只更新线程配置,不发消息、不启动 Turn
/// 与 UserInput 走同一个 submission queue,保证调用顺序
ThreadSettings {
    #[serde(flatten)]
    thread_settings: ThreadSettingsOverrides,
},

/// 用户用 !cmd 直接执行 shell 命令
/// 使用用户默认 shell,支持管道/重定向等 shell 语法
/// 输出通过 ExecCommand* 事件流式推送
RunUserShellCommand {
    command: String,
},

UserInput 通过 #[serde(flatten)]ThreadSettingsOverrides 的字段直接平铺进来,意味着用户每次发消息可以顺带修改模型、沙箱策略、审批模式等配置,不需要单独发一个 ThreadSettings Op——两步合一,减少一次往返。

final_output_json_schema 允许调用方约束 Agent 本次输出的 JSON 结构,这是结构化输出场景(比如让 Agent 返回一个格式化的 JSON 对象)的入口。

审批回应(与 EQ 审批请求配对)

EQ 里有若干"Agent 需要人类介入"的事件,客户端处理后通过对应的 Op 回应:

/// 审批/拒绝 Agent 要执行的命令
ExecApproval { id: String, turn_id: Option<String>, decision: ReviewDecision },

/// 审批/拒绝 Agent 要应用的代码补丁
PatchApproval { id: String, decision: ReviewDecision },

/// 回应 request_user_input 工具调用(Agent 向用户提问)
#[serde(rename = "user_input_answer", alias = "request_user_input_response")]
UserInputAnswer { id: String, response: RequestUserInputResponse },

/// 回应 request_permissions 工具调用(运行时申请额外权限)
RequestPermissionsResponse { id: String, response: RequestPermissionsResponse },

/// 回应 dynamic tool 调用(外部系统处理完工具请求后返回结果)
DynamicToolResponse { id: String, response: DynamicToolResponse },

/// 回应 MCP elicitation 请求(MCP 服务器向用户征集输入)
ResolveElicitation {
    server_name: String,
    request_id: RequestId,
    decision: ElicitationAction,
    content: Option<Value>,
    meta: Option<Value>,
},

/// 批准被 Guardian 系统拒绝的操作(一次性覆盖)
ApproveGuardianDeniedAction { event: GuardianAssessmentEvent },

这七个 Op 构成了"Agent 暂停等待人类决策"的完整机制。它们都携带一个 id,与 EQ 里的请求事件一一对应。UserInputAnswer 有一个历史兼容 alias(request_user_input_response),说明这个 Op 经历过重命名。

会话管理

/// 手动触发上下文压缩(对话历史超长时也会自动触发)
Compact,

/// 从内存上下文中删除最近 N 个用户 Turn
/// 注意:不会撤销磁盘上的文件改动,客户端需自行处理
ThreadRollback { num_turns: u32 },

/// 设置线程是否允许生成记忆(Memory)
SetThreadMemoryMode { mode: ThreadMemoryMode },

/// 重新加载用户配置层(不重启线程)
/// 用于运行时热更新 app 启用/禁用状态等配置
ReloadUserConfig,

/// 重新初始化 MCP 服务器并刷新工具列表缓存
RefreshMcpServers { config: McpServerRefreshConfig },

/// 请求 Agent 对当前代码做 Review
Review { review_request: ReviewRequest },

ThreadRollback 的注释特别值得关注:它只修改内存里的对话历史,不会撤销 Agent 已经写到磁盘上的文件改动。这是一个有意为之的边界——Codex 认为文件系统状态由客户端或用户自行管理,Agent 只负责管自己的会话状态。

实时对话(语音)

RealtimeConversationStart(ConversationStartParams),  // 开启语音流
RealtimeConversationAudio(ConversationAudioParams),  // 发送音频数据
RealtimeConversationText(ConversationTextParams),    // 发送文字到语音流
RealtimeConversationClose,                          // 关闭语音流
RealtimeConversationListVoices,                     // 获取支持的声音列表

这五个 Op 支持语音交互模式(通过 WebRTC)。在普通的文字 coding 场景下不会用到,但它们的存在说明 Codex 的协议层从设计上就不局限于文字。

多 Agent 协作

InterAgentCommunication {
    communication: InterAgentCommunication,
},

这是整个 Op 枚举里最值得深入看的一个。InterAgentCommunication 结构体如下:

pub struct InterAgentCommunication {
    pub author: AgentPath,     // 发送方 Agent 的路径标识
    pub recipient: AgentPath,  // 主接收方
    pub other_recipients: Vec<AgentPath>,  // 抄送方
    pub content: String,
    pub encrypted_content: Option<String>,  // 加密内容(与 content 互斥)
    pub trigger_turn: bool,  // 收到消息后是否立即触发一个新 Turn
}

trigger_turn 这个字段说明了两种 Agent 间通信模式:true 时消息到达后立即驱动接收方 Agent 开始一个新的 Turn(主动推送,适合任务委派);false 时只是把消息写进历史,等接收方下次被其他原因触发时自然读到(被动通知)。

消息通过 to_model_input_item() 转换成模型的输入格式:未加密的走普通 assistant message;加密的走专门的 AgentMessage item 类型,内容是 EncryptedContent。加密消息说明 Codex 的多 Agent 场景设计里考虑了 Agent 之间信息隔离的需求。

三、EventMsg 枚举:Agent 会说什么

EventMsg 是 EQ 的 payload,定义了 Agent 能产生的所有事件类型。数量很多,但有清晰的分组逻辑。

Turn 生命周期事件

// v1 wire format 用 task_started,v2 用 turn_started,两者都接受
#[serde(rename = "task_started", alias = "turn_started")]
TurnStarted(TurnStartedEvent),

#[serde(rename = "task_complete", alias = "turn_complete")]
TurnComplete(TurnCompleteEvent),

TurnAborted(TurnAbortedEvent),

这里有一个向后兼容的处理:v1 协议叫 task_started,v2 改名叫 turn_started,通过 alias 两者都能反序列化。

命令执行事件(Begin/Delta/End 三段式):

ExecCommandBegin(ExecCommandBeginEvent),
ExecCommandOutputDelta(ExecCommandOutputDeltaEvent),
ExecCommandEnd(ExecCommandEndEvent),

命令输出是流式推送的,OutputDelta 会持续发送增量内容,直到 End。这让终端 UI 能实时展示命令输出,而不是等命令跑完再一次性显示。

审批请求事件

ExecApprovalRequest(ExecApprovalRequestEvent),
ApplyPatchApprovalRequest(ApplyPatchApprovalRequestEvent),
RequestUserInput(RequestUserInputEvent),
RequestPermissions(RequestPermissionsEvent),
ElicitationRequest(ElicitationRequestEvent),

这五种都是"Agent 需要人类介入"的场景。它们通过 EQ 推给客户端,客户端通过对应的 OpExecApprovalPatchApprovalUserInputAnswerRequestPermissionsResponseResolveElicitation)回应。

上下文压缩事件

ContextCompacted(ContextCompactedEvent),
ThreadRolledBack(ThreadRolledBackEvent),

压缩完成后推这个事件,客户端据此更新 UI 状态(比如清空或折叠历史显示)。

Token 计数

TokenCount(TokenCountEvent),

TokenCountEvent 包含 total_token_usagelast_token_usagemodel_context_windowTokenUsage 结构体里有一个 percent_of_context_window_remaining 方法,它从总窗口里减去一个 BASELINE_TOKENS = 12000 的基线(用于系统提示和工具描述),让显示给用户的"剩余容量"更准确。

四、CodexThread:会话模型

CodexThread 是 Codex 里一个会话(Thread)的句柄。它封装了一个 Codex(底层会话对象),对外提供操作接口。

pub struct CodexThread {
    pub(crate) codex: Codex,
    pub(crate) session_source: SessionSource,
    session_configured: SessionConfiguredEvent,
    rollout_path: Option<PathBuf>,
    out_of_band_elicitation_count: Mutex<u64>,
}

rollout_path 是会话历史的持久化路径——Codex 把每个会话的事件流写到磁盘,叫做 rollout,可以用来恢复历史或 fork 新会话。

out_of_band_elicitation_count 是一个引用计数,用于跟踪"带外 elicitation"的数量。当有未处理的带外审批请求时,会暂停正常的 Turn 流程,等审批完成再恢复。

ThreadConfigSnapshot 是线程配置的快照,包含了一次 Turn 需要的所有运行时参数:

pub struct ThreadConfigSnapshot {
    pub model: String,
    pub model_provider_id: String,
    pub service_tier: Option<String>,
    pub approval_policy: AskForApproval,
    pub approvals_reviewer: ApprovalsReviewer,
    pub permission_profile: PermissionProfile,
    pub active_permission_profile: Option<ActivePermissionProfile>,
    pub cwd: AbsolutePathBuf,
    pub workspace_roots: Vec<AbsolutePathBuf>,
    pub profile_workspace_roots: Vec<AbsolutePathBuf>,
    pub ephemeral: bool,
    pub reasoning_effort: Option<ReasoningEffort>,
    pub reasoning_summary: Option<ReasoningSummary>,
    pub personality: Option<Personality>,
    pub collaboration_mode: CollaborationMode,
    pub session_source: SessionSource,
    pub forked_from_thread_id: Option<ThreadId>,
    pub parent_thread_id: Option<ThreadId>,
    pub thread_source: Option<ThreadSource>,
}

这个快照里有几个有意思的字段:workspace_rootsprofile_workspace_roots 分开存——前者是运行时实际用于沙箱权限的路径,后者是 profile 配置里的路径,两者可以不同。forked_from_thread_idparent_thread_id 说明 Codex 支持会话 fork,一个会话可以从另一个会话的历史分叉出来。

TryStartTurnIfIdleRejectionReason 是一个细节,但说明了 Codex 对 Turn 启动有严格的并发控制:

pub enum TryStartTurnIfIdleRejectionReason {
    /// 有用户/客户端触发的任务在队列里,优先级更高
    PendingTriggerTurn,
    /// 线程处于 Plan 模式,不允许自动启动新 Turn
    PlanMode,
    /// 有其他 Turn 或任务正在运行,或 idle 预留在启动前已丢失
    Busy,
}

这个枚举用于扩展(Extension)想在线程空闲时自动启动工作的场景。三种拒绝原因清楚地描述了什么情况下不允许自动 Turn:有更高优先级的用户输入在等待、Plan 模式(规划模式下 Agent 只做规划不执行)、或者线程本身就忙。

CodexThread 的核心方法是 submit,它把一个 Op 放进 SQ:

pub async fn submit(&self, op: Op) -> CodexResult<String> {
    self.codex.submit(op).await
}

返回值是 Submission ID,客户端可以用这个 ID 关联后续的 EQ 事件。

五、Turn 执行流程:从 submit 到 TurnComplete

上面描述的 SQ/EQ 协议是通信骨架,真正让 Agent 工作的是 run_turn 函数(位于 session/turn.rs)。一次 Turn 的执行分为若干阶段:

阶段一:预压缩检查

Turn 开始前,先检查是否需要在采样(调用模型)之前就做一次压缩。如果上次 Turn 结束时上下文已接近窗口上限,下一次 Turn 带来的新用户消息会直接撑破限制——这时需要提前压缩:

if let Err(err) = run_pre_sampling_compact(&sess, &turn_context, &mut client_session).await {
    // ...推送错误事件,如果是 UsageLimitExceeded 则通知 goal 系统
    return None;
}

如果预压缩失败(比如 token 用量已超出限制),Turn 直接终止,不进行模型调用。

阶段二:工具和插件注入

接下来构建本次 Turn 需要的所有工具和插件:

let (injection_items, explicitly_enabled_connectors) =
    build_skills_and_plugins(&sess, turn_context.as_ref(), &input, &cancellation_token).await?;

build_skills_and_plugins 做的事情包括:扫描用户消息里是否有 @mention(显式提及某个工具或 skill)、加载 MCP 工具、构建 connector 工具、处理插件注入。返回的 injection_items 会在用户消息之前插入到对话历史,让模型知道有哪些工具可用。

阶段三:Hook 执行

Codex 有一套 hook 机制,允许在 Turn 的关键节点插入自定义逻辑。Turn 前有 run_pending_session_start_hooks(只在会话首次 Turn 触发)和 run_hooks_and_record_inputs(记录用户输入并执行 pre-turn hooks)。Turn 完成后有 run_turn_stop_hooksrun_legacy_after_agent_hook。如果任何 hook 返回"应该终止",Turn 提前退出。

阶段四:采样循环(核心)

真正的模型调用在 run_sampling_request 里,它是一个带重试的循环:

let max_retries = turn_context.provider.info().stream_max_retries();
let mut retries = 0;
loop {
    let prompt = build_prompt(prompt_input, router.as_ref(), turn_context.as_ref(), base_instructions.clone());
    let err = match try_run_sampling_request(tool_runtime.clone(), /* ... */).await {
        Ok(output) => { return Ok(output); }
        Err(CodexErr::ContextWindowExceeded) => {
            sess.set_total_tokens_full(&turn_context).await;
            return Err(CodexErr::ContextWindowExceeded);
        }
        Err(err) => err,
    };
    // 处理可重试的流式错误...
    retries += 1;
    if retries > max_retries { return Err(err); }
}

每次循环调用一次模型,如果返回 ContextWindowExceeded,直接返回错误让上层触发压缩后重试;如果是网络抖动等可重试的流式错误,在重试次数内重新发起请求,重试时重新读取历史构建 prompt(而不是复用上次的)。

ToolRouter:工具调度

每次采样前,built_tools 构建一个 ToolRouter

pub struct ToolRouter {
    registry: ToolRegistry,
    model_visible_specs: Vec<ToolSpec>,
}

model_visible_specs 是发给模型的工具描述列表,registry 负责实际调用。工具来源有四类:内置工具(shell 命令执行、文件读写、apply_patch 等)、MCP 工具(通过 MCP 协议接入的外部工具)、Dynamic 工具(运行时动态注册)、Extension 工具(插件提供的)。

模型返回工具调用请求后,ToolCallRuntime 并发执行多个工具调用(如果模型支持 parallel_tool_calls),执行结果写回历史,触发下一轮采样——这个"采样→工具调用→采样"的循环一直持续到模型不再输出工具调用为止。

流式事件处理

模型的流式响应通过 ResponseEvent 枚举分发:

ResponseEvent::OutputTextDelta(delta) => {
    // 增量文字,推送 AgentMessageContentDeltaEvent 给 EQ
}
ResponseEvent::ToolCallInputDelta { call_id, delta } => {
    // 工具调用参数增量,累积到 ToolCallRuntime
}
ResponseEvent::OutputItemDone(item) => {
    // 一个完整的输出项(文字或工具调用)完成
    // 如果是工具调用,提交到 ToolCallRuntime 执行
}
ResponseEvent::Completed { usage, .. } => {
    // 模型输出完成,记录 token 用量
}

文字增量通过 AssistantTextStreamParser 解析,负责从流式文字中识别引用(citations)、Plan 区块(规划模式下)等结构。工具调用参数通过 FuturesOrdered 并发处理,保证执行顺序和结果收集的正确性。

整个 Turn 的事件流是:TurnStartedAgentMessage delta × N → ExecCommandBegin/OutputDelta/End(如果有命令执行)→ TokenCountTurnComplete。如果中途触发了上下文压缩,会插入 ContextCompacted 事件,然后 Turn 继续。

六、上下文压缩:两种实现

Agent 对话越来越长,迟早会触及模型的上下文窗口限制。Codex 的解决方案是压缩(Compact):把历史对话总结成一段文字,替换掉原始历史,释放上下文空间。

Codex 有两种压缩实现:inline 压缩compact.rs)和 remote v2 压缩compact_remote_v2.rs)。选择哪种取决于模型提供商是否支持服务端压缩:

pub(crate) fn should_use_remote_compact_task(provider: &ModelProviderInfo) -> bool {
    provider.supports_remote_compaction()
}

Inline 压缩

Inline 压缩的流程:用模型本身来总结历史,再用总结替换原始历史。

压缩完成后,build_compacted_history 构建新的历史:

fn build_compacted_history_with_limit(
    mut history: Vec<ResponseItem>,
    user_messages: &[String],
    summary_text: &str,
    max_tokens: usize,
) -> Vec<ResponseItem> {
    // 从最新的用户消息往前选,直到 token 预算耗尽
    let mut selected_messages: Vec<String> = Vec::new();
    if max_tokens > 0 {
        let mut remaining = max_tokens;
        for message in user_messages.iter().rev() {
            // ...选取或截断消息...
        }
        selected_messages.reverse();
    }
    // 把选中的用户消息加进去,最后加上总结
    for message in &selected_messages {
        history.push(ResponseItem::Message { role: "user".to_string(), /* ... */ });
    }
    history.push(ResponseItem::Message {
        role: "user".to_string(),
        content: vec![ContentItem::InputText { text: summary_text.to_string() }],
        /* ... */
    });
    history
}

压缩后的历史结构是:保留最近的用户消息(最多 COMPACT_USER_MESSAGE_MAX_TOKENS = 20000 tokens),再加上一条总结消息。总结以 SUMMARY_PREFIX 开头,后续代码通过 is_summary_message 识别它,避免把总结消息当成真实用户消息再次压缩。

压缩有两种触发时机,对应不同的 InitialContextInjection 策略:

pub(crate) enum InitialContextInjection {
    // 中途压缩:在最后一条用户消息前注入初始上下文
    BeforeLastUserMessage,
    // 手动/Turn 前压缩:不注入,下一个 Turn 会重新注入
    DoNotInject,
}

这个区分是必要的:中途压缩发生在一个 Turn 正在进行时(比如上下文窗口快满了),模型需要看到完整的初始上下文(系统提示、工具描述等),所以要注入;而 Turn 前的手动压缩不需要,因为下一个 Turn 开始时会自动重建初始上下文。

Remote V2 压缩

Remote v2 压缩把压缩工作交给服务端处理。它在发送给模型的 prompt 末尾追加一个特殊的触发项:

let mut input = prompt_input.clone();
input.push(ResponseItem::CompactionTrigger);

服务端识别到 CompactionTrigger,返回一个 ResponseItem::Compaction(包含加密的压缩内容)。客户端收到后,用 build_v2_compacted_history 构建新历史:只保留 user/developer/system 角色的消息,过滤掉 assistant 消息和 function call,再加上服务端返回的压缩项,最终保留不超过 RETAINED_MESSAGE_TOKEN_BUDGET = 64000 tokens 的内容。

Remote v2 的优势是压缩质量更高(服务端有更多上下文信息),缺点是依赖服务端支持。两种实现共用同一套 hook 机制(run_pre_compact_hooks/run_post_compact_hooks),用户可以在压缩前后插入自定义逻辑。

压缩完成后,无论哪种实现,都会调用 sess.replace_compacted_history 替换历史,然后 sess.recompute_token_usage 重新计算 token 用量,最后通过 EQ 推送 ContextCompacted 事件。

七、整体架构观察

把这三块放在一起看,Codex 的架构有几个一致的设计原则:

异步优先。SQ/EQ 的分离、Begin/Delta/End 的事件三段式、审批的异步回路,都是为了让 Agent 的长时间运行不阻塞客户端,让 UI 能实时响应。

状态可恢复。rollout 持久化、Thread fork、ThreadRollback Op,说明 Codex 把会话状态的持久化和可回溯当成一等公民。这对一个会执行文件修改和命令的 Agent 来说是必要的安全网。

权限分层AskForApproval 有五种策略(UnlessTrustedOnFailureOnRequestGranularNever),SandboxPolicy 有四种(DangerFullAccessReadOnlyExternalSandboxWorkspaceWrite),GranularApprovalConfig 可以细粒度控制哪类操作需要审批。权限控制的粒度远超"要不要审批"这个二元问题。

协议向后兼容task_started/turn_started 的双名支持,Op 枚举的 #[non_exhaustive],这些细节说明这套协议在演进中,且设计者有意维护兼容性。