Skip to content

Session输入队列

分析普通输入、steer、inject 和 Agent mailbox 的存储位置、消费时机、唤醒规则与跨 Turn 隔离。

基于rust-v0.150.0
CodexRustSessionQueue

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 运行时
普通 UserInputTurnInput::UserInput创建新 RegularTask,作为首批参数尝试 steer 当前 Turn拒绝 steer
显式 steerUserInput + additional context返回 NoActiveTurn写入当前 TurnState返回 ActiveTurnNotSteerable
injectResponseIteminject_if_running 原样退回转成 TurnInput::ResponseItem只要求 active,不检查 TaskKind
Agent mailboxInterAgentCommunicationtrigger 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与队列存储

rust
/// 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

rust
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

rust
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(核心路径)

rust
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

rust
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

rust
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

rust
// :: 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

rust
// :: 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

rust
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消费节选)

rust
// 有首批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

rust
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(队列收尾节选)

rust
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(清理节选)

rust
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

rust
// :: 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

rust
// :: 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 主线落回源码:

bash
rg -n "TurnInput|MailboxDeliveryPhase|get_pending_input|inject" codex-rs/core/src