Skip to content

Core异步任务拓扑

逐项梳理 Codex Core 的常驻与临时异步任务、channel 类型、容量、背压、取消和关闭顺序。

基于rust-v0.150.0
CodexRustTokioConcurrency

Core异步任务拓扑 ​

Codex Core 不是一条从 submit() 调到模型再返回的同步调用栈。每个 loaded Thread 至少拥有一个长期 submission loop;活动 Turn 在独立 Tokio task 中执行;模型 response stream、工具调用、rollout writer、MCP 预热、unified exec 输出和 Realtime 都有自己的 task 或 channel。

这些异步边界的语义并不相同:bounded queue 会给生产者背压,unbounded event queue 把压力转移到 内存,watch 只保存最新状态,broadcast 允许多个订阅者但慢消费者可能 lag,oneshot 表达一个等待中 的交互,CancellationToken 只发取消信号而不负责 join。理解拓扑必须同时记录创建者、容量、消费者和 关闭者。

阅读前先看 运行时任务拓扑 区分 OS 进程、Tokio task 与协议连接,再用 Core运行时架构总览 定位 Core 职责域。本文从 SessionIo/submission_loop 开始,只展开 Core 内部异步 owner。

1. Task队列通知 ​

原语Core 中的用途是否保存所有消息满载/重复时的语义
async_channel::boundedsubmission、MCP refresh、delegate bridge、Realtime input是,直到容量send().await 背压;try_send 可合并/丢弃
async_channel::unboundedCore → host Event、部分 delegate event是不阻塞生产者,压力进入内存
tokio::sync::mpscmodel stream、rollout writer是,直到容量sender await,形成背压
watchagent/process/input activity 最新状态否,只保留最新值中间状态可被覆盖
broadcastthread-created、多订阅者 exec output环形缓冲慢 receiver 得到 Lagged
oneshotapproval、flush、last response、一次性完成单值sender/receiver drop 表示取消
Notifytask done、output drained、数据到达不保存业务消息permit/唤醒语义,不是队列
CancellationTokentask、sampling、tool、MCP、Realtime 取消树不保存消息cooperative;仍需 await/abort 清理

channel 容量不能脱离发送方式解释。例如 MCP prewarm 的容量是 1,但生产者使用 try_send(()),满时 忽略新请求,从而把多次 refresh 合并成“一次待处理”;Realtime audio 同样使用 try_send,满时直接 丢 frame。相反,submission 和 model event 使用 send().await,容量满会暂停生产者。

2. Thread 级总拓扑 ​

图中没有把每个函数画成 task。run_turn()、run_sampling_request() 等 async 函数通常运行在当前 SessionTask 的 Tokio task 内;只有明确 tokio::spawn、runtime spawn 或独立 writer/transport owner 的位置才算新的并发执行主体。

总拓扑里的 queue 必须能追溯到 owner。下面的类图只画关键所有权关系,区分“持有 sender/receiver”与 “拥有后台 task”两种不同责任。

SessionIo 的 completion future 只是观察 session loop;真正拥有 Turn task 的是 Session.active_turn, 真正拥有文件 writer 的则是 RolloutRecorder。当前 Session 还拥有 MCP prewarm task、网络 proxy refresh 和 hook/Guardian 相关状态;关闭顺序必须沿各 owner 关系反向展开。

3. Channel容量 ​

