Session输入队列
Codex 的 Session 输入不是一个按到达时间排列的全局队列。新 Turn 的首批输入直接交给 SessionTask; 运行中 steer 和 inject 存在当前 TurnState;跨 Turn 的 Agent 通信存在 Session 级 mailbox;一个 watch channel 只通知“有活动”,并不承载消息本体。
这种拆分解决的是所有权,而不只是性能:普通用户输入必须优先完成本 Turn 的第一次 sampling,steer 只能进入匹配的 Regular Turn,迟到的 child mail 不能在最终回答已经展示后悄悄延长旧 Turn,而需要主动 唤醒的新消息又不能因为 Session 恰好处于 idle 而永久滞留。
阅读前先理解 CodexThread公共API 的 submit/steer入口,以及 Session启动上下文 的 sampling边界。本文只研究输入所有权与消费时机, 不逐分支解释 submission loop 中的所有 Op。
1. 输入路径
下图中的实线表示输入对象的移动,虚线表示活动通知。InputQueue 这个名字容易让人误以为它持有全部 输入;实际上它只直接拥有 mailbox,Turn-local pending items 位于 TurnState。
四条路径的差异可以压缩成下面这张表:
| 输入路径 | 数据形态 | idle 时 | Regular Turn 运行时 | Review/Compact 运行时 |
|---|---|---|---|---|
| 普通 UserInput | TurnInput::UserInput | 创建新 RegularTask,作为首批参数 | 尝试 steer 当前 Turn | 拒绝 steer |
| 显式 steer | UserInput + additional context | 返回 NoActiveTurn | 写入当前 TurnState | 返回 ActiveTurnNotSteerable |
| inject | ResponseItem | inject_if_running 原样退回 | 转成 TurnInput::ResponseItem | 只要求 active,不检查 TaskKind |
| Agent mailbox | InterAgentCommunication | trigger mail 可启动 synthetic Turn | 按 mailbox phase 决定当前或下一 Turn | 保留在 Session mailbox |
inject 与 steer 的 TaskKind 约束不同是源码事实:inject_if_running() 只检查 ActiveTurn 是否存在, steer_input() 则明确只接受 TaskKind::Regular。调用方不能把两者当作同义 API。
2. TurnInput
队列元素是一个三分枚举。用户输入保留多模态 UserInput 和 client ID;内部/扩展输入可以直接携带 ResponseItem;Agent 通信保留结构化 sender、recipient、content 与 trigger_turn,直到记录历史时才投影 成模型输入。
源码位置:codex-rs/core/src/session/input_queue.rs :: TurnInput与队列存储
/// Input consumed by a regular turn.
#[derive(Clone, Debug, PartialEq)]
pub enum TurnInput {
UserInput {
content: Vec<UserInput>,
// client_id用于事件/UI关联,不需要塞进模型可见文本。
client_id: Option<String>,
},
// Hook、additional context和inject可以直接保留协议级item结构。
ResponseItem(ResponseItem),
// Agent mail在消费前仍保留trigger_turn与拓扑身份。
InterAgentCommunication(InterAgentCommunication),
}
/// Turn-local pending input storage owned by the input queue flow.
#[derive(Default)]
pub(crate) struct TurnInputQueue {
// Vec维护同一Turn内追加顺序,所有权随TurnState结束。
items: Vec<TurnInput>,
}
/// Session-scoped pending input storage and active-turn mailbox delivery coordination.
pub(crate) struct InputQueue {
// watch只传活动类型,不传输入本体。
activity_tx: watch::Sender<InputQueueActivity>,
// mailbox必须跨Turn存活,因此不放进TurnState。
mailbox_pending_mails: Mutex<VecDeque<PendingMailboxCommunication>>,
}对象所有权如下。ActiveTurn 的 task 可以暂时为 None,这个中间态用于 idle turn reservation; TurnState 仍可在 reservation 和 RunningTask 安装之间保存 pending items。
TurnState 不只存输入,还存审批、权限请求、用户问答、MCP elicitation 和动态工具的 oneshot sender。 因此 abort 清理 pending input 与关闭 pending waiters 必须属于同一个 Turn 生命周期,不能只清一个全局消息 数组。
协议层的 TurnInputMode 现在明确区分三种入口:StartOrSteer、StartIfIdle 和带 expected_turn_id 的 Steer。handle_recovery() 复用 StartIfIdle,但用 is_recovery=true 允许空 用户输入恢复;普通空闲唤醒在 Plan mode 或存在 trigger-turn mailbox 时返回 NotSubmittedReason::PlanMode / PendingTriggerTurn,不会创建空 Turn。所有成功结果都只是“已接受输入”, 不等待 hook、history 写入或模型采样完成。
源码位置:codex-rs/core/src/state/turn.rs :: ActiveTurn与TurnState
pub(crate) struct ActiveTurn {
// reservation阶段task为None,正式start_task后才安装RunningTask。
pub(crate) task: Option<RunningTask>,
pub(crate) turn_state: Arc<Mutex<TurnState>>,
}
#[derive(Default)]
pub(crate) struct TurnState {
pending_approvals: HashMap<String, oneshot::Sender<ReviewDecision>>,
pending_request_permissions: HashMap<String, PendingRequestPermissions>,
pending_user_input: HashMap<String, oneshot::Sender<RequestUserInputResponse>>,
pending_elicitations: HashMap<(String, RequestId), oneshot::Sender<ElicitationResponse>>,
pending_dynamic_tools: HashMap<String, oneshot::Sender<DynamicToolResponse>>,
// steer/inject与已转移进来的mailbox items都归当前TurnState所有。
pub(crate) pending_input: TurnInputQueue,
mailbox_delivery_phase: MailboxDeliveryPhase,
granted_permissions_by_environment_id: HashMap<String, AdditionalPermissionProfile>,
pub(crate) tool_calls: u64,
pub(crate) has_memory_citation: bool,
pub(crate) token_usage_at_turn_start: TokenUsage,
// 其余MCP审批metadata与auto-review字段省略。
}3. 普通用户输入
当前普通输入、steer、start-if-idle 和 recovery 都通过 TurnInputRequest + TurnInputMode 路由。 StartOrSteer 先尝试 steer;只有返回 NotSubmittedReason::NoActiveTurn 才应用 start settings 并 启动 RegularTask。其他拒绝原因会转换为 typed TurnInputSubmission::NotSubmitted,不会悄悄创建新 Turn。
正常新 Turn 的首批输入不先进入 TurnInputQueue,而是成为 SessionTask::run(input) 的参数。这样首轮 sampling 可以明确区分“启动 Turn 的输入”和运行中后来到的 pending input。
源码位置:codex-rs/core/src/session/turn_input.rs :: start_or_steer
match session.steer_input(
&mut items,
additional_context.clone(),
/*expected_turn_id*/ None,
settings.required_active_final_output_json_schema(),
client_id.clone(),
responsesapi_client_metadata.clone(),
incoming_root_turn_id,
).await {
Ok(turn_id) => {
settings.apply_steered(session, submission_id).await?;
Ok(TurnInputSubmission::Steered { turn_id })
}
Err(NotSubmittedReason::NoActiveTurn) => {
let turn_context = settings.apply_started(session, submission_id.clone()).await?;
let mut task_input = merge_additional_context_input(session, additional_context).await;
if !items.is_empty() {
task_input.push(TurnInput::UserInput { content: items, client_id });
}
session.start_task(
turn_context,
task_input,
RegularTask::new(),
MailboxParentProvenance::Ignore,
).await;
Ok(TurnInputSubmission::Started { turn_id: submission_id })
}
Err(reason) => Ok(TurnInputSubmission::NotSubmitted { reason }),
}这里没有“所有 steer 失败都新建 Turn”的策略。只有 NoActiveTurn 才能走 Started 分支;如果存在 Review 或 Compact task,普通用户输入会得到明确错误,避免与不可 steer 的任务并发运行。
4. steer 锁范围
显式 steer_input() 支持 expected_turn_id,用于防止客户端看到 Turn A 后,网络延迟导致 steer 误投给 刚替换上来的 Turn B。函数持有 active_turn 锁完成 task identity、TaskKind、空输入检查和入队,源码用 clippy expect 明确承认这是有意的跨 await 原子区。
源码位置:codex-rs/core/src/session/mod.rs :: Session::steer_input(核心路径)
let mut active = self.active_turn.lock().await;
let Some(active_turn) = active.as_mut() else {
// 原input随错误退回,调用方可决定启动新Turn或保留。
return Err(NotSubmittedReason::NoActiveTurn);
};
let Some(active_task) = active_turn.task.as_ref() else {
return Err(NotSubmittedReason::NoActiveTurn);
};
let active_turn_id = &active_task.turn_context.sub_id;
if let Some(expected_turn_id) = expected_turn_id
&& expected_turn_id != active_turn_id
{
// 防止旧客户端请求跨Turn串入新的active task。
return Err(NotSubmittedReason::ExpectedTurnMismatch {
expected: expected_turn_id.to_string(),
actual: active_turn_id.clone(),
});
}
match active_task.kind {
TaskKind::Regular => {}
TaskKind::Review => {
return Err(NotSubmittedReason::ActiveTurnNotSteerable {
turn_kind: NonSteerableTurnKind::Review,
});
}
TaskKind::Compact => {
return Err(NotSubmittedReason::ActiveTurnNotSteerable {
turn_kind: NonSteerableTurnKind::Compact,
});
}
}
if input.is_empty() {
return Err(NotSubmittedReason::EmptyInput);
}
let additional_context_input = {
let mut state = self.state.lock().await;
state.additional_context.merge(additional_context)
};
let mut pending_input = additional_context_input
.into_iter()
.map(ResponseItem::from)
.map(TurnInput::ResponseItem)
.collect::<Vec<_>>();
pending_input.push(TurnInput::UserInput {
content: input,
client_id: client_user_message_id,
});
self.input_queue
// 入队steer同时把mailbox phase重新打开到CurrentTurn。
.extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
active_turn.turn_state.as_ref(),
pending_input,
)
.await;
Ok(active_turn_id.clone())steer 没有抢占正在执行的模型请求。它只是追加 pending items 并发送 activity signal;run_turn() 在安全的 sampling 边界读取这些输入,再经过 UserPromptSubmit hook、历史记录和新的 StepContext capture 后进入下次 模型请求。
5. inject
inject_if_running() 接受已经是 ResponseItem 的模型可见对象。只要存在 ActiveTurn 就转成 TurnInput::ResponseItem 放入当前 Turn queue,并重新允许 mailbox 当前轮投递;没有 active 时把原 Vec 原样放进 Err,不会自行写历史或启动任务。
源码位置:codex-rs/core/src/session/inject.rs :: Session::inject_if_running
pub async fn inject_if_running(
&self,
input: Vec<ResponseItem>,
) -> Result<(), Vec<ResponseItem>> {
let mut active = self.active_turn.lock().await;
match active.as_mut() {
Some(active_turn) => {
self.input_queue
.extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
active_turn.turn_state.as_ref(),
// 保留ResponseItem结构,只包一层TurnInput来源标签。
input.into_iter().map(TurnInput::ResponseItem).collect(),
)
.await;
Ok(())
}
// 所有权退给调用方,避免idle时输入被悄悄丢弃。
None => Err(input),
}
}内部还有 inject_no_new_turn():若运行中则走同一 pending queue;若 idle,则创建一个默认 TurnContext 仅用于 给 item 补元数据并直接记录历史,不启动模型 Turn。CodexThread::inject_user_message_without_turn() 的 session-prefix 行为就建立在这条路径上。
源码位置:codex-rs/core/src/session/inject.rs :: Session::inject_no_new_turn
pub(crate) async fn inject_no_new_turn(
&self,
items: Vec<ResponseItem>,
current_turn_context: Option<&TurnContext>,
) {
let Err(items) = self.inject_if_running(items).await else {
// active时已转移到TurnState,等待当前Task消费。
return;
};
let default_turn_context;
let turn_context = match current_turn_context {
Some(turn_context) => turn_context,
None => {
default_turn_context = self.new_default_turn().await;
default_turn_context.as_ref()
}
};
// idle时只写history,不产生TurnStarted或sampling。
self.record_conversation_items(turn_context, &items).await;
}6. Mailbox跨Turn
Agent mail 存在 Session 级 VecDeque,enqueue 始终 push_back,drain 使用 drain(..),因此同一 Session 观察到 FIFO 顺序。trigger_turn 决定 idle 时是否唤醒 RegularTask;queue-only mail 默认等待下一次由其他 输入启动的 Turn,但 durable sleep 存在时任何 mail 都能唤醒。
源码位置:codex-rs/core/src/session/input_queue.rs
// :: enqueue_mailbox_communication与drain_mailbox_input_items
pub(crate) async fn enqueue_mailbox_communication(
&self,
communication: InterAgentCommunication,
parent_turn_id: Option<String>,
) {
self.mailbox_pending_mails
.lock()
.await
// 只从尾部追加,drain时保持到达顺序。
.push_back(PendingMailboxCommunication {
communication,
parent_turn_id,
});
self.activity_tx.send_replace(InputQueueActivity::Mailbox);
}
pub(crate) async fn drain_mailbox_input_items(
&self,
) -> (Vec<TurnInput>, Option<String>) {
let pending_mails = self
.mailbox_pending_mails
.lock()
.await
.drain(..)
.collect::<Vec<_>>();
let parent_turn_id = pending_mails
.iter()
// queue-only mail不参与synthetic Turn的parent归因。
.filter(|mail| mail.communication.trigger_turn)
.map(|mail| mail.parent_turn_id.as_deref())
// 所有trigger mail必须一致指向同一非空parent,否则放弃归因。
.reduce(|expected, candidate| expected.filter(|id| candidate == Some(*id)))
.and_then(|id| id.filter(|id| !id.trim().is_empty()).map(str::to_string));
let items = pending_mails
.into_iter()
.map(|mail| TurnInput::InterAgentCommunication(mail.communication))
.collect();
(items, parent_turn_id)
}inter_agent_communication() 在 enqueue 后只对 trigger mail 或 durable sleep 调度 pending-work scheduler。 真正启动前仍会复查 mailbox 和 idle 状态,避免两个并发通知各自创建一个 Turn。
synthetic Turn 的 start_task() 使用 MailboxParentProvenance::Attribute,普通用户启动和 replace task 使用 Ignore。parent turn ID 因而只会附着到确由 mailbox 唤醒、且来源没有歧义的 Turn。
7. Mailbox边界
Mailbox 是否能进入当前 Turn 不是简单布尔配置,而是 MailboxDeliveryPhase 状态机。Turn 默认 CurrentTurn;记录被判定为用户可见终态的 assistant item 后切到 NextTurn;显式 steer、模型要求 follow-up、工具调用或 stop-hook continuation 又会重新打开 CurrentTurn。
状态切换不是为了阻止 mailbox 入队,而是控制 has_pending_input() 和 get_pending_input() 是否把它视作 当前 Turn 的 follow-up。消息本体始终留在 Session mailbox,下一 Turn 仍能读取。
源码位置:codex-rs/core/src/session/input_queue.rs
// :: defer与accept mailbox delivery
pub(crate) async fn defer_mailbox_delivery_to_next_turn(
&self,
active_turn: &Mutex<Option<ActiveTurn>>,
sub_id: &str,
) {
let turn_state = self.turn_state_for_sub_id(active_turn, sub_id).await;
let Some(turn_state) = turn_state else {
// 过期TurnId不能改变新Turn的mailbox phase。
return;
};
let mut turn_state = turn_state.lock().await;
if turn_state.pending_input.items.iter().any(|input| {
!matches!(
input,
TurnInput::InterAgentCommunication(communication)
if !communication.trigger_turn
)
}) {
// 已有steer、inject或trigger mail意味着明确same-turn工作,旧defer不能覆盖它。
return;
}
turn_state.set_mailbox_delivery_phase(MailboxDeliveryPhase::NextTurn);
}
pub(super) async fn extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
&self,
turn_state: &Mutex<TurnState>,
input: Vec<TurnInput>,
) {
{
let mut turn_state = turn_state.lock().await;
turn_state.pending_input.items.extend(input);
// 显式same-turn输入到达后,之前的answer boundary不再关闭mailbox。
turn_state.accept_mailbox_delivery_for_current_turn();
}
self.activity_tx.send_replace(InputQueueActivity::Steer);
}turn_state_for_sub_id() 同时检查 active task 的 turn_context.sub_id。这使异步流处理器晚到的 defer/accept 操作无法误改刚替换的新 TurnState。
8. 消费顺序由边界
get_pending_input() 在一个原子区内读取 ActiveTurn、TurnState items 和 mailbox phase。phase 为 CurrentTurn 时先 split_off(0) 取走 Turn-local items,再 drain mailbox 并追加到尾部;所以同一次消费 的顺序是 steer/inject 在前、mailbox 在后。phase 为 NextTurn 时两者都不交给当前 Turn。
源码位置:codex-rs/core/src/session/input_queue.rs :: InputQueue::get_pending_input
let (pending_input, accepts_mailbox_delivery) = {
let mut active = active_turn.lock().await;
match active.as_mut() {
Some(active_turn) => {
let mut turn_state = active_turn.turn_state.lock().await;
let accepts_mailbox_delivery =
turn_state.accepts_mailbox_delivery_for_current_turn();
let pending_input = if accepts_mailbox_delivery {
// split_off(0)一次性转移Turn-local Vec并留下空队列。
turn_state.pending_input.items.split_off(0)
} else {
Vec::new()
};
(pending_input, accepts_mailbox_delivery)
}
// idle时没有旧Turn边界限制,mailbox可由新Turn drain。
None => (Vec::new(), true),
}
};
if !accepts_mailbox_delivery {
return (pending_input, None);
}
let (mailbox_items, parent_turn_id) = self.drain_mailbox_input_items().await;
if pending_input.is_empty() {
(mailbox_items, parent_turn_id)
} else {
let mut pending_input = pending_input;
// 同一批消费中mailbox排在已经等待的steer/inject之后。
pending_input.extend(mailbox_items);
(pending_input, parent_turn_id)
}但“同批顺序”不等于所有输入按全局到达时间排序。新 Turn 的 direct input 不在 queue 中,run_turn() 刻意先 sampling 它;只有第一次 sampling 结束后才允许 drain 运行期间到达的 steer/mailbox。自动 compact 后 也可能暂缓 drain,让模型/工具 continuation 先恢复。
对应的 run_turn() 循环把 drain 点放在每次 StepContext capture 之前:
源码位置:codex-rs/core/src/session/turn.rs :: run_turn(pending消费节选)
// 有首批direct input时先采样它;空输入synthetic Turn可立即drain mailbox。
let mut can_drain_pending_input = input.is_empty();
// 此处省略首批input的hook记录和Skill、Plugin注入。
let mut next_step_context = Some(first_step_context);
loop {
let pending_input = if can_drain_pending_input {
sess.input_queue
.get_pending_input(&sess.active_turn)
.await
.0
} else {
Vec::new()
};
if run_hooks_and_record_inputs(&sess, &turn_context, &pending_input).await {
break;
}
let step_context = match next_step_context.take() {
Some(step_context) => step_context,
None if pending_input.is_empty() => {
sess.capture_step_context(
Arc::clone(&turn_context),
&cancellation_token,
)
.await?
}
None => {
// 新pending input可能显式提到MCP server,因此消费后重新解析required servers。
let pending_user_input = turn_user_input(&pending_input);
let (required_servers, _) = required_mcp_servers_for_input(
&sess,
turn_context.as_ref(),
&pending_user_input,
)
.or_cancel(&cancellation_token)
.await?;
sess.capture_step_context_with_required_mcp_servers(
Arc::clone(&turn_context),
&cancellation_token,
&required_servers,
)
.await?
}
};
// sampling完成后can_drain变为true,下一循环才接纳后来输入。
// 此处省略run_sampling_request及follow-up判定。
}Streaming 期间还有一条 mailbox 低延迟路径:当模型刚输出 commentary 或 reasoning,若 mailbox 已有消息, Core 提前结束当前 stream 消费并令 needs_follow_up=true,尽快进入下一次循环。最终 AgentMessage 不触发 这类 preemption,仍由 answer-boundary phase 决定是否留给下一 Turn。
9. Watch Channel
activity_tx 使用 Tokio watch。连续多次 send_replace() 会合并为最新活动值,因此消费者不能通过 changed 次数推断消息数量。订阅时 Core 会额外检查现有 TurnState 和 mailbox,弥补“先入队、后订阅”不会重放 每次历史通知的问题;已有 steer 优先报告为 Steer,否则才报告 Mailbox。
源码位置:codex-rs/core/src/session/input_queue.rs :: InputQueue::subscribe_activity
pub(crate) async fn subscribe_activity(
&self,
turn_state: Option<&Mutex<TurnState>>,
) -> (
watch::Receiver<InputQueueActivity>,
Option<InputQueueActivity>,
) {
let activity_rx = self.activity_tx.subscribe();
let has_pending_steer = if let Some(turn_state) = turn_state {
// 只把真实UserInput视为pending steer;内部ResponseItem不冒充用户steer状态。
turn_state.lock().await.pending_input.has_user_input()
} else {
false
};
let pending_activity = if has_pending_steer {
Some(InputQueueActivity::Steer)
} else if self.has_pending_mailbox_items().await {
Some(InputQueueActivity::Mailbox)
} else {
None
};
(activity_rx, pending_activity)
}这个 API 的正确用法是:收到 activity 后回到队列查询真实状态。把 watch 当作计数 channel 会漏事件, 把 InputQueueActivity::Mailbox 当作“队列中只有一个 mail”同样错误。
10. Task 结束与中断
Task 正常结束时可能仍有 race 中迟到的 Turn-local input。on_task_finished() 先从 ActiveTurn 取走 task, 再 drain TurnInputQueue,逐项经过 hook 并记录历史;它不会为这些 leftover items 再发模型请求,但也不会 静默丢弃用户输入。随后只有确认 active_turn 仍指向同一个 TurnState 且 task 为 None,才清空 active slot。
源码位置:codex-rs/core/src/tasks/mod.rs :: Session::on_task_finished(队列收尾节选)
let turn_state = {
let mut active = self.active_turn.lock().await;
active.as_mut().and_then(|active_turn| {
let task = active_turn.task.take()?;
task.handle.detach();
Some(Arc::clone(&active_turn.turn_state))
})
};
let Some(turn_state) = turn_state else {
return;
};
let pending_input = self
.input_queue
.take_pending_input_for_turn_state(turn_state.as_ref())
.await;
if !pending_input.is_empty() {
for pending_input_item in pending_input {
let hook_outcome =
inspect_pending_input(self, &turn_context, &pending_input_item).await;
if hook_outcome.should_stop {
record_additional_contexts(
self,
&turn_context,
hook_outcome.additional_contexts,
)
.await;
} else {
// leftover仍进入history和用户消息事件,但不会重新调用sampling。
record_pending_input(
self,
&turn_context,
pending_input_item,
hook_outcome.additional_contexts,
)
.await;
}
}
}
let cleared_active_turn = {
let mut active = self.active_turn.lock().await;
if let Some(active_turn) = active.as_ref()
&& active_turn.task.is_none()
&& Arc::ptr_eq(&active_turn.turn_state, &turn_state)
{
// 指针身份检查避免旧Task完成回调清掉后来安装的新ActiveTurn。
*active = None;
true
} else {
false
}
};显式中断实际 RunningTask 时,Core 取消 task、发送 abort lifecycle,再清除该 Turn 的 waiter 与 pending input;Session mailbox 不在 clear_pending() 的清理范围内。Interrupted 之后 scheduler 会复查 mailbox, trigger mail 仍可启动下一 Turn。
源码位置:codex-rs/core/src/tasks/mod.rs :: Session::abort_all_tasks(清理节选)
if let Some(mut active_turn) = self.take_active_turn().await {
let task = active_turn.task.take();
aborted_turn = task.is_some();
turn_context = task.as_ref().map(|task| Arc::clone(&task.turn_context));
if let Some(task) = task {
self.handle_task_abort(task, reason.clone()).await;
}
if aborted_turn {
active_turn_to_clear = Some(active_turn);
}
}
if let Some(active_turn) = active_turn_to_clear {
// 这里只清TurnState waiters和Turn-local items,不清Session mailbox。
self.input_queue.clear_pending(&active_turn).await;
}
if reason == TurnAbortReason::Interrupted && aborted_turn {
// 中断旧Turn后,trigger mailbox仍有机会唤醒新Turn。
self.maybe_start_turn_for_pending_work().await;
}一个特殊测试锁定了 reservation 语义:ActiveTurn 存在但 task 仍为 None 时,abort_all_tasks(Replaced) 不会把预先缓存在该 TurnState 的输入清掉。否则自动 idle work 在 reservation 到 task 安装之间被替换时, 输入会无声消失。
11. 用竞态测试
输入队列的关键错误大多来自时序,因此测试不是只检查单次 enqueue/dequeue。当前覆盖包括 expected turn 不匹配、Review/Compact 拒绝、leftover 落历史、answer boundary 后 mailbox 延迟、steer 重新开放 mailbox, 以及 stale defer 不能覆盖新 steer。
下面的测试直接验证 NextTurn phase 不删除 queue-only mail,只让当前 Turn 暂时看不见它:
源码位置:codex-rs/core/src/session/tests.rs
// :: queue_only_mailbox_mail_waits_for_next_turn_after_answer_boundary(节选)
sess.input_queue
.defer_mailbox_delivery_to_next_turn(&sess.active_turn, &tc.sub_id)
.await;
sess.input_queue
.enqueue_mailbox_communication(
communication.clone(),
/*parent_turn_id*/ None,
)
.await;
assert!(
// has_pending_input对当前Turn返回false,避免已展示答案后自动follow-up。
!sess.input_queue.has_pending_input(&sess.active_turn).await,
);
assert_eq!(
sess.input_queue.get_pending_input(&sess.active_turn).await,
(Vec::new(), None),
);
sess.abort_all_tasks(TurnAbortReason::Replaced).await;
assert_eq!(
// active Turn结束后,同一mail仍在Session mailbox并可被下一Turn取得。
(sess.input_queue.get_pending_input(&sess.active_turn).await).0,
vec![TurnInput::InterAgentCommunication(communication)],
);steer 到来后,消费顺序应是 steer 在前、此前等待的 mailbox 在后:
源码位置:codex-rs/core/src/session/tests.rs
// :: steered_input_reopens_mailbox_delivery_for_current_turn(断言节选)
sess.steer_input(
vec![UserInput::Text {
text: "follow up".to_string(),
text_elements: Vec::new(),
}],
/*additional_context*/ Default::default(),
Some(&tc.sub_id),
/*client_user_message_id*/ None,
/*responsesapi_client_metadata*/ None,
)
.await
.expect("steered input should be accepted");
assert_eq!(
(sess.input_queue.get_pending_input(&sess.active_turn).await).0,
vec![
TurnInput::UserInput {
content: vec![UserInput::Text {
text: "follow up".to_string(),
text_elements: Vec::new(),
}],
client_id: None,
},
// reopen后mailbox追加到同批Turn-local steer之后。
TurnInput::InterAgentCommunication(communication),
],
);排查输入“丢失、插错轮次或重复 sampling”时,可以按所有权定位:
- 首条用户消息未出现:检查新 Task 的 direct
input,它本来就不在TurnInputQueue。 - steer 返回成功但没有立即打断模型:这是预期行为;它在 sampling 边界消费,不是硬抢占。
- steer 投到错误 Turn:检查调用方是否提供
expected_turn_id,以及是否处理ExpectedTurnMismatch。 - inject 在 idle 时无效果:
inject_if_running会原样退回,调用方需要显式选择记录历史或启动 Turn。 - child mail 在最终回答后未被当前 Turn 消费:检查
MailboxDeliveryPhase::NextTurn;消息可能仍安全地留在 Session mailbox。 - mailbox 顺序异常:区分 Turn-local items 与 mailbox FIFO;Core 只保证 mailbox 内部 FIFO,以及同批消费 时 local items 先于 mailbox,并不存在跨三种存储的全局到达顺序。
- 活动通知次数少于消息数:watch signal 会合并,必须回查真实队列。
- 中断后 steer 消失但 mail 仍在:RunningTask abort 会清 TurnState,Session mailbox 具有独立生命周期。
Session 输入队列真正维护的是三条边界:输入属于哪个 Turn、何时允许形成下一次 sampling,以及哪些消息 必须跨 Turn 保存。理解 direct task input、Turn-local pending queue、Session mailbox 和 activity signal 的 区别,才能正确解释 steer、inject 与多 Agent 通信在竞态下的行为。
控制面如何调用 steer、inject 和 mailbox handler 见 Session运行时处理; 中断与关闭时哪些输入保留或清理见 Session关闭流程。
可以用下面的只读搜索把本文的 input ownership 主线落回源码:
rg -n "TurnInput|MailboxDeliveryPhase|get_pending_input|inject" codex-rs/core/src