ThreadManager 发布的不是一个裸 Session,而是 Arc<CodexThread>。这个类型看起来像一层很薄的
包装,实际却定义了 Core 对上层最重要的并发合同:哪些动作必须排队、哪些信息按事件流返回、哪些
状态只能读取最新快照,以及调用方怎样确认后台循环已经真正停止。
如果把它简单理解成“给 Agent 发消息的对象”,就会混淆提交成功、Turn 启动、Turn 完成和 Thread 关闭四个完全不同的时刻。
CodexThread 是门面,不是第二套 Runtime
CodexThread 的核心成员可以分为三组:
| 类别 | 成员 | 用途 |
|---|---|---|
| Runtime 所有权 | Arc<Session> | 活动 Turn、历史、服务、工具状态的真实所有者 |
| 外部通信 | SessionIo | Submission、Event、AgentStatus 与终止等待 |
| 发布快照 | Session source、首个 SessionConfigured、rollout path | 启动后稳定公开的识别与配置信息 |
因此,CodexThread 不会复制 Session 的状态机。它做的是缩小可见面:上层不必拿到 Session 的所有 内部锁和服务,也能完成提交、读取、观察与关闭。
SessionIo 不是一条双向队列
四种并发原语承担不同语义:
- Submission 使用容量 512 的有界 Channel。发送端在队列满时等待,避免调用者无限制造待处理操作;
- Event 使用无界 Channel。运行时产生的流式输出、审批请求和生命周期事件不会因前端短暂变慢而 直接堵住模型处理;代价是消费者必须持续排空;
- AgentStatus 使用
watch。新订阅者只需要当前最新状态,不需要重放每一次 Running/Completed 变化; - termination 是共享 Future。多个所有者可以等待同一个 submission loop 结束。
这里有一个刻意的不对称: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 更新等完整序列。
配置快照与动态查询不能混为一谈
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 时,都应先确认它们具体走的是哪条通路。
评论
登录后即可评论