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::bounded | submission、MCP refresh、delegate bridge、Realtime input | 是,直到容量 | send().await 背压;try_send 可合并/丢弃 |
async_channel::unbounded | Core → host Event、部分 delegate event | 是 | 不阻塞生产者,压力进入内存 |
tokio::sync::mpsc | model stream、rollout writer | 是,直到容量 | sender await,形成背压 |
watch | agent/process/input activity 最新状态 | 否,只保留最新值 | 中间状态可被覆盖 |
broadcast | thread-created、多订阅者 exec output | 环形缓冲 | 慢 receiver 得到 Lagged |
oneshot | approval、flush、last response、一次性完成 | 单值 | sender/receiver drop 表示取消 |
Notify | task done、output drained、数据到达 | 不保存业务消息 | permit/唤醒语义,不是队列 |
CancellationToken | task、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 created | broadcast | 1024 | ThreadId | ThreadManager::new | App Server/listener;manager drop 关闭 |
| Session submission | async_channel::bounded | 512 | Submission | Session::spawn_internal | submission loop;Shutdown 或所有 sender drop |
| Session event | async_channel::unbounded | 无界 | Event | Session::spawn_internal | CodexThread::next_event;Session sender drop |
| Agent status | watch | 最新 1 个 | AgentStatus | Session::spawn_internal | status subscriber;Session sender drop |
| Input activity | watch | 最新 1 个 | Mailbox/Steer | InputQueue::new | sampling wait;InputQueue drop |
| MCP prewarm | async_channel::bounded | 1 | () | Session::new | prewarm worker;cancel token + join |
| Auth changes | watch | 最新 1 个 | revision u64 | AuthManager | MCP prewarm worker;sender/worker drop |
| Model events | Tokio mpsc | 1600 | Result<ResponseEvent> | map_response_events | Turn sampling;stream drop cancels mapper |
| Last response | oneshot | 1 | LastResponse | model mapper | ModelClientSession sticky state |
| Rollout commands | Tokio mpsc | 256 | RolloutCmd | RolloutRecorder::new | rollout writer;Shutdown command + ack |
| Unified exec output | broadcast | 64 | Vec<u8> | UnifiedExecProcess::new | output consumers/watchers;process drop |
| Unified exec state | watch | 最新 1 个 | ProcessState | UnifiedExecProcess::new | wait/read/write paths;process drop |
| Realtime audio | async_channel::bounded | 256 | RealtimeAudioFrame | conversation start | input task;full 时丢 frame |
| Realtime text | async_channel::bounded | 64 | ConversationTextParams | conversation start | input task;sender await |
| Realtime handoff | async_channel::bounded | 64 | RealtimeOutbound | conversation start | input task;sender await |
| Realtime events | async_channel::bounded | 256 | RealtimeEvent | conversation start | fanout/client;sender await |
| Transcript tail | async_channel::bounded | 1 | String | conversation start | close/final consumer |
| Turn interaction | oneshot | 1/请求 | 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
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
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
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() 的顺序是:
- cancel task token;
- 等待
done.notified()或 graceful timeout; - abort Tokio handle;
- 调用具体
SessionTask::abort()做补偿; - 记录 interrupted marker、flush,并发送
TurnAborted; - 再 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
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;
ProcessStatewatch:给 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 input | 256 | try_send,满时丢 frame | 音频实时性高于完整性 |
| text input | 64 | send().await | 文本不可静默丢失 |
| handoff output | 64 | send().await | 后台 agent 输出回灌 realtime |
| output events | 256 | producer await | 向 client fanout RealtimeEvent |
| transcript tail | 1 | 单次 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 是否卡在长 dispatch | bounded 512、session loop span |
| 内存随 UI 卡顿增长 | event consumer 是否停止 | unbounded Event channel |
| 状态跳过中间值 | 是否把 watch 当 event log | agent/process/input watch |
| 新 thread 通知缺失 | broadcast receiver 是否 Lagged | thread-created 1024 |
| 模型请求已取消仍读流 | ResponseStream 是否 drop、mapper token 是否 cancel | mapper task trace |
| 工具完成顺序异常 | parallel gate 与 FuturesOrdered | call ID、execution start time |
| TurnAborted 后仍有子进程 | task token、runtime teardown、process Job/PGID | handle_task_abort + exec manager |
| ExecCommandEnd 早于最后输出 | output_closed/output_drained ordering | output/exit watcher |
| rollout flush 一直失败 | writer terminal failure 还是 pending retry | RolloutCmd ack、pending_items |
| MCP 重复刷新或状态旧 | prewarm coalescing 与 step capture | bounded 1、dirty revision |
| realtime 音频缺帧 | audio queue 满是否触发 drop warning | capacity 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 主线落回源码:
rg -n "submission_loop|mpsc::|broadcast::|CancellationToken|shutdown" codex-rs/core/src