Channel类型容量消息创建者主要消费者/关闭者
Thread createdbroadcast1024ThreadIdThreadManager::newApp Server/listener;manager drop 关闭
Session submissionasync_channel::bounded512SubmissionSession::spawn_internalsubmission loop;Shutdown 或所有 sender drop
Session eventasync_channel::unbounded无界EventSession::spawn_internalCodexThread::next_event;Session sender drop
Agent statuswatch最新 1 个AgentStatusSession::spawn_internalstatus subscriber;Session sender drop
Input activitywatch最新 1 个Mailbox/SteerInputQueue::newsampling wait;InputQueue drop
MCP prewarmasync_channel::bounded1()Session::newprewarm worker;cancel token + join
Auth changeswatch最新 1 个revision u64AuthManagerMCP prewarm worker;sender/worker drop
Model eventsTokio mpsc1600Result<ResponseEvent>map_response_eventsTurn sampling;stream drop cancels mapper
Last responseoneshot1LastResponsemodel mapperModelClientSession sticky state
Rollout commandsTokio mpsc256RolloutCmdRolloutRecorder::newrollout writer;Shutdown command + ack
Unified exec outputbroadcast64Vec<u8>UnifiedExecProcess::newoutput consumers/watchers;process drop
Unified exec statewatch最新 1 个ProcessStateUnifiedExecProcess::newwait/read/write paths;process drop
Realtime audioasync_channel::bounded256RealtimeAudioFrameconversation startinput task;full 时丢 frame
Realtime textasync_channel::bounded64ConversationTextParamsconversation startinput task;sender await
Realtime handoffasync_channel::bounded64RealtimeOutboundconversation startinput task;sender await
Realtime eventsasync_channel::bounded256RealtimeEventconversation startfanout/client;sender await
Transcript tailasync_channel::bounded1Stringconversation startclose/final consumer
Turn interactiononeshot1/请求approval/input/permission/elicitation response发起请求的 handler当前 Turn waiter;终止时 clear/drop

Notify、shared completion future 和 cancellation token 没有“容量”,所以不应为了凑表格把它们写成 channel。它们的创建与所有权在后文单独说明。

4. Session常驻task ​

Session::spawn_internal() 建立两个方向相反的 channel:

源码位置:codex-rs/core/src/session/mod.rs :: SessionIo, Session::spawn_internal

rust
pub(crate) const SUBMISSION_CHANNEL_CAPACITY: usize = 512;

pub(crate) struct SessionIo {
    pub(crate) tx_sub: Sender<Submission>,
    pub(crate) rx_event: Receiver<Event>,
    pub(crate) agent_status: watch::Receiver<AgentStatus>,
    // Shared future 允许多个持有者等待同一次 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();

let io = SessionIo {
    tx_sub,
    rx_event,
    agent_status: agent_status_rx,
    session_loop_termination: session_loop_termination_from_handle(session_loop_handle),
};

随后 Core spawn 一个 session_loop task,持有 Arc<Session>、configured config 和 rx_sub。 SessionIo 把 tx_sub、rx_event、agent status receiver 和一个共享 termination future 交给 CodexThread。

4.1 submission ​

submission queue 的发送使用 send().await。如果宿主在 Session 来不及处理时连续提交超过 512 个 操作,生产者会等待,防止控制消息无限占用内存。Session loop 仍逐条 dispatch;它不会为每个 Op 无条件 spawn 新 task。

event queue 则是无界的。Core 发 event 时不会因为 UI 或 App Server 暂时不消费而阻塞 Turn,避免 持久化/工具 task 与客户端形成直接死锁;代价是慢消费者会让内存增长。send_event_raw 先按 policy 写 rollout,再向 channel 发送,所以关闭 event receiver 只会记录 debug,不会反向取消 Session。

submission sender 的失败被明确转换为“内部 agent 已死亡”,而不是继续接受无法消费的 Op:

源码位置:codex-rs/core/src/session/mod.rs :: SessionIo::submit_with_trace, submit_with_id

rust
pub(crate) async fn submit_with_trace(
    &self,
    op: Op,
    trace: Option<W3cTraceContext>,
    parent_turn_id: Option<String>,
) -> CodexResult<String> {
    let id = new_submission_id();
    let sub = Submission {
        id: id.clone(),
        op,
        client_user_message_id: None,
        trace,
        parent_turn_id,
    };
    // 只有成功进入有界队列后才把 ID 返回调用方。
    self.submit_with_id(sub).await?;
    Ok(id)
}

pub(crate) async fn submit_with_id(&self, mut sub: Submission) -> CodexResult<()> {
    // 调用方未显式传 trace 时,在跨异步队列前捕获当前 span 上下文。
    if sub.trace.is_none() {
        sub.trace = current_span_w3c_trace_context();
    }
    self.tx_sub
        .send(sub)
        .await
        // receiver 已关闭表示 session loop 消失,不能伪装成普通队列拥塞。
        .map_err(|_| CodexErr::InternalAgentDied)?;
    Ok(())
}

ID、trace 和 parent Turn 关系都在发送前固化。队列满只会让 send().await 等待;只有 receiver drop 才 进入 InternalAgentDied 错误分支。

无界不等于可靠持久化。进程崩溃仍会丢内存中的 event;delta 是否写 rollout又由 persistence policy 决定。事件背压和持久化是两条独立轴。

4.2 Session Loop ​

活动 Turn 在另一个 task 中运行时,submission loop 仍要处理:

