雨天小六

读懂 Codex(2.5):CodexThread 的提交端、事件端与状态查询端

· 更新于 2026-08-02 · 专栏:读懂 Codex

#Codex#Agent Runtime#软件架构#CodexThread#Async Runtime

ThreadManager 发布的不是一个裸 Session,而是 Arc<CodexThread>。这个类型看起来像一层很薄的 包装,实际却定义了 Core 对上层最重要的并发合同:哪些动作必须排队、哪些信息按事件流返回、哪些 状态只能读取最新快照,以及调用方怎样确认后台循环已经真正停止。

如果把它简单理解成“给 Agent 发消息的对象”,就会混淆提交成功、Turn 启动、Turn 完成和 Thread 关闭四个完全不同的时刻。

CodexThread 是门面,不是第二套 Runtime

CodexThread 的核心成员可以分为三组:

类别成员用途
Runtime 所有权Arc<Session>活动 Turn、历史、服务、工具状态的真实所有者
外部通信SessionIoSubmission、Event、AgentStatus 与终止等待
发布快照Session source、首个 SessionConfigured、rollout path启动后稳定公开的识别与配置信息

因此,CodexThread 不会复制 Session 的状态机。它做的是缩小可见面:上层不必拿到 Session 的所有 内部锁和服务,也能完成提交、读取、观察与关闭。

SessionIo 不是一条双向队列

四种并发原语承担不同语义:

  1. Submission 使用容量 512 的有界 Channel。发送端在队列满时等待,避免调用者无限制造待处理操作;
  2. Event 使用无界 Channel。运行时产生的流式输出、审批请求和生命周期事件不会因前端短暂变慢而 直接堵住模型处理;代价是消费者必须持续排空;
  3. AgentStatus 使用 watch。新订阅者只需要当前最新状态,不需要重放每一次 Running/Completed 变化;
  4. termination 是共享 Future。多个所有者可以等待同一个 submission loop 结束。
CodexThread 包围 Session,并通过有界 Submission 通道、无界 Event 通道、状态 watch 和共享终止 future 连接调用方
图 2.5-1:四条通路的背压与保留语义不同,不能合并成一个“消息总线”。

这里有一个刻意的不对称:Submission 需要背压,Event 更重视 Runtime 不被 UI 阻塞。无界 Event 并不意味着无限内存安全;它把“及时消费事件”变成客户端责任。

submit 返回的是关联 ID,不是完成结果

普通提交会生成 UUIDv7 形式的 ID,将 Op 包进 Submission。Submission 还可以携带 trace、父 Turn ID 和客户端用户消息 ID。若调用方没有显式 trace,发送端会从当前 tracing span 捕获 W3C 上下文。

async def submit(io: SessionIo, op: Op) -> str:
    submission_id = uuid7()
    envelope = Submission(
        id=submission_id,
        op=op,
        trace=current_w3c_trace_context(),
        parent_turn_id=None,
        client_user_message_id=None,
    )
    await io.submission_tx.send(envelope)  # 只保证入队
    return submission_id

返回 ID 后,submission loop 可能尚未收到它,Turn 也可能因参数、容量或活动任务限制而失败。完成语义 必须从同 ID 相关的 Event 或上层映射出的 Turn 状态判断。

submit_with_id 允许调用方自己构造 Submission,但源码明确要求谨慎使用。统一生成 ID 能维持时间有序 特性和关联规则,避免不同 Surface 各自发明 ID。

UserInput 有一条带容量检查的专用入口

submit_user_input_with_client_user_message_id 在入队前调用 AgentControl 的 execution-capacity 检查, 然后用专用字段保留客户端消息身份。这使 App Server 能把一次客户端 turn/start 与 Core 事件相关联, 同时在多 Agent 容量已经耗尽时尽早拒绝。

该入口用 debug assertion 约束 Op 必须是 UserInput。也就是说,这个字段不是任意 Op 的通用幂等键。

steer、inject 与普通 submit 的差别

CodexThread 还暴露三条不走普通 Submission 队列的活动 Turn 路径:

入口条件成功效果失败时输入
steer_input存在可 steer 的活动 Turn,expected Turn ID 相容把用户输入送进当前 Turn返回结构化错误
inject_if_running存在活动 Turn注入模型可见 ResponseItem原样返还 Vec,调用者可重试
try_start_turn_if_idle无排队触发、无活动任务、非 Plan mode扩展启动自动 Regular Turn返回稳定拒绝原因和原输入

Review 与 Compact 属于不可 steer Turn。自动 idle 工作还必须给用户/客户端排队的工作让路,避免扩展在 竞态中抢先开一个模型 Turn。

