# Codex 源码阅读与分析(一):客户端主流程——从用户输入到事件循环

Codex 源码阅读与分析(一):从用户输入到事件循环

本篇技术博客主要解读 Codex 命令行版本的运行机制。我会沿着一次 codex exec 任务的主流程,从用户输入开始,一直看到事件循环结束以及最终结果输出。

Codex 面向用户的输入方式大体可以分为两类:

  • 通过 codex exec 运行非交互式任务;
  • 通过 TUI 进行交互。

TUI 有自己的入口,位于 codex-rs/tui/src/lib.rs。本篇文章主要关注 codex exec,对应代码位于 codex-rs/exec/src/lib.rs。两者虽然共享底层的 App Server 等基础设施,但入口并不相同。
codex exec 的主入口是:

1
pub async fn run_main(cli: Cli, arg0_paths: Arg0DispatchPaths) -> anyhow::Result<()>

run_main 完成命令行参数、配置、日志以及运行环境等初始化后,最终进入 run_exec_session。我们主要关注的会话逻辑就在这个函数里面。

1. Prompt 的解析

用户的 Prompt 主要通过 codex-rs/exec/src/lib.rs 中的 resolve_root_promptresolve_prompt 两个函数进行解析:

1
2
3
fn resolve_root_prompt(prompt_arg: Option<String>) -> String

fn resolve_prompt(prompt_arg: Option<String>) -> String

普通的 codex exec 根路径使用 resolve_root_promptcodex exec resume 路径则直接使用 resolve_prompt。当命令行没有提供 Prompt,或者显式传入 - 时,Codex 会从 stdin 读取输入:

1
let Some(prompt) = read_prompt_from_stdin(behavior)

resolve_root_prompt 除了处理命令行参数,也可以把管道输入的 stdin 内容作为额外上下文追加进去。

总体来说,这几个函数解决的是同一件事:确定本次任务最终使用的 Prompt 文本来自命令行参数、标准输入,还是两者的组合。

2. 使用 UserInput 存储用户输入

Prompt 解析完成后会变成 prompt_text。对于用户输入,Codex 准备了一个 UserInput 枚举进行统一存储:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
pub enum UserInput {
Text {
text: String,
/// UI-defined spans within `text` that should be treated as special elements.
/// These are byte ranges into the UTF-8 `text` buffer and are used to render
/// or persist rich input markers (e.g., image placeholders) across history
/// and resume without mutating the literal text.
#[serde(default)]
text_elements: Vec<TextElement>,
},
/// Pre-encoded data: URI image.
Image {
image_url: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional)]
detail: Option<ImageDetail>,
},

/// Local image path provided by the user. This will be converted to an
/// `Image` variant (base64 data URL) during request serialization.
LocalImage {
path: std::path::PathBuf,
#[serde(default, skip_serializing_if = "Option::is_none")]
#[ts(optional)]
detail: Option<ImageDetail>,
},

/// Skill selected by the user (name + path to SKILL.md).
Skill {
name: String,
path: std::path::PathBuf,
},
/// Explicit structured mention selected by the user.
///
/// `path` identifies the exact mention target, for example
/// `app://<connector-id>` or `plugin://<plugin-name>@<marketplace-name>`.
Mention { name: String, path: String },
}

这里的 TextElement 设计得很好。它不会直接修改原始文本,而是通过 UTF-8 字节范围记录图片占位符等特殊元素在文本中的位置。
不过,本篇文章暂不对这部分展开分析。

3. 封装为 InitialOperation

解析后的输入会进一步封装到 InitialOperation 中:

1
2
3
4
5
6
7
8
9
enum InitialOperation {
UserTurn {
items: Vec<UserInput>,
output_schema: Option<Value>,
},
Review {
review_request: ReviewRequest,
},
}

InitialOperation 主要包括两种操作:

  • UserTurn:普通用户任务,用户输入就存放在这里;
  • Review:启动专门的代码审查任务。

我们这里主要关注 Codex 的普通执行流程,因此暂时不展开 Review

二、会话启动前的检查与初始化

1. 检查当前工作目录

封装完成后,Codex 会检查当前工作目录是否位于一个 Git 仓库中:

1
2
3
4
5
6
7
if !skip_git_repo_check
&& !dangerously_bypass_approvals_and_sandbox
&& get_git_repo_root(&default_cwd).is_none()
{
eprintln!("Not inside a trusted directory and --skip-git-repo-check was not specified.");
std::process::exit(1);
}

只有在没有跳过 Git 检查、也没有开启绕过审批与沙箱的危险模式时,这项检查才会生效。

2. 启动进程内 App Server Client

接下来开始建立客户端与 App Server 之间的通信:

1
2
let mut request_ids = RequestIdSequencer::new();
let mut client = InProcessAppServerClient::start(in_process_start_args).await?;

这里使用的是 InProcessAppServerClient,也就是运行在同一进程内部的 App Server Client。它会在调用方和底层 App Server 之间转发请求、响应以及事件。

3. 创建或恢复 Thread

启动客户端以后,Codex 会确定本次任务使用哪个 Thread:

  • 如果执行的是 codex exec resume,就通过现有的 thread/listthread/resume API 找到并恢复旧 Thread;
  • 否则通过 thread/start 创建一个新 Thread。

得到 Thread 信息后,Codex 会把 Thread ID 记录到当前的 tracing span 中:

1
exec_span.record("thread.id", primary_thread_id_for_span.as_str());

三、启动 Turn 与中断监听

1. 监听 Ctrl+C

Thread 准备完成后,Codex 会单独启动一个异步任务,用来监听用户是否按下 Ctrl+C

1
2
3
4
5
6
7
let (interrupt_tx, mut interrupt_rx) = mpsc::unbounded_channel::<()>();
tokio::spawn(async move {
if tokio::signal::ctrl_c().await.is_ok() {
tracing::debug!("Keyboard interrupt");
let _ = interrupt_tx.send(());
}
});

这里没有直接结束进程,而是向 interrupt_rx 对应的通道发送一个中断信号。真正的中断请求会在后面的事件循环里处理。

2. 根据 InitialOperation 启动任务

接下来,Codex 根据 InitialOperation 的类型启动普通 Turn 或 Review。我们仍然只看普通的 UserTurn 路径。

Turn ID 由 App Server 返回:

1
2
3
4
5
// 省略 send_request_with_response 的具体参数。
let response: TurnStartResponse = send_request_with_response(/* ... */).await?;
// 首先通过 client 开启一个 Turn,
// 然后从服务端返回的 response 中取得 ID,作为 task_id。
let task_id = response.turn.id;

到这里,我们已经拿到了两个关键 ID:

  • primary_thread_id:定位会话;
  • task_id:定位本次 Turn。

后续不论是处理中断、过滤通知还是判断任务完成状态,基本都会围绕这两个 ID 进行。

自此,任务的初始化环节结束,程序正式进入事件循环。

四、事件循环

事件循环的完整代码比较长,不适合全部放到正文中。这里主要拆解其中最关键的几类分支。

1. 处理 Ctrl+C

循环首先同时监听两类消息:

  • 用户是否按下 Ctrl+C
  • App Server 是否产生了新事件。

如果收到中断信号,客户端会向 App Server 发送 TurnInterrupt 请求:

1
2
3
4
send_request_with_response::<TurnInterruptResponse>(
&client,
ClientRequest::TurnInterrupt(...),
)

这个请求中会携带当前的 thread_idturn_id,用于中断当前 Turn。需要注意的是,发送请求后程序不会马上退出事件循环,而是会继续接收 App Server 的后续事件,从而确认任务最终是完成、失败还是已经中断。

如果 maybe_interrupt.is_none(),代表中断监听通道已经关闭。此时程序只会停止监听这个通道,并继续处理服务器事件。

2. 处理 ServerRequest

第一类服务端事件是 ServerRequest

1
2
3
InProcessServerEvent::ServerRequest(request) => {
handle_server_request(&client, request, &mut error_seen).await;
}

ServerRequest 表示服务端主动向客户端发起、并且需要客户端返回响应的请求,例如命令执行审批、文件修改审批、工具向用户请求输入、MCP elicitation 以及动态工具调用等。

这些请求统一交给 handle_server_request 处理。这个流程很重要,但是我们暂时不展开。

3. 处理 ServerNotification

第二类事件是 ServerNotification

1
2
3
InProcessServerEvent::ServerNotification(mut notification) => {
// ...
}

它表示服务端主动向客户端发送的通知,其中包括运行过程中出现的错误、Item 状态变化、流式输出以及 Turn 完成信息等。

首先,事件循环会根据 ErrorTurnCompleted 的内容判断本次任务是否已经失败或中断:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
InProcessServerEvent::ServerNotification(mut notification) => {
if let ServerNotification::Error(payload) = &notification {
if payload.thread_id == primary_thread_id_for_requests
&& payload.turn_id == task_id
&& !payload.will_retry
{
error_seen = true;
}
} else if let ServerNotification::TurnCompleted(payload) = &notification
&& payload.thread_id == primary_thread_id_for_requests
&& payload.turn.id == task_id
&& matches!(
payload.turn.status,
codex_app_server_protocol::TurnStatus::Failed
| codex_app_server_protocol::TurnStatus::Interrupted
)
{
error_seen = true;
}

这里匹配到的 notification 使用了 mut,是因为后面还可能对它进行补全:

1
2
3
4
5
6
7
maybe_backfill_turn_completed_items(
config.ephemeral,
&client,
&mut request_ids,
&mut notification,
)
.await;

这一部分属于尽力而为的数据修复。当非临时 Thread 的 TurnCompleted 没有携带完整 Items 时,Codex 会通过一次 thread/read 尝试补回最终结果。

注意,这种补全恢复的是 Thread 中已经持久化的最终状态,并不等于把之前丢失的每一条流式事件重新播放一遍。

补全结束后,should_process_notification 会判断当前通知是否属于本次 CLI 任务。因为事件流中可能存在其他 Thread、其他 Turn 或全局级别的通知,所以这里需要根据 thread_idturn_id 进行过滤。

最后,符合条件的通知会交给:

1
event_processor.process_server_notification(notification)

它负责输出事件,并根据返回的 CodexStatus 决定继续运行,还是开始结束整个执行会话。

4. 处理 Lagged

最后一类事件是 Lagged

1
2
3
4
5
InProcessServerEvent::Lagged { skipped } => {
let message = lagged_event_warning_message(skipped);
warn!("{message}");
event_processor.process_warning(message);
}

Lagged 表示内部事件消费者跟不上事件生成速度,已经有部分事件因为背压被跳过。这里会生成一条警告,既写入日志,也交给 event_processor 展示。

需要特别注意:这里不会逐条重新获取被跳过的原始事件。 原因不在这个 match 分支本身,而在更上游的事件转发逻辑。

4.1 事件队列本身是有界的

codex-rs/app-server-client/src/lib.rs 中,事件通道按照 channel_capacity 创建:

1
2
3
4
5
6
7
pub async fn start(args: InProcessClientStartArgs) -> IoResult<Self> {
let channel_capacity = args.channel_capacity.max(1);
let mut handle =
codex_app_server::in_process::start(args.into_runtime_start_args()).await?;
let request_sender = handle.sender();
let (command_tx, mut command_rx) = mpsc::channel::<ClientCommand>(channel_capacity);
let (event_tx, event_rx) = mpsc::channel::<InProcessServerEvent>(channel_capacity);

也就是说,队列最多只能保存 channel_capacity 个尚未被消费的事件。这样可以让内存占用保持有界,而不是在消费者变慢时无限堆积消息。

4.2 事件被分为无损和尽力而为两类

上游首先通过 event_requires_delivery 判断事件是否必须送达:

1
2
3
4
5
6
7
8
9
10
11
12
fn event_requires_delivery(event: &InProcessServerEvent) -> bool {
// These transcript and terminal events must remain lossless. Dropping
// streamed assistant text or the authoritative completed item can leave
// the TUI with permanently corrupted markdown, while dropping completion
// notifications can leave surfaces waiting forever.
match event {
InProcessServerEvent::ServerNotification(notification) => {
server_notification_requires_delivery(notification)
}
_ => false,
}
}

必须送达的事件包括助手文本、计划、推理增量,以及 ItemCompletedTurnCompleted 等关键通知。这些事件使用等待式的 send:队列满时会等待消费者腾出空间,而不是直接丢弃。

其他事件则属于尽力而为类型,例如高频的命令输出增量或进度通知。它们使用非阻塞的 try_send

4.3 队列满时,尽力而为事件会被明确丢弃

最直接的证据仍然位于 codex-rs/app-server-client/src/lib.rs

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
if event_requires_delivery(&event) {
// Block until the consumer catches up for transcript/completion notifications; this
// preserves the visible assistant output even when the queue is otherwise saturated.
if event_tx.send(event).await.is_err() {
return ForwardEventResult::DisableStream;
}
return ForwardEventResult::Continue;
}

match event_tx.try_send(event) {
Ok(()) => ForwardEventResult::Continue,
Err(mpsc::error::TrySendError::Full(event)) => {
*skipped_events = skipped_events.saturating_add(1);
warn!("dropping in-process app-server event because consumer queue is full");
if let InProcessServerEvent::ServerRequest(request) = event {
reject_server_request(request);
}
ForwardEventResult::Continue
}
Err(mpsc::error::TrySendError::Closed(_)) => ForwardEventResult::DisableStream,
}

这里的执行流程非常明确:

  1. 使用 try_send(event) 尝试把尽力而为事件放入队列;
  2. 队列已满时进入 TrySendError::Full(event) 分支;
  3. skipped_events 加一;
  4. 输出 dropping ... because consumer queue is full 警告;
  5. 返回 ForwardEventResult::Continue,继续处理后续事件,并没有重试或保存当前事件。

如果无法入队的是 ServerRequest,代码还会调用 reject_server_request(request) 明确拒绝它,避免 App Server 永远等待一个不可能到来的响应。

因此,当下游最终收到 Lagged { skipped } 时,它拿到的只有“丢失了多少条事件”这个计数,并没有那些事件的完整内容、事件 ID 或可供重放的游标。此时自然无法在 Lagged 分支中把它们逐条取回来。

这套设计的核心取舍是:让可降级的实时通知在过载时可以被丢弃,对无法入队的 ServerRequest 给出明确拒绝,换取内存有界,并避免慢消费者拖住整个执行流程;同时对文本内容和完成状态等关键事件保持可靠送达。

4.4 Lagged 上报的是累计丢失数量

被跳过的事件数量会累计在 skipped_events 中。转发下一个事件时,如果队列已经有空间,上游会尝试先发送 Lagged;如果下一个事件属于必须送达类型,则会等待容量并保证先发出这个标记。计数会被封装成:

1
2
3
InProcessServerEvent::Lagged {
skipped: *skipped_events,
}

对应的事件定义位于 codex-rs/app-server/src/in_process.rs

1
2
3
4
5
6
7
8
9
10
11
12
13
/// Event emitted from the app-server to the in-process client.
///
/// [`Lagged`](Self::Lagged) is a transport health marker, not an application
/// event — it signals that the consumer fell behind and some events were dropped.
#[derive(Debug, Clone)]
pub enum InProcessServerEvent {
/// Server request that requires client response/rejection.
ServerRequest(ServerRequest),
/// App-server notification directed to the embedded client.
ServerNotification(ServerNotification),
/// Indicates one or more events were dropped due to backpressure.
Lagged { skipped: usize },
}

源码注释也明确说明,Lagged 是一个传输健康标记,而不是原始业务事件。它只能告诉下游“消费者落后了,并且已经丢失若干事件”,无法携带和恢复那些事件的具体内容。

五、会话结束

当事件处理器根据完成通知返回 CodexStatus::InitiateShutdown 后,Codex 会请求取消 Thread 订阅,并退出事件循环。随后关闭进程内客户端:

1
2
3
if let Err(err) = client.shutdown().await {
warn!("in-process app-server shutdown failed: {err}");
}

最后通过:

1
event_processor.print_final_output();

输出本次任务的最终结果。如果此前记录到了不可重试的错误、失败或中断状态,进程还会以非零状态码退出,方便脚本和自动化系统判断执行结果。