具体问题与边界
为什么队列容量、消费者数量和关闭信号是 Runtime 协议,而不是随手选择的 asyncio 参数?
输入/输出各一个 Queue,默认容量 32;Session 是输入的唯一消费者,客户端是事件流的单消费者;关闭哨兵只存在于事件队列。 本节不是把 Rust 改写成 Python;它从锁定提交的字段、调用顺序和测试行为提炼实现合同,再检查 Mini Codex 是否以 Python 的并发原语保持同一条不变量。
协议、类型与状态所有权
| 对象 | 创建/所有者 | 生命周期与作用域 | 是否持久化 |
|---|---|---|---|
_submissions | Session | 容量 32;只由主循环 get | 否 |
_events | Session | 容量 32;由事件迭代器 get | 否 |
_EVENT_STREAM_CLOSED | Session._run Shutdown 分支 | 只投递一次 | 否 |
| idle Event | 活动 Turn 所有者 | 提交前 clear,终态后 set | 否 |
正常路径
asyncio.Queue(maxsize=32)把内存上界变成协议:第 33 个未消费项会让生产者在put()等待。- 输入 Queue 只有 Session 主循环消费,因此 Operation 顺序与提交顺序一致;不能让多个 worker 竞争后再猜测 Interrupt 属于哪个 Turn。
- 输出 Queue 的反压会传回
_emit(),继而暂停 Turn。客户端若不持续排空事件,Runtime 不承诺继续推进。 - 关闭时先取消并等待活动 Turn,再 flush/close Rollout,最后投递私有哨兵;事件迭代器看到哨兵后正常结束,不把它暴露为 AgentEvent。
- idle 不是队列空的同义词。用户输入入队前就 clear,只有拥有 Turn 终态的任务才能重新 set。
机制调用链
1. `asyncio.Queue(maxsize=32)` 把内存上界变成协议:第 33 个未消费项会让生产者在 `put()` 等待
→ 2. 输入 Queue 只有 Session 主循环消费,因此 Operation 顺序与提交顺序一致;不能让多个 worker 竞争后再猜测 Interrupt 属于哪个 Turn
→ 3. 输出 Queue 的反压会传回 `_emit()`,继而暂停 Turn
→ 4. 关闭时先取消并等待活动 Turn,再 flush/close Rollout,最后投递私有哨兵;事件迭代器看到哨兵后正常结束,不把它暴露为 AgentEvent
→ 5. idle 不是队列空的同义词
这里最重要的不是类名,而是控制权何时转移:创建者决定 ID 和初值,状态所有者决定何时修改,跨越 await 的调用必须明确取消、失败和可见性边界。任何绕过这些边界的“便捷调用”都会让恢复或并发测试失去确定答案。
Python 风格伪代码
CLOSED = object()
input_queue = asyncio.Queue(maxsize=32)
event_queue = asyncio.Queue(maxsize=32)
async def submit_user(text):
idle.clear() # before put; queue wait is part of submit
await input_queue.put(Submission(UserInput(text)))
async def events():
while True:
value = await event_queue.get()
if value is CLOSED:
return
yield cast(AgentEvent, value)
async def shutdown():
cancel_active_turn()
await join_active_turn()
await rollout.close()
await event_queue.put(CLOSED) # may backpressure until client drains
伪代码只保留设计职责;Mini Codex 的可运行版本见下方实现导航。它没有伪造官方源码中不存在的 Python API,也没有把路径策略写成 OS 沙箱。
失败、取消与恢复
| 故障或错误设计 | 会留下什么 | Mini Codex 的处理 |
|---|---|---|
| 事件消费者停止读取 | _emit 在满队列阻塞,Turn 暂停 | 这是显式反压;Host 必须持续消费或取消 |
用 None 当哨兵 | 未来合法载荷可能与哨兵冲突 | 使用模块私有 object identity |
| Shutdown 先发哨兵 | 客户端看见流结束,但活动 Turn 仍写日志/事件 | 先 join 与 close,再结束事件流 |
| 把 Queue.empty 当 idle | 活动任务可能仍运行 | idle 由 Turn 终态所有者管理 |
必须保持的不变量
队列满时等待而不是丢事件;关闭事件流之前,活动 Turn 与耐久写入必须已经收束。
这条不变量同时约束正常路径、异常路径和恢复路径。只在 happy path 里得到正确输出,不足以证明该模块边界成立。
设计取舍
官方 Codex 的 Event channel 当前可采用不同容量策略;Mini Codex 故意让输入和输出都 bounded,以便把消费者停滞暴露成测试可见行为。这是教学选择,不是一比一端口。
源码能够直接证明类型、分支、调用顺序和测试期望;关于工程动机的解释是基于这些事实的设计归纳,不冒充未公开承诺。
测试与复现实验
cd examples/mini-codex
uv run pytest -q -k 'test_wait_until_idle_after_submit_waits_for_new_turn' -k 'test_bounded_mailbox_applies_backpressure'
uv run mypy src
本节对应的关键断言:
test_wait_until_idle_after_submit_waits_for_new_turntest_bounded_mailbox_applies_backpressure
全量离线基线为 27 项通过;真实 Responses 测试需要显式环境变量,默认跳过。单项测试名用于定位,不代替对断言内容的解释。
官方源码导航
- codex-rs/core/src/session/mod.rs:Submission bounded channel、Event channel 与 Session loop
- codex-rs/core/src/codex_thread.rs:submit/next_event 的 Channel 边界
Mini Codex 对照
src/mini_codex/runtime/session.py:容量、put/get、哨兵和 idlesrc/mini_codex/runtime/thread.py:单消费者 AsyncIterator APIsrc/mini_codex/runtime/agents.py:Mailbox 容量与 put 反压的第二个实例
Mini Codex 保留本节的状态所有权、顺序和失败反馈;省略的产品能力会在边界处明确列出,不能由测试通过外推为生产等价。
本节边界
已经证明:队列满时等待而不是丢事件;关闭事件流之前,活动 Turn 与耐久写入必须已经收束。
尚未覆盖的生产问题由后续单元继续展开;公开版保留官方源码链接和可运行测试合同。
评论
登录后即可评论