async def extension_idle_work(thread: CodexThread, items: list[ResponseItem]):
    result = await thread.try_start_turn_if_idle(items)
    if result.rejected:
        # reason ∈ pending_trigger_turn / plan_mode / busy
        return RetryDecision(result.reason, result.original_items)
    return Started()

这些直达 Session 的入口是有意设计的控制面,不应被理解为绕过所有规则:它们内部仍检查活动任务、 Turn 类型和期望 ID。

事件读取与状态订阅解决不同问题

next_event() 从 Event Channel 取走一个事件。多个消费者若同时读取同一 Receiver,是竞争消费,并不会 让每个消费者都得到副本;正常情况下应由一个 Surface 读取后再自行广播或投影。

AgentStatus 则是观察最新值。调用方可以:

  • 一次读取 agent_status()
  • clone/subscribe watch Receiver,等待后续变化;
  • 在错过中间状态时仍取得最终值。

状态由事件派生但不是事件的替代品。Completed 状态可能带最终消息,却不包含工具输出 delta、审批请求 或 token 更新等完整序列。

一次提交从入队到 Turn 事件、状态更新以及 Shutdown 等待 loop termination 的时序
图 2.5-2:Submission ID 只建立关联;TurnComplete 与 loop termination 分别结束一次 Turn 和整个 Session。

配置快照与动态查询不能混为一谈

ThreadManager 在发布时缓存首个 SessionConfiguredEvent,其中记录模型、Provider、审批策略、权限、 cwd、初始消息、network proxy 和 rollout path。它适合表达“这个 Session 启动时是什么”。

运行期间设置可以更新,因此 CodexThread 另有 config_snapshot()config()、instruction sources、 当前 MCP config/runtime context 等动态查询。调用方若始终读取启动事件,就可能展示过期模型或权限。

持久化操作受 ephemeral 边界约束

通过活 Thread 可以读取/更新 ThreadStore metadata、加载历史、追加 rollout item,也可以 materialize 或 flush。Session 若以 ephemeral 模式启动,没有 LiveThread;要求持久化的操作会返回明确的“persistence is disabled”错误,而不是偷偷建立一个文件。

inject_response_items 是一个特殊管理入口:

  • 空数组直接拒绝;
  • 没有当前 TurnContext 时构造默认 Turn/StepContext;
  • 把 ResponseItem 写入模型历史;
  • 不创建普通 User Turn;
  • 最后 flush rollout,给调用者一个持久化屏障。

这适合恢复或系统注入,不适合模拟用户发消息。

shutdown_and_wait 为什么要做两步

把 Shutdown 入队只表示关闭请求被接收。真正的 Session teardown、持久化关闭与 lifecycle hook 都发生在 submission loop 中或 loop 退出收尾里。因此关闭 API 先保存共享 termination future,再提交 Shutdown, 最后等待 future。

async def shutdown_and_wait(io: SessionIo) -> None:
    terminated = io.session_loop_termination.clone()
    try:
        await io.submit(Op.Shutdown)
    except InternalAgentDied:
        pass  # loop 已结束,也符合关闭目标
    await terminated

如果所有 Submission Sender 被丢弃导致 Channel 关闭,loop 也会执行 teardown。is_running 只以发送端 是否关闭为近似判断,它不是“当前是否有活动 Turn”的同义词。

失败矩阵

场景可观察结果调用方正确处理
Submission 队列满send 等待容量不要在 UI 主循环持锁等待
submission loop 已死InternalAgentDied停止继续提交,检查终止状态
submit 返回 ID 后业务拒绝后续 Error Event不能把 ID 当成功结果
多消费者读 Event单个事件只被其中一个取走建立单读者 fan-out
status 订阅较晚只获得最新状态需要审计时读取 Event/Rollout
ephemeral Thread 调持久化 API显式错误不假设存在 rollout
Shutdown 提交时 loop 已结束关闭路径容忍 InternalAgentDied仍等待共享 termination

可验证的设计不变式

async def test_submit_is_not_completion(thread):
    sid = await thread.submit(user_input("inspect"))
    assert sid
    event = await next_correlated_terminal_event(thread, sid)
    assert event.type in {"turn_complete", "error", "turn_aborted"}


async def test_watch_keeps_latest_value(thread):
    status_rx = thread.subscribe_status()
    await run_one_turn(thread)
    assert status_rx.borrow().kind in {"completed", "errored", "interrupted"}


async def test_shutdown_waits_for_loop(thread):
    await thread.shutdown_and_wait()
    assert not thread.is_running()

CodexThread 的价值不是“少写几个函数调用”,而是把排队、流式输出、状态观察和生命终止分成四个可 推理的合同。后续分析 Op 与 Event 时,都应先确认它们具体走的是哪条通路。

评论


← 返回文章列表