雨天小六

读懂 Codex(3.3):CodexThread 的 Channel、状态 Watcher 与关闭协议

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

#Codex#Agent Runtime#生命周期#软件架构

CodexThread 是调用方接触 Session 的受控门面。它没有把 Session 的锁和内部方法全部暴露出去,而是组合四条不同语义的通道:提交、事件、最新状态和终止通知。把它们统一成一个队列会同时破坏背压、广播和关闭等待。

四条路径各自解决什么问题

路径实现容量/语义适合的数据
Submissionbounded async channel容量 512,发送者受背压UserInput、Interrupt、审批、Shutdown
Eventunbounded async channel单一接收流,不阻塞 Session 产出增量、工具事件、终态
AgentStatusTokio watch只保存最新值,可多订阅Pending、Running、Completed 等状态
Terminationshared future多个等待者共享同一完成结果Session loop 已完全退出
CodexThread 四条输入输出和状态通道
图 3.3-1:命令流需要背压,事件流需要避免阻塞 Runtime,状态观察只关心最新值,终止等待必须允许多个调用者;它们的协议目标不同。

Submission 有界意味着恶意或失控客户端不能无限堆积控制操作。Event 无界则是另一种取舍:模型流不能因为 UI 暂时没消费而握着 Session 内部锁停住;代价是调用者必须持续排空事件,外层也需要防止长期不读造成内存增长。

Watch 不重放每次变化。新订阅者立刻得到当前 AgentStatus,并在值变化时醒来,适合列表页或 Agent 控制器,不适合构建完整审计日志。完整历史来自 Event/Rollout。

class SessionIo:
    submissions: BoundedSender[Submission]  # 512
    events: UnboundedReceiver[Event]
    status: WatchReceiver[AgentStatus]
    terminated: SharedFuture[None]

class CodexThread:
    session: Arc[Session]
    io: SessionIo

    async def submit(self, op):
        await self.io.submissions.send(Submission(uuid_v7(), op))

    async def next_event(self):
        return await self.io.events.recv()

状态 Watcher 是事件投影,不是第二套状态机

AgentStatus 初值是 PendingInit。Session 发出生命周期 Event 时,同时把其中可映射的事件投影为新状态并 send_replace。因此状态和事件共享事实来源:状态订阅者看到最新摘要,事件消费者看到带顺序和负载的原始变化。

如果 Watch 接收端落后,它不会逐项补发中间状态。这正是预期行为。例如一个 UI 只需知道 Agent 现在已经 Idle,无需先消费 Running、Waiting、Completed 的每个瞬间;需要这些细节的观察者应读 Event。

Shutdown 不是关掉 sender 就结束

shutdown_and_wait 先克隆共享终止 Future,再尝试提交 Op::Shutdown,最后等待 submission loop 完成完整 teardown。先取得 Future 很重要:如果 Shutdown 触发 loop 很快退出并关闭内部对象,等待者仍然持有可解析的完成信号。

async def shutdown_and_wait(io):
    terminated = io.termination.clone()
    try:
        await io.submit(Shutdown)
    except InternalAgentDied:
        pass                 # channel 已关闭,teardown 可能正在进行
    await terminated

多个调用者可以并发等待同一个终止结果。关闭已经开始时,第二个 Shutdown 发送失败不等于第二个等待者可以提前返回;它仍须等共享 Future 完成。

CodexThread shutdown_and_wait 的多等待者关闭协议
图 3.3-2:Shutdown 是请求,termination 才是完成证明。即使提交端已经关闭,所有等待者仍等到活动 Turn、子进程和 ThreadStop 收尾结束。

Submission 通道自然关闭也必须 teardown

最后一个 sender 被丢弃后,submission_loop 的 receive 返回关闭。源码没有把它当成“无需清理的异常退出”,而是走与显式 Shutdown 相同的 runtime teardown:先中止活动 Turn,再停止线程级资源,发 ThreadStop 生命周期并关闭 LiveThread。

这个保证覆盖了调用方崩溃、门面被整体释放等路径。测试还专门验证 Channel close 会先 abort active turn,再发 thread stop,避免后台模型流或工具进程成为孤儿。

设计权衡

Event 无界和 Submission 有界看似不对称,实际上分别保护生产者活性和输入内存上限。Shared termination future 则把“发出命令”与“确认完成”分开。若 API 只提供 shutdown() 而没有 wait,ThreadManager 就无法可靠判断何时能从 live registry 删除对象。

评论


← 返回文章列表