  • Interrupt;
  • exec/patch approval response;
  • request_user_input、permission、dynamic tool 和 MCP elicitation response;
  • steer/inter-agent input;
  • realtime audio/text/close;
  • config/MCP refresh;
  • shutdown。

如果 Turn 在 loop 内被完整 await,所有这些控制消息都会饿死。loop 只负责短 dispatch 或启动 task, 耗时模型/工具执行由 active SessionTask 承担。

4.3 关闭入口汇合 ​

显式 Op::Shutdown 令 dispatch 返回 true,loop 退出前执行 runtime shutdown、extension stop、thread persistence shutdown,并发送 ShutdownComplete。如果所有 submission sender 被 drop,recv() 返回 错误,loop 仍执行 teardown,只是没有对应 shutdown submission。

SessionIo::shutdown_and_wait() 先提交 Shutdown,再等待由 session loop JoinHandle 转成的 Shared<BoxFuture<()>>。Shared future 允许多个持有者等待同一个 loop;它不是另一个后台 task,也不 传递业务消息。

5. 状态通知 ​

5.1 AgentStatus ​

agent status channel 初始为 PendingInit。deliver_event_raw() 从部分 EventMsg 推导新状态并 send_replace()。watch receiver 读取的是最近状态,不保证看到每次中间变化;适合回答“现在是什么 状态”,不适合作为完整事件日志。

5.2 InputQueue唤醒 ​

InputQueue 的 watch value 只有 Mailbox 与 Steer。真实输入仍保存在 mailbox deque 或 TurnState.pending_input,watch 只唤醒 sampling wait。订阅时还会检查已有队列,避免先入队后订阅而 错过状态。

5.3 Thread-created ​

ThreadManager 的 broadcast 容量为 1024,允许 App Server 和其他观察者各自订阅。慢 receiver 落后超过 buffer 会收到 Lagged;App Server 对此不会假装补齐逐条通知,而是继续监听并依赖 thread list/live state 做权威查询。

unified exec output broadcast 容量为 64。慢输出消费者也可能 lag,watcher 遇到 lag 会继续读取后续 chunk。最终 transcript 另存于 HeadTailBuffer,因此实时 delta 丢失和最终聚合结果是两种不同保证。

6. Turn task ​

Session::start_task() 创建:

  • 一个 root CancellationToken,传给 task 的是 child token;
  • Arc<Notify> 作为 graceful completion 信号;
  • 带 thread/turn/model 字段的 OTel span;
  • tokio::spawn 返回的 handle,再包装成 AbortOnDropHandle;
  • RunningTask,安装到 ActiveTurn.task。

承载 task 的 Tokio task 执行 SessionTask::run(),结束后 flush rollout,未取消时调用 on_task_finished(),最后 done.notify_waiters()。

安装任务前,Core 先终止旧任务,再把新任务的取消、完成和 Turn 状态绑定到同一个 active slot:

源码位置:codex-rs/core/src/tasks/mod.rs :: Session::spawn_task, Session::start_task

rust
pub async fn spawn_task<T: SessionTask>(
    self: &Arc<Self>,
    turn_context: Arc<TurnContext>,
    input: Vec<TurnInput>,
    task: T,
) {
    // 先结束旧 task,保证 ActiveTurn 的 task slot 不会同时出现两个 owner。
    self.abort_all_tasks(TurnAbortReason::Replaced).await;
    self.clear_connector_selection().await;
    self.start_task(
        turn_context,
        input,
        task,
        MailboxParentProvenance::Ignore,
    )
    .await;
}

pub(crate) async fn start_task<T: SessionTask>(
    self: &Arc<Self>,
    turn_context: Arc<TurnContext>,
    input: Vec<TurnInput>,
    task: T,
    mailbox_parent_provenance: MailboxParentProvenance,
) {
    let task: Arc<dyn AnySessionTask> = Arc::new(task);
    let task_kind = task.kind();
    let started_at = Instant::now();
    // 新 task 获得独立取消根和完成通知,分别服务取消与 graceful wait。
    let cancellation_token = CancellationToken::new();
    let done = Arc::new(Notify::new());

    let (pending_items, parent_turn_id) =
        self.input_queue.get_pending_input(&self.active_turn).await;
    if let (MailboxParentProvenance::Attribute, Some(id)) =
        (mailbox_parent_provenance, parent_turn_id)
    {
        turn_context.turn_metadata_state.set_parent_turn_id(id);
    }

    let turn_state = {
        // 只在安装/取得 active slot 时持锁,后续输入归并不占用该锁。
        let mut active = self.active_turn.lock().await;
        let turn = active.get_or_insert_with(ActiveTurn::default);
        debug_assert!(turn.task.is_none());
        Arc::clone(&turn.turn_state)
    };
    self.input_queue
        .extend_pending_input_for_turn_state(turn_state.as_ref(), pending_items)
        .await;
    // ...
}

这里的顺序防止两个竞态:旧任务未结束就覆盖 handle,以及 queued input 在新 TurnState 安装前丢失。 debug_assert! 只是开发期检查,真正的互斥来自先执行 abort_all_tasks() 和持有 active_turn 锁。

6.1 取消路径 ​

handle_task_abort() 的顺序是:

  1. cancel task token;
  2. 等待 done.notified() 或 graceful timeout;
  3. abort Tokio handle;
  4. 调用具体 SessionTask::abort() 做补偿;
  5. 记录 interrupted marker、flush,并发送 TurnAborted;
  6. 再 flush terminal event。

CancellationToken 只通知协作式退出;AbortOnDropHandle 才是强制兜底;Notify 用于判断 task 是否已 自行完成。三者缺一会造成不同的泄漏或错误终态。

取消不是一次 abort() 调用,而是一条先给任务自清理机会、再强制终止、最后持久化终态的时序。

事件只有在 interrupted marker 已进入 rollout 后才发送,尽量让实时客户端和下一次 resume 得到一致 终态;graceful timeout 则防止不响应 token 的 task 永久占据 active slot。

6.2 Task内调用 ​

run_turn() 和 run_sampling_request() 是同一 RegularTask 内的 async 调用,不是每轮都 spawn 新 Turn task。工具续轮、retry、pending input 和 compaction 可以让同一 task 多次进入模型采样,同时 保持一个 TurnContext 和 cancellation subtree。

7. 模型流 ​

provider 返回的 API stream 不直接暴露给 Turn。map_response_events() spawn 一个 mapper task,将 API error 转成 CodexErr、记录 telemetry/trace、累计 added items,再向容量 1600 的 Tokio mpsc 发送 ResponseEvent。

该边界还有两个辅助原语:

  • oneshot<LastResponse>:Completed 时返回 response ID 与 added items,用于 WebSocket sticky state;
  • consumer_dropped: CancellationToken:ResponseStream::drop() 时取消 mapper,避免 Turn 已停止而 producer 继续读取 provider stream。

1600 项缓冲允许吸收 SSE/WebSocket burst;满时 mapper 的 send().await 会把背压传回 provider stream 读取。它不是丢弃型队列。若 receiver 已关闭,mapper 记录 cancelled/failed trace 后退出。

8. 工具调用 ​

一次模型 response 可包含多个 tool call。try_run_sampling_request() 用 FuturesOrdered<BoxFuture<...>> 保存 in-flight output,ToolCallRuntime 对每个 call spawn dispatch task,并用 AbortOnDropHandle 管理。

并发准入不是 channel,而是 parallel_execution 读写锁:

  • 支持 parallel 的 tool 获取 read guard,可以并发;
  • 不支持 parallel 的 tool 获取 write guard,与所有其他调用互斥;
  • FuturesOrdered 保持结果按加入顺序交回模型,即使内部完成顺序不同。

每个 dispatch task 持有 step 对应的 ToolRouter、TurnContext 和 invocation cancellation token。取消 时,普通 runtime 可直接 abort;声明需要 teardown 的 runtime 会等待其完成清理,再生成 aborted tool response。于是“工具 task 结束”和“模型收到 terminal tool output”仍是两个阶段。

9. Rollout writer ​

RolloutRecorder 为每个 live rollout spawn 一个 writer task,独占异步文件 handle。调用方通过容量 256 的 Tokio mpsc 发送四种命令:

源码位置:codex-rs/rollout/src/recorder.rs :: RolloutCmd

rust
enum RolloutCmd {
    // 批量 item 不逐次确认;需要 durability/关闭语义的命令才携带 oneshot ack。
    AddItems(Vec<RolloutItem>),
    Persist { ack: oneshot::Sender<std::io::Result<()>> },
    Flush { ack: oneshot::Sender<std::io::Result<()>> },
    Shutdown { ack: oneshot::Sender<std::io::Result<()>> },
}

AddItems 不需要逐次 ack;Persist/Flush/Shutdown 用 oneshot 把文件操作结果返回调用方。channel 满时 send future yield,让写入压力回传到 Session,但不会在 Tokio worker 上执行 blocking file I/O。

普通 I/O 失败不会立刻杀掉 writer:未写 items 留在 pending buffer,后续 flush/shutdown 可以重试。 只有 writer task 的 terminal failure 会保存到 RolloutWriterTask.terminal_failure,供后续 API 返回。 Shutdown 成功才 break;若 drain 失败,writer 保持存活,允许调用者再次尝试。

10. MCP 与启动预热 ​

10.1 MCP预热队列 ​

Session 创建容量 1 的 async_channel<()>。schedule_mcp_prewarm() 使用 try_send(()),所以 worker 忙时 的多次请求只保留“仍有一次刷新需要处理”这个事实。worker 同时 select:

  • shutdown token;
  • prewarm request;
  • AuthManager revision watch;
  • refresh_mcp_if_dirty()。

它持有 Weak<Session>,不会仅因后台 worker 延长 Session 生命周期。停止时先 cancel token,再 take 并 await JoinHandle。预热结果允许过时,因为真正 sampling 会在 capture_step_context() 中重新获得 正确 MCP binding;prewarm 只优化延迟。

10.2 模型预热 ​

启用 Responses WebSocket 时,Session spawn startup prewarm task,结果是 ModelClientSession。SessionStartupPrewarmHandle 用 AbortOnDropHandle 持有它;首个 RegularTask 在 timeout/cancellation 约束下消费结果。未启用 WebSocket 时只异步预热 auth。

prewarm 不是主 Turn。它使用 INITIAL_SUBMIT_ID 构造临时 Turn/Step/Prompt,不写用户 Turn lifecycle;成功后把已热连接交给第一个 RegularTask,失败或超时则正常创建新的 client session。

11. Unified exec ​

每个 UnifiedExecProcess 持有:

  • output broadcast 64:给实时 output watcher;
  • ProcessState watch:给 write/wait/read 查询最新 starting/exited/failed 状态;
  • output task:从 local PTY 或 exec-server process 持续读 chunk;
  • cancellation token:表示进程/输出生命周期结束;
  • Notify:output closed 与 output drained;
  • interaction mutex:序列化 terminal interaction 和最终事件。

Core 再 spawn output watcher 将 UTF-8 chunk 聚合到 transcript 并发送 ExecCommandOutputDelta;exit watcher 等进程结束、output drained 和可选 network denial monitor 后,只发一次 ExecCommandEnd。 这个次序避免“终态先到、最后输出后到”。

普通非 PTY exec 也会分别 spawn stdout/stderr reader task,与 child wait/timeout/cancellation 并发, 最后 join reader,保证管道被排空或按 I/O drain timeout 收尾。

12. Realtime队列 ​

Realtime conversation 启动时创建五条 bounded queue:

队列容量发送语义设计目的
audio input256try_send,满时丢 frame音频实时性高于完整性
text input64send().await文本不可静默丢失
handoff output64send().await后台 agent 输出回灌 realtime
output events256producer await向 client fanout RealtimeEvent
transcript tail1单次 close 结果session end 尾部转交

一个 input task 拥有 WebSocket/WebRTC sideband writer 和输入 receivers;另一个可注册 fanout task 将 events 投影到 Core event。ConversationState 保存两个 JoinHandle 与 stop token。

重新 start 会先取出并停止旧 state;shutdown cancel stop token、关闭 transport,并根据路径 await 或 detach fanout。audio queue 满只产生日志并返回成功,防止网络抖动把麦克风 producer 阻塞到不可恢复; 这与文本队列的可靠语义有意不同。

13. Delegate协作 ​

review/guardian 和 legacy delegate 可能在 parent Session 旁启动 child Session。interactive delegate 额外建立容量 512 的 caller-op 与 caller-event bridge,并 spawn forward_ops、forward_events 两个 task;parent cancellation 的 child token 同时终止两条桥。

approval 不直接转发给 caller,而由 event bridge 路由到 parent Session 决策,再把 response 提交给 child。one-shot delegate 另建容量 512 的 event bridge,在 TurnComplete/TurnAborted 后主动发送 child Shutdown 并关闭对外 submission sender。

这些桥解释了为什么日志里可能同时出现 parent task、child session loop 和两个 forwarding task。它们 不是同一个 Turn 的普通工具并行,而是嵌套 runtime。

14. 关闭顺序由所有权 ​

显式 Session shutdown 的主要顺序是:

该顺序先阻止新工作和外部 I/O,再关闭扩展,最后持久化并发 terminal event。关闭 channel 并不能 自动完成这些步骤:drop submission senders 只能让 loop 退出,teardown 仍必须显式 cancel/join 各个 owner。反过来,只 cancel child token 而不等待 handle,也可能让 cleanup 尚未完成就释放 store。

15. 异步边界定位 ​

症状首先检查关键观察点
submit 长时间等待submission queue 是否满、loop 是否卡在长 dispatchbounded 512、session loop span
内存随 UI 卡顿增长event consumer 是否停止unbounded Event channel
状态跳过中间值是否把 watch 当 event logagent/process/input watch
新 thread 通知缺失broadcast receiver 是否 Laggedthread-created 1024
模型请求已取消仍读流ResponseStream 是否 drop、mapper token 是否 cancelmapper task trace
工具完成顺序异常parallel gate 与 FuturesOrderedcall ID、execution start time
TurnAborted 后仍有子进程task token、runtime teardown、process Job/PGIDhandle_task_abort + exec manager
ExecCommandEnd 早于最后输出output_closed/output_drained orderingoutput/exit watcher
rollout flush 一直失败writer terminal failure 还是 pending retryRolloutCmd ack、pending_items
MCP 重复刷新或状态旧prewarm coalescing 与 step capturebounded 1、dirty revision
realtime 音频缺帧audio queue 满是否触发 drop warningcapacity 256 + try_send

验证异步修改时,应使用能够控制时序的测试:paused Tokio time、oneshot gate、Notify、受控 mock stream 和明确 timeout。用固定 sleep 猜测 task 已运行只会制造新的竞态。Core 已有的 channel capacity、 consumer-drop、shutdown、rollout retry、MCP refresh 和 unified exec tests 都采用“先阻塞一个边界,再 观察另一个边界是否仍可推进”的方式验证拓扑。

按 Op 分派继续阅读 Session运行时处理;按资源逆序关闭继续阅读 Session关闭流程。

几条容量和关闭结论可以直接由测试反向验证:response mapper 的回压测试填满容量为 1600 的队列, 再让 consumer drop,断言已产生的 output item 仍进入 Cancelled rollout;submission channel close 测试安排活动 Turn 和 thread lifecycle recorder,断言 turn_abort 先于 thread_stop;rollout retry 测试则让一次写入失败后再次触发 flush,断言 pending items 没有被静默丢弃。它们分别覆盖满队列、接收者 终止和持久化重试;没有覆盖每一条 Realtime 队列在真实网络断开下的所有交错。

可以用下面的只读搜索把本文的 task topology 主线落回源码:

bash
rg -n "submission_loop|mpsc::|broadcast::|CancellationToken|shutdown" codex-rs/core/src