阶段 4:事件总线

📅 2026/8/18 9:47:07
阶段 4:事件总线
阶段 4事件总线对应代码stage04_events.py学习目标理解事件既是给消费者的消息也是驱动派生状态的输入并实现terminal_error的单写多读模式。源码锚点概念位置说明EventMsg枚举protocol/src/protocol.rs约 L1900-3700数十种事件11.1 节的五种投影形态三层发送链core/src/session/mod.rs:1828-2077send_event → send_event_raw_with_persistence → deliver_event_rawterminal_error 锁存mod.rs:1828-1841Error事件副作用写入终态读取14.6 节单写多读AgentStatus watchmod.rs:2069-2077事件同时更新进程内派生状态代码走读EventBus.send()的顺序就是真实send_event的骨架ifmsgERRORandaffects_turn_status:self.terminal_errorpayload# 先副作用锁存错误if(s:self.agent_status_from(msg)):# 再更新派生状态self.agent_statuss self._queue.put_nowait(Event(...))# 最后投递两个教学要点锁存的过滤条件是affects_turn_status——真实版里 16 个CodexErrorInfo变体只有 2 个返回 false14.1 节。演示里的stream_error不锁存重试中间态error锁存——这正是willRetry语义的分界11.7 节。队列 unbounded发送方永不阻塞但消费者必须持续 drain3.7 节的警告同样适用。运行与预期输出python stage04_events.py# 打印 7 条事件最后两行# 最终 agent_status AgentStatus.completed - 最后一个生命周期事件决定# 锁存的 terminal_error boom - 终态投影 failed 的依据注意agent_status是 completed 而terminal_error是 boom——两者不冲突因为真实系统的AgentStatus描述 Core 生命周期Turn 失败是 App Server 的投影5.3 节的两层状态。练习给send()加legacy双发某些事件同时投递旧名字对应 11.2 节的四条扇出路径之一。把队列改成maxsize10并一次塞 20 条观察背压——然后回答为什么真实设计选 unbounded 事件通道 有界提交通道实现affects_turn_status(msg)函数表把过滤条件从布尔参数改成查表。与真实实现的差距真实send_event在 raw 层之前还有 trace、父代理转发、realtime 镜像、legacy 双发四条扇出11.2 节四条扇出路径。持久化过滤在 raw 层内send_event_raw_with_persistence阶段 9 会补。代码积木 4事件总线。 真实对应物 - codex-rs/protocol/src/protocol.rs 的 EventMsg 枚举数十种事件 - codex-rs/core/src/session/mod.rs:1828-2077 的三层发送链 send_event - send_event_raw_with_persistence - deliver_event_raw - terminal_error 锁存mod.rs:1828-1841分析文档 14.6 节单写多读 - AgentStatus watchmod.rs:2069-2077 核心认知事件既是给消费者的消息也是驱动派生状态的输入 terminal_error、AgentStatus 都由事件副作用维护。 运行python stage04_events.py from__future__importannotationsimportasynciofromdataclassesimportdataclassfromenumimportEnumclassTurnAbortReason(str,Enum):对应 protocol.rs:4209-4214。interruptedinterruptedreplacedreplacedclassAgentStatus(str,Enum):对应 core/src/agent/status.rs。pending_initpendingInitrunningrunningcompletedcompletedinterruptedinterruptederrorederroreddataclassclassEvent:对应 Event { id: turn 子 id, msg }。payload 是教学版附加的文本。id:strmsg:strpayload:str# —— 教学版认得的全部事件真实版每个都是独立 dataclass——TURN_STARTEDturn_startedITEM_STARTEDitem_startedITEM_COMPLETEDitem_completedAGENT_DELTAagent_message_deltaERRORerrorTURN_COMPLETEturn_completeTURN_ABORTEDturn_abortedclassEventBus:对应 tx_event(unbounded) deliver_event_raw 的 AgentStatus watch。def__init__(self):self._queue:asyncio.Queue[Event]asyncio.Queue()# 真实unboundedself.agent_statusAgentStatus.pending_init# 对应 watch 值self.terminal_error:str|NoneNone# 对应 TurnContext.terminal_errordefagent_status_from(self,msg:str)-AgentStatus|None:对应 agent_status_from_event() 的最小子集。return{TURN_STARTED:AgentStatus.running,TURN_COMPLETE:AgentStatus.completed,TURN_ABORTED:AgentStatus.interrupted,ERROR:AgentStatus.errored,}.get(msg)asyncdefsend(self,turn_id:str,msg:str,affects_turn_status:boolFalse,payload:str,)-None:对应 send_event()先副作用锁存错误再投递再更新派生状态。 真实版在 raw 层之前还有 trace / 父代理转发 / realtime 镜像 / legacy 双发 分析文档 11.2 节的四条扇出路径。 ifmsgERRORandaffects_turn_status:self.terminal_errorpayloadormsg# mod.rs:1828-1841 的锁存语义if(s:self.agent_status_from(msg))isnotNone:self.agent_statuss self._queue.put_nowait(Event(turn_id,msg,payload))asyncdefnext_event(self)-Event:returnawaitself._queue.get()asyncdefdemo()-None:busEventBus()turn_idt1awaitbus.send(turn_id,TURN_STARTED)awaitbus.send(turn_id,ITEM_STARTED,payloadc1)awaitbus.send(turn_id,AGENT_DELTA,payload正在跑命令...)awaitbus.send(turn_id,ITEM_COMPLETED,payloadc1)# 一次不影响 Turn 状态的警告真实版StreamError - willRetry不锁存awaitbus.send(turn_id,stream_error,payloadReconnecting... 1/5)# 一次影响 Turn 状态的错误 - 锁存进 terminal_errorawaitbus.send(turn_id,ERROR,affects_turn_statusTrue,payloadboom)awaitbus.send(turn_id,TURN_COMPLETE)whilenotbus._queue.empty():evawaitbus.next_event()print(f *{ev.id}{ev.msg:18}{ev.payload})print(最终 agent_status ,bus.agent_status)print(锁存的 terminal_error ,bus.terminal_error)# 终态投影 failed 的依据if__name____main__:asyncio.run(demo())