CodexThread 是调用方接触 Session 的受控门面。它没有把 Session 的锁和内部方法全部暴露出去,而是组合四条不同语义的通道:提交、事件、最新状态和终止通知。把它们统一成一个队列会同时破坏背压、广播和关闭等待。
四条路径各自解决什么问题
| 路径 | 实现 | 容量/语义 | 适合的数据 |
|---|---|---|---|
| Submission | bounded async channel | 容量 512,发送者受背压 | UserInput、Interrupt、审批、Shutdown |
| Event | unbounded async channel | 单一接收流,不阻塞 Session 产出 | 增量、工具事件、终态 |
| AgentStatus | Tokio watch | 只保存最新值,可多订阅 | Pending、Running、Completed 等状态 |
| Termination | shared future | 多个等待者共享同一完成结果 | Session loop 已完全退出 |
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 完成。
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 删除对象。
评论
登录后即可评论