InputQueue 不是 Submission channel 的别名。它把“已经由主循环接收、但尚未进入模型历史”的内容保存在 TurnState 中,并另设 Session 级 mailbox deque。两级存储解决了活动 Turn 和空闲 Session 对输入归属不同的问题。
TurnInput 的三种内容
TurnInput 可以是 UserInput、已经构造好的 ResponseItem,或 InterAgentCommunication。Steer 产生前两类:additional context 先转换为 ResponseItem,用户正文带 client id 追加在后。Agent mailbox 产生第三类。
class InputQueue:
activity: Watch[Mailbox | Steer]
mailbox: Mutex[Deque[PendingMailboxCommunication]]
class TurnInputQueue:
items: list[TurnInput]
async def enqueue_steer(turn_state, context_items, user_input):
async with turn_state.lock() as state:
state.pending_input.items.extend(context_items)
state.pending_input.items.append(user_input)
state.accept_mailbox_delivery_for_current_turn()
activity.send_replace(STEER)
activity watch 用于通知等待输入的组件,只表达“最近发生 Mailbox 或 Steer 活动”,不承载消息本身。订阅时会主动检查已有 pending steer/mail,补偿 watch 订阅前已经发生的更新。
排空为什么要同时锁 ActiveTurn 和 TurnState
get_pending_input 需要原子判断当前 ActiveTurn、读取 mailbox delivery phase,并从 TurnState split_off 全部 items。若先释放 ActiveTurn 锁再取 TurnState,Turn 可能已完成并被替换,旧消费者会偷走新 Turn 的输入。
只有当前 phase 接受 mailbox delivery 时,函数才进一步 drain Session mailbox。若 Turn pending 非空,mailbox items 追加在后,保持 steer 已进入的先后语义。Mailbox deque 本身按接收顺序 drain。
async def get_pending_input(active_turn):
async with active_turn.lock() as active:
if active is None:
turn_items, accepts_mail = [], True
else:
async with active.turn_state.lock() as state:
accepts_mail = state.accepts_mailbox_current_turn()
turn_items = state.pending.split_all() if accepts_mail else []
if not accepts_mail:
return turn_items, None
mail_items, parent = await drain_mailbox_in_order()
return turn_items + mail_items, parent
MailboxDeliveryPhase 防止错误续采样
并非所有 Agent mail 都应延长当前 Turn。没有 trigger_turn 的普通子 Agent 通知,在模型已经给出最终回答时可以留给下一轮。defer_mailbox_delivery_to_next_turn 检查当前 pending items:只要存在显式用户输入或 trigger mail,仍须当前轮 follow-up;只有 queue-only 非触发 mail 才把 phase 改成 NextTurn。
一旦模型产生工具调用、Steer 到达或其他明确继续信号,accept_mailbox_delivery_for_current_turn 可以把 phase 拉回当前轮。这是一个交付策略状态机,不是简单 bool。
完成与中断的队列处理不同
正常完成时 on_task_finished 取走残留 pending input,逐项运行 inspect hook并记录到历史,保证完成竞态到达的输入不会无声消失。Interrupt 则在任务观察取消后清空 pending waiter 和 items;用户中断的含义是停止旧 Turn,而不是偷偷把旧 steer 搬进下一个 Turn。
评论
登录后即可评论