运行时任务拓扑
启动 codex 后,CLI、TUI、App Server 和 Core 不一定分别占用一个进程。默认交互模式下, 它们实际共享同一个 OS 进程和 Tokio runtime;连接本地 daemon 或远端 App Server 时,协议边界 才同时成为进程边界。codex exec 也使用嵌入式 App Server,并不会为了进入非交互模式而再 启动一个 App Server 子进程。
理解这套拓扑时要始终区分五种对象:
| 边界 | 源码中的表现 | 失败或退出意味着什么 |
|---|---|---|
| OS 进程 | codex、daemon、被执行命令、sandbox/helper | 地址空间、环境变量和退出码彼此隔离 |
| OS 线程 | codex-main、Tokio worker | 同一进程内执行同步栈或轮询 future |
| Tokio task | client worker、processor、listener、Turn task | 可独立等待、取消或 abort,但不形成内存隔离 |
| channel | mpsc、async_channel、watch、broadcast、oneshot | 决定背压、丢弃、通知和完成确认语义 |
| 协议连接 | 内存中的 typed JSON-RPC、UDS WebSocket、TCP WebSocket、stdio | 决定请求与事件的编码和连接生命周期 |
Crate依赖边界 描述编译期依赖;本文描述运行时 所有权。两者的箭头不能互换:crate A 依赖 crate B,不代表 A 和 B 一定运行在不同进程,也不代表事件只会从 A 流向 B。
本文负责产品入口到 Core、App Server、执行服务和 OS 子进程的跨层拓扑;Core 内部 Session loop → SessionTask → Turn task 的字段和派生关系由 Core异步任务拓扑 展开。读者定位到 Core task owner 后应转入 Core异步任务拓扑,不在本篇继续下钻每种 Task 实现。
1. 所有常规入口
codex 多工具入口、独立 codex-tui、codex-exec 和 codex-app-server 的 main 都调用 arg0_dispatch_or_else。它先处理 helper 重入,再创建名为 codex-main 的 OS 线程,并在线程内 建立多线程 Tokio runtime。顶层 future 由该线程执行,Tokio worker 负责调度其派生 task。
下面节选自 codex-rs/arg0/src/lib.rs 的 arg0_dispatch_or_else 和 build_runtime;中文注释为 本文补充:
源码位置:codex-rs/arg0/src/lib.rs :: arg0_dispatch_or_else, build_runtime
// 所有常规入口先进入 codex-main 线程,再由同一个 Tokio runtime 承载异步任务。
const TOKIO_WORKER_STACK_SIZE_BYTES: usize = 16 * 1024 * 1024;
pub fn arg0_dispatch_or_else<F, Fut>(main_fn: F) -> anyhow::Result<()>
where
F: FnOnce(Arg0DispatchPaths) -> Fut + Send + 'static,
Fut: Future<Output = anyhow::Result<()>>,
{
// 在创建线程前处理 alias/helper,并保留临时 PATH 目录的生命周期。
let path_entry_guard = arg0_dispatch();
let current_exe = std::env::current_exe().ok();
// 顶层 future 不使用调用 main 的原始线程栈,而是进入 16 MiB 的 codex-main。
let handle = std::thread::Builder::new()
.name("codex-main".to_string())
.stack_size(TOKIO_WORKER_STACK_SIZE_BYTES)
.spawn(move || {
// runtime 和后续产品 task 都仍属于当前 OS 进程。
let runtime = build_runtime()?;
runtime.block_on(run_main_with_arg0_guard(
path_entry_guard,
current_exe,
main_fn,
))
})?;
// main 线程同步等待 codex-main,panic 也按原载荷继续传播。
match handle.join() {
Ok(result) => result,
Err(payload) => std::panic::resume_unwind(payload),
}
}
fn build_runtime() -> anyhow::Result<tokio::runtime::Runtime> {
let mut builder = tokio::runtime::Builder::new_multi_thread();
builder.enable_all(); // 启用 I/O、时间和信号驱动。
builder.thread_stack_size(TOKIO_WORKER_STACK_SIZE_BYTES);
Ok(builder.build()?)
}因此,CLI 的 subcommand dispatch 主要是同一进程内的 Rust 函数分派。例如 cli/src/main.rs::cli_main 直接 await codex_exec::run_main、 codex_app_server::run_main_with_transport_options 或 TUI 的 run_main。只有 daemon 管理器明确 Command::spawn,或者执行、沙箱与 helper 路径创建新程序时,才新增 OS 进程。
arg0_dispatch 是这条常规路径之前的例外。Unix exec helper、文件系统 helper、 apply_patch 和 sandbox alias 可以让同一可执行文件以另一种身份重入。某些 helper 建立 current-thread runtime,Unix exec helper 更进一步调用 exec 替换自身进程映像;这些路径不能 画成常规 multi-thread runtime 下的普通 Tokio task。
2. TUI部署形态
TUI 使用 AppServerTarget 明确表示三种运行位置。Embedded 在当前 TUI 进程中启动 App Server task;LocalDaemon 通过 Unix socket 连接本机独立进程;Remote 通过 UDS 或 WebSocket 连接显式端点,并把工作区视为服务端工作区。
codex-rs/tui/src/lib.rs 中的选择与启动代码如下:
源码位置:codex-rs/tui/src/lib.rs :: AppServerTarget
pub(crate) enum AppServerTarget {
Embedded,
// daemon 位于本机其他进程,但使用 RemoteAppServerClient 传输实现。
LocalDaemon { endpoint: RemoteAppServerEndpoint },
// 显式远端可能在本机,也可能在另一台主机。
Remote { endpoint: RemoteAppServerEndpoint },
}
fn app_server_target_for_launch(
explicit_remote_endpoint: Option<RemoteAppServerEndpoint>,
default_daemon_socket: Option<AbsolutePathBuf>,
can_reuse_implicit_local_daemon: bool,
) -> AppServerTarget {
match explicit_remote_endpoint {
// 显式 --remote 的优先级最高。
Some(endpoint) => AppServerTarget::Remote { endpoint },
None if can_reuse_implicit_local_daemon => {
// 仅在 daemon socket 实际可连接时复用;否则回退到嵌入式。
default_daemon_socket.map_or(AppServerTarget::Embedded, |socket_path| {
AppServerTarget::LocalDaemon {
endpoint: RemoteAppServerEndpoint::UnixSocket { socket_path },
}
})
}
None => AppServerTarget::Embedded,
}
}
match target {
// 创建当前进程内的 runtime host 和桥接 task。
AppServerTarget::Embedded => start_embedded_app_server(
arg0_paths,
config,
cli_kv_overrides,
loader_overrides,
strict_config,
cloud_config_bundle,
feedback,
log_db,
state_db,
environment_manager,
)
.await
.map(AppServerClient::InProcess),
// 两种跨进程部署共用 WebSocket/UDS remote client。
AppServerTarget::LocalDaemon { endpoint } | AppServerTarget::Remote { endpoint } => {
connect_remote_app_server(endpoint.clone()).await
}
}隐式 daemon 复用还有四项保护条件:当前 invocation 没有 -c 覆盖、loader override 为默认值、 没有启用 strict config,也没有 hook trust bypass 或 PSP 这类无法重放的启动覆盖。原因不是连接 能力不足,而是已运行 daemon 无法重新采用本次 TUI 的完整进程级配置。Unix 上默认 socket 探测超时只有 50 ms;非 Unix 实现直接返回 None,因此不会隐式采用本地 daemon。
三种 target 最终都被包装为 AppServerSession,TUI 上层始终发 typed App Server request、接收 AppServerEvent。这使 UI 不需要为部署方式维护两套业务逻辑,但资源所有权不同:关闭 Embedded 会结束同进程 App Server;关闭 remote client 只会关闭 WebSocket,不会停止 daemon 或远端服务。
3. 嵌入式桥接队列
嵌入路径仍保留 JSON-RPC 的 request、notification、server request 和 response envelope, 只是把字节传输替换成 typed value 与内存 channel。它不是 TUI → Core 的直接方法调用。
第一层位于 codex-app-server-client。InProcessAppServerClient 为调用方建立一个 command queue、 一个 event queue 和一个 worker task。worker 同时消费调用方命令与底层 runtime 事件;请求等待 被拆到独立 task,防止一个需要用户响应的请求阻塞事件排空。
源码位置:codex-rs/app-server-client/src/lib.rs :: InProcessAppServerClient
pub struct InProcessAppServerClient {
// TUI 或 exec 把请求、通知、审批回复和 Shutdown 发到这里。
command_tx: mpsc::Sender<ClientCommand>,
// 调用方从这里顺序消费服务端通知与请求。
event_rx: mpsc::Receiver<InProcessServerEvent>,
// 该 task 桥接公开 facade 与底层 InProcessClientHandle。
worker_handle: tokio::task::JoinHandle<()>,
}
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();
// 默认容量来自 App Server transport,目前为 128。
let (command_tx, mut command_rx) = mpsc::channel(channel_capacity);
let (event_tx, event_rx) = mpsc::channel(channel_capacity);
let worker_handle = tokio::spawn(async move {
loop {
tokio::select! {
command = command_rx.recv() => {
if let Some(ClientCommand::Request { request, response_tx }) = command {
let request_sender = request_sender.clone();
// 等待 response 的工作脱离桥接循环,循环可继续接收事件。
tokio::spawn(async move {
let result = request_sender.request(*request).await;
let _ = response_tx.send(result);
});
}
// 其他命令分支负责通知、server request 回复和 Shutdown。
}
event = handle.next_event() => {
// 实际源码在这里执行 lossless/best-effort 分流和 Lagged 记账。
// 消费者关闭时,worker 停止继续转发事件。
let _ = event;
}
}
}
});
Ok(Self { command_tx, event_rx, worker_handle })
}上面的事件分支为突出拓扑省略了完整 match,真实分流由 forward_in_process_event 完成。第二层 位于 codex-app-server::in_process:runtime host 再建立 client、event、processor、outgoing 和 writer 队列,并派生 processor task 与 outbound router task。
图中的所有单元仍在一个 OS 进程、一个 Tokio runtime 内。箭头标注的是主要 queue;初始化、 状态数据库、watch channel、background task 和审批 oneshot 没有全部展开。128 是默认容量, 调用方可传入其他正数,源码会用 max(1) 防止零容量。
4. 背压策略
不同 channel 有意采用不同策略。最重要的分界是:文本与终态必须保持无损,进度类通知可以 在过载时丢弃,要求客户端回答的 server request 不能静默消失。
| 位置 | Channel | 容量 | 满载策略 |
|---|---|---|---|
| client facade command/event | Tokio mpsc | 默认 128 | command 异步等待;event 分为无损与 best-effort |
| in-process runtime 各队列 | Tokio mpsc | 默认 128 | request 返回 overload;部分 notification 丢弃 |
| Core submission | async_channel | 512 | send().await 施加背压 |
| Core event | async_channel | 无界 | producer 不因 UI 变慢而阻塞,listener 必须持续排空 |
| Thread created | Tokio broadcast | 1024 | 落后者收到 Lagged,当前实现记录警告但不重同步 |
| remote client command | Tokio mpsc | 配置值 | 异步等待 |
| remote client event | Tokio unbounded mpsc | 无界 | 连接 task 不因消费速度限流,内存压力由消费者承担 |
codex-rs/app-server-client/src/lib.rs 把以下通知列为 lossless:
源码位置:codex-rs/app-server-client/src/lib.rs :: server_notification_requires_delivery
pub(crate) fn server_notification_requires_delivery(
notification: &ServerNotification,
) -> bool {
matches!(
notification,
// Turn 和 Item 终态是调用方结束等待的权威信号。
ServerNotification::TurnCompleted(_)
| ServerNotification::ItemCompleted(_)
| ServerNotification::ThreadSettingsUpdated(_)
| ServerNotification::ExternalAgentConfigImportCompleted(_)
// 文本、计划与 reasoning delta 丢失会永久破坏可见内容。
| ServerNotification::AgentMessageDelta(_)
| ServerNotification::PlanDelta(_)
| ServerNotification::ReasoningSummaryTextDelta(_)
| ServerNotification::ReasoningTextDelta(_)
)
}
if event_requires_delivery(&event) {
// 无损事件等待队列腾出空间,把压力反向传给 producer。
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)) => {
// best-effort 事件计入 skipped,后续向调用方发送 Lagged 标记。
*skipped_events = skipped_events.saturating_add(1);
if let InProcessServerEvent::ServerRequest(request) = event {
// 审批类请求若无法投递,必须显式拒绝,避免 Core 永久等待。
reject_server_request(*request);
}
ForwardEventResult::Continue
}
Err(mpsc::error::TrySendError::Closed(_)) => {
ForwardEventResult::DisableStream
}
}这段设计保证了“终态优先于吞吐”,但不保证所有遥测或命令输出 delta 都到达 UI。消费者看到 Lagged { skipped } 时,只能知道至少有若干 best-effort 事件丢失,不能据此重建具体内容。 远端 client 当前使用无界 event channel,也没有复用这套本地背压;因此跨进程部署把风险从 “丢 best-effort 事件”换成了“慢消费者可能积累内存”。
5. 独立 App Server
独立服务的 run_main_with_transport_options 首先创建三条容量为 128 的 channel:transport event、outgoing envelope 和 outbound control。stdio、Unix socket、WebSocket 与 remote control acceptor 把连接事件送给 processor task;outbound task 独占连接 writer map,避免慢写操作阻塞 JSON-RPC 请求处理。
源码位置:codex-rs/app-server/src/lib.rs :: run_main transport channels
// transport task 报告连接建立、关闭和收到的 JSON-RPC message。
let (transport_event_tx, mut transport_event_rx) =
mpsc::channel::<TransportEvent>(CHANNEL_CAPACITY);
// MessageProcessor 产生的 response、request 和 notification 进入此队列。
let (outgoing_tx, mut outgoing_rx) =
mpsc::channel::<OutgoingEnvelope>(CHANNEL_CAPACITY);
// processor 通知 outbound task 注册、移除或断开连接。
let (outbound_control_tx, mut outbound_control_rx) =
mpsc::channel::<OutboundControlEvent>(CHANNEL_CAPACITY);
let outbound_handle = tokio::spawn(async move {
let mut outbound_connections = HashMap::<ConnectionId, OutboundConnectionState>::new();
loop {
tokio::select! {
biased;
event = outbound_control_rx.recv() => {
let Some(event) = event else { break; };
// 原始 match 还负责 Opened 与 DisconnectAll;这里展示移除分支。
if let OutboundControlEvent::Closed { connection_id } = event {
outbound_connections.remove(&connection_id);
}
}
envelope = outgoing_rx.recv() => {
let Some(envelope) = envelope else { break; };
// 潜在慢 I/O 只占用 outbound task,不占用 processor loop。
route_outgoing_envelope(&mut outbound_connections, envelope).await;
}
}
}
});这段节选保留了源码中的 channel、task 与 Closed 控制分支;为避免复制大段字段装配, Opened 和 DisconnectAll 分支用注释折叠。processor task 则同时监听 transport event、 shutdown signal、running Turn 数量、remote-control 状态、Thread 创建广播和连接清理 task。
每个已加载 Core Thread 另有一个 listener task。它在 conversation.next_event()、listener control command、取消 oneshot 和自动卸载计时器之间 select!;事件先更新 ThreadState,再 转换为 App Server notification。Thread 同时无订阅者且不活跃 30 分钟后,listener 会触发 卸载。因而“一个连接一个 Thread task”并不准确:同一 Thread 的 listener 可以向多个已订阅 connection 投影事件。
6. Core与Thread
ThreadManager 维护 Arc<CodexThread> map。创建 Thread 时,Session::spawn 建立容量 512 的 submission channel、无界 event channel 和 agent status watch channel,然后派生长期存在的 submission_loop。CodexThread 只是持有 Arc<Session> 与这些 I/O endpoint,不等于 OS 线程。
源码位置:codex-rs/core/src/session/mod.rs :: SessionIo
pub(crate) const SUBMISSION_CHANNEL_CAPACITY: usize = 512;
pub(crate) struct SessionIo {
// App Server 把 Op 封装为 Submission 后写入有界队列。
pub(crate) tx_sub: Sender<Submission>,
// Core Event 使用无界队列,由 App Server listener 持续读取。
pub(crate) rx_event: Receiver<Event>,
// watch 只保留最新 AgentStatus,不保存完整状态历史。
pub(crate) agent_status: watch::Receiver<AgentStatus>,
// Shared future 允许多个调用方等待同一个 session loop 结束。
pub(crate) session_loop_termination: SessionLoopTermination,
}
let (tx_sub, rx_sub) = async_channel::bounded(SUBMISSION_CHANNEL_CAPACITY);
let (tx_event, rx_event) = async_channel::unbounded();
// 中间的 Session::new 装配完整服务依赖,并持有 tx_event;此处折叠参数列表。
let (agent_status_tx, agent_status_rx) = watch::channel(AgentStatus::PendingInit);
// let session = Session::new(..., tx_event, agent_status_tx, ...).await?;
let thread_id = session.thread_id;
let session_for_loop = Arc::clone(&session);
// 一个已运行 Thread 对应一个长期 submission loop task。
let session_loop_handle = tokio::spawn(async move {
submission_loop(session_for_loop, configured_config, rx_sub)
.instrument(info_span!("session_loop", thread_id = %thread_id))
.await;
});
let io = SessionIo {
tx_sub,
rx_event,
agent_status: agent_status_rx,
session_loop_termination: session_loop_termination_from_handle(session_loop_handle),
};submission loop 串行接收 Op,但不会把整个 Turn 内联执行到底。UserInput 最终调用 Session::start_task,它为当前 Turn 创建 CancellationToken、完成通知和 Tokio task,并把 RunningTask 存入 active_turn。实现位于 core/src/tasks/mod.rs,RunningTask 字段定义位于 core/src/state/turn.rs。同一个 Session 正常情况下只持有一个当前 task;新的替换型任务会先 abort_all_tasks。
源码位置:codex-rs/core/src/tasks/mod.rs :: Session::start_task
let cancellation_token = CancellationToken::new();
let done = Arc::new(Notify::new());
// child token 传入具体 Turn 实现,使模型流与工具链可协作取消。
let task_cancellation_token = cancellation_token.child_token();
let session = Arc::clone(self);
let ctx = Arc::clone(&turn_context);
let done_clone = Arc::clone(&done);
let handle = tokio::spawn(async move {
let ctx_for_finish = Arc::clone(&ctx);
let task_result = task_for_run
.run(
Arc::clone(&session),
ctx,
task_input,
task_cancellation_token.child_token(),
)
.await;
// terminal event 之前先尽力刷新 rollout。
let _ = session.flush_rollout().await;
if !task_cancellation_token.is_cancelled() {
session.on_task_finished(Arc::clone(&ctx_for_finish), task_result).await;
}
done_clone.notify_waiters();
});
// AbortOnDropHandle 防止 RunningTask 被丢弃后 task 继续成为孤儿。
let running_task = RunningTask {
done,
handle: AbortOnDropHandle::new(handle),
kind: task_kind,
task,
cancellation_token,
turn_context: Arc::clone(&turn_context),
_agent_execution_guard: agent_execution_guard,
_timer: timer,
};
turn.task = Some(running_task);Core 的“串行”只限定 submission dispatch 和单个 Session 的 active Turn 所有权。Turn task 内部 仍会执行模型流、工具、MCP 与其他异步工作;多个 Thread 也可以各自运行 Turn。业务事件角度的完整 链路见 Turn端到端链路,本文只确定这些task的父子关系。
7. Exec子进程
codex exec-server 由 CLI 在当前 codex 进程内直接进入 run_main_with_telemetry;若它由外部 服务或独立命令启动,整个 codex 自然就是一个单独服务进程。服务支持单连接 stdio 和多连接 WebSocket。每条连接拥有 JSON-RPC reader/writer task、outbound task 和 request dispatcher; concurrent dispatch 启用后,请求再进入受 semaphore 限制的 task lane。
收到 exec 请求后,local backend 调用 codex_sandboxing::spawn_process 创建真正的 OS 子进程, 再分别派生 stdout、stderr 和 exit watcher 三个 task:
源码位置:codex-rs/exec-server/src/local_process.rs :: LocalProcessServer::start_process
async fn start_process(&self, params: ExecParams) -> Result<ExecResponse, JSONRPCErrorError> {
// 请求先完成环境、路径和网络策略准备,成功后才真正创建子进程。
let prepared = prepare_exec_request(
¶ms,
child_env(¶ms),
self.runtime_paths.as_ref(),
network_policy_decider,
).await?;
// 这里跨越 OS 进程边界;sandbox transform 已包含在 SpawnRequest 中。
let spawned = codex_sandboxing::spawn_process(codex_sandboxing::SpawnRequest {
command: &prepared.command,
cwd: prepared.cwd.as_path(),
env: &prepared.env,
arg0: &prepared.arg0,
sandbox: prepared.sandbox,
windows_sandbox: prepared.windows_sandbox_spawn_request(),
tty: params.tty,
stdin_open: params.tty || params.pipe_stdin,
inherited_fds: &[],
}).await.map_err(|err| internal_error(err.to_string()))?;
// 两个输出 task 把字节写入进程事件日志并通知订阅者。
tokio::spawn(stream_output(
process_id.clone(),
if params.tty { ExecOutputStream::Pty } else { ExecOutputStream::Stdout },
spawned.stdout_rx,
Arc::clone(&self.inner),
Arc::clone(&output_notify),
));
tokio::spawn(stream_output(
process_id.clone(),
if params.tty { ExecOutputStream::Pty } else { ExecOutputStream::Stderr },
spawned.stderr_rx,
Arc::clone(&self.inner),
Arc::clone(&output_notify),
));
// 第三个 task 独立等待退出码,避免 read 请求承担 wait 职责。
tokio::spawn(watch_exit(
process_id.clone(), spawned.exit_rx,
Arc::clone(&self.inner), output_notify,
));
Ok(ExecResponse { process_id, sandbox_type })
}这条服务路径不能被概括成“所有 Codex shell 命令都会启动 exec-server 进程”。EnvironmentManager 可以提供本地或远端执行环境;本地 backend 能在宿主进程内管理命令 child,远端环境才跨越 exec-server 协议边界。无论由谁管理,最终用户命令都是独立 OS 进程,Tokio task 只负责控制、 I/O 和状态投影。
daemon 是另一种明确的进程创建。Unix PID backend 以 codex app-server --listen unix:// 启动 detached session,stdin/stdout 置空,stderr 写入日志,并用 PID 与进程启动时间防止误杀 复用 PID 的其他进程。TUI 连接该 socket 后只拥有自己的 WebSocket,不拥有 daemon 进程。
8. 关闭传播
嵌入式关闭有多层有界等待。TUI 或 Exec 先关闭 client facade;facade 在发送 Shutdown 前主动 丢弃自己的 event receiver,以解除可能阻塞在 lossless send().await 上的 worker。上层最多 等待 45 秒。底层 runtime 最多等待 35 秒的关闭确认,再给 runtime task 5 秒 join 时间。
runtime host 接到 Shutdown 后按顺序关闭请求入口、取消未完成 server request、等待 processor、 停止 outbound router、刷新 analytics。processor 清理连接和 listener,等待 background task, 最后并发关闭全部 Core Thread;每个 Thread 的等待上限为 10 秒。
Core Thread 收到 Op::Shutdown 后,session loop 先执行 runtime teardown。活动 Turn 先取消 token, 最多给 100 ms 协作退出时间,然后 abort task;随后终止 unified exec 子进程、关闭 code mode、MCP runtime、Guardian session 与持久化 writer。
独立 App Server 的退出规则不同:stdio 是单客户端模式,连接关闭会结束 processor loop;Unix socket、WebSocket 和 remote-control 模式可安装信号处理。第一次 SIGINT/SIGTERM 请求 graceful restart drain,SIGHUP 只请求 graceful drain;服务仍接受请求,直到 running assistant Turn 数量 降到零。第二个可强制信号会跳过 Thread drain,取消 acceptor 并断开连接。
remote TUI 的 AppServerClient::shutdown 只发送 WebSocket close 并等待本地 worker,最多 5 秒后 abort worker。它不会向 daemon 发送“停止整个服务”的生命周期命令;这是共享服务能够继续服务 其他客户端的必要条件。
9. 卡死定位
运行时故障可以按所有权快速缩小范围:
| 现象 | 优先检查 | 原因 |
|---|---|---|
| TUI 无响应但 daemon 仍可被其他客户端使用 | TUI remote worker 与本连接 | Core 和服务进程可能正常 |
文本完整但命令输出片段缺失并出现 Lagged | in-process event consumer | 命令输出属于可丢弃的 best-effort 层 |
| Turn 已结束但入口仍等待 | TurnCompleted 投影、listener 和 lossless queue | 终态按设计不能丢弃,缺失通常是链路关闭或 listener 故障 |
| Thread 不响应新 Op | 512 容量 submission queue 与 session loop | request 成功不代表 Core loop 已消费 |
| 退出停在数十秒 | App Server background drain、Thread shutdown、analytics flush | 嵌入式关闭是多层嵌套超时,不是单个 5 秒 timeout |
| 命令结束但状态仍 Running | stdout/stderr task 与 exit watcher | 子进程退出、输出排空和状态投影由不同 task 完成 |
| daemon socket 存在但 TUI 回退嵌入式 | 50 ms probe 与配置可重放条件 | socket 文件存在不等于连接成功,也不等于 invocation 可复用 daemon |
升级到后续稳定 tag 时,至少复核以下变化:AppServerTarget 的选择条件、默认 channel 容量、 lossless notification 集合、remote event queue 是否仍无界、Thread 自动卸载时长、Turn 协作取消 窗口,以及 App Server graceful shutdown 的信号与 running-Turn 判定。它们都影响运行时行为, 却不一定体现在 crate 依赖图或公开协议 schema 中。
最终可以把这套拓扑压缩成一句话:产品入口通常共享进程,App Server 统一协议语义,Core 以 Thread session loop 和 Turn task 承载工作,执行层才把命令落实为 OS 子进程;每一层通过不同 channel 和超时拥有自己的背压与回收责任。
10. 拓扑契约测试
静态 spawn() 和 channel 清单只能说明 task 可能存在。下面的测试分别验证过载、drain、生命周期顺序 和卡住后台工作的关闭语义;更完整的 Session 资源账本见 Session关闭流程。
10.1 Best-effort丢弃
forward_in_process_event_preserves_transcript_notifications_under_backpressure 先用 stdout delta 填满 容量 1 队列,再送第二条 stdout delta,要求计入一次 skip。随后并发排空队列,并依次投递 agent text、 item completed 和 turn completed;这些 lossless 事件必须全部进入接收序列。
源码位置:codex-rs/app-server-client/src/lib.rs
// :: forward_in_process_event_preserves_transcript_notifications_under_backpressure(关键路径)
let (event_tx, mut event_rx) = mpsc::channel(1);
event_tx
.send(InProcessServerEvent::ServerNotification(Box::new(
command_execution_output_delta_notification("stdout-1"),
)))
.await
.expect("initial event should enqueue");
let mut skipped_events = 0usize;
let result = forward_in_process_event(
&event_tx,
&mut skipped_events,
InProcessServerEvent::ServerNotification(Box::new(
command_execution_output_delta_notification("stdout-2"),
)),
|_| {},
)
.await;
// best-effort输出在满载时丢弃并记账,producer不等待队列。
assert_eq!(result, ForwardEventResult::Continue);
assert_eq!(skipped_events, 1);
let receive_task = tokio::spawn(async move {
let mut events = Vec::new();
for _ in 0..5 {
events.push(
timeout(Duration::from_secs(2), event_rx.recv())
.await
.expect("event should arrive before timeout")
.expect("event stream should stay open"),
);
}
events
});
for notification in [
agent_message_delta_notification("hello"),
item_completed_notification("hello"),
turn_completed_notification(),
] {
let result = forward_in_process_event(
&event_tx,
&mut skipped_events,
InProcessServerEvent::ServerNotification(Box::new(notification)),
|_| {},
)
.await;
// lossless事件会等待receiver腾出容量,不能走skip分支。
assert_eq!(result, ForwardEventResult::Continue);
}
assert_eq!(skipped_events, 0);测试的完整接收断言还验证 stdout-1 先到、Lagged marker随后出现、三类 transcript/terminal notification 保持顺序。事件分类和远端 transport 对照见 CodexThread背压。
10.2 进程内关闭
shutdown_waits_for_in_process_drain 在暂停时间的 Tokio runtime 中让 worker 收到 Shutdown 后等待 30 秒, 再设置完成标志并回复 oneshot。client.shutdown().await 返回时标志必须已经为 true。
源码位置:codex-rs/app-server-client/src/lib.rs
// :: shutdown_waits_for_in_process_drain(关键路径)
let (command_tx, mut command_rx) = mpsc::channel(1);
let (_event_tx, event_rx) = mpsc::channel(1);
let completed = Arc::new(AtomicBool::new(false));
let worker_completed = Arc::clone(&completed);
let worker_handle = tokio::spawn(async move {
let response_tx = match command_rx.recv().await {
Some(ClientCommand::Shutdown { response_tx }) => response_tx,
_ => panic!("expected shutdown command"),
};
tokio::time::sleep(Duration::from_secs(30)).await;
worker_completed.store(true, Ordering::Release);
let _ = response_tx.send(Ok(()));
});
let client = InProcessAppServerClient {
command_tx,
event_rx,
worker_handle,
};
// shutdown确认来自worker完成,不是command成功入队。
client.shutdown().await.expect("shutdown should complete");
assert!(completed.load(Ordering::Acquire));这证明第8节的外层箭头是 completion barrier。45/35/5 秒等 timeout 是异常兜底,正常路径仍应等待每层 owner 完成自己的 drain。
10.3 Submission关闭
submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle 启动一个永不结束但监听取消 token 的 RegularTask,随后丢弃 submission sender。最终 lifecycle recorder 必须严格得到 turn_abort → thread_stop。
源码位置:codex-rs/core/src/session/tests.rs
// :: submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle(结果路径)
let session = Arc::new(session);
session
.spawn_task(
Arc::new(turn_context),
Vec::new(),
NeverEndingTask {
kind: TaskKind::Regular,
listen_to_cancellation_token: true,
},
)
.await;
let (tx_sub, rx_sub) = async_channel::bounded(1);
// 所有Submission sender消失,session loop把它解释为Thread控制面结束。
drop(tx_sub);
submission_loop(Arc::clone(&session), session.get_config().await, rx_sub).await;
assert_eq!(
vec!["turn_abort", "thread_stop"],
*calls
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
);sender drop 不是普通 Op,但它仍复用完整 runtime teardown。这说明“channel close”属于生命周期输入, 不能只在日志中当成一次 receive error。
10.4 MCP启动阻塞
shutdown_cancels_startup_prewarm_waiting_for_mcp_startup 让 MCP HTTP listener 接受连接但不完成协议, 确保 startup prewarm 被卡住。随后 shutdown_and_wait() 必须在 2 秒内返回,并且再等待 100ms 后模型 WebSocket server 仍没有连接。
源码位置:codex-rs/core/tests/suite/rmcp_client.rs
// :: shutdown_cancels_startup_prewarm_waiting_for_mcp_startup(关键断言)
let (_pending_mcp_connection, _) =
tokio::time::timeout(Duration::from_secs(5), pending_mcp_listener.accept())
.await
.context("startup prewarm should start the MCP connection")??;
// shutdown不能等待永远不ready的MCP startup/prewarm链。
tokio::time::timeout(Duration::from_secs(2), fixture.codex.shutdown_and_wait())
.await
.context("shutdown should not wait for startup prewarm MCP startup")??;
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
server.connections().is_empty(),
"startup prewarm should not send a websocket request after shutdown"
);该测试带 skip_if_no_network!,所以 CI 绿色结果需要结合是否实际执行解释。模型预热与 MCP worker 的 职责差异见 Session启动预热。
10.5 任务拓扑排查
读者应能从一个“退出卡住”现象依次指出:入口 client、App Server host/processor、Thread listener、 Session loop、Turn token 和 OS child 的 owner;并解释哪一层等待 ack、哪一层允许 abort、remote client 为何不能关闭共享 daemon。若已经进入某个 Core task 的字段和状态机,请停止使用本总览,转入 Core异步任务拓扑 或 Session关闭流程。
可以用下面的只读搜索把本文的 task topology and shutdown 主线落回源码:
rg -n "submission_loop|shutdown|capacity" codex-rs/core/src codex-rs/app-server/src