Skip to content

Turn中断与运行中注入

从 Session handler 到 pending input 和 mailbox,相比 interrupt、steer、inject 与空闲启动的并发边界。

基于rust-v0.150.0
CodexRustRuntime

Turn中断与运行中注入 ​

本文回答一个具体问题:一个 Turn 正在运行时,新的用户输入、扩展上下文和中断分别如何进入 Session,什么时候仍属于当前 Turn,什么时候只能等待下一 Turn。重点是 Session、ActiveTurn、RunningTask 和 InputQueue 之间的边界;本文不把 mailbox 当作普通的无序消息队列,也不把“函数返回成功”解释为“模型已经消费了输入”。目标竞态测试覆盖了这些投递时机。

阅读前应先掌握 Session输入队列 中 pending input 与 mailbox 的区别。 本文不重新介绍 Session 字段,只追踪 interrupt、steer、inject 和 idle admission 如何改变活动 Turn 与 下一次 sampling。读完后应能从返回错误反查状态边界,并用目标测试验证投递时机。

1. 输入入口分流 ​

Op::Interrupt、typed turn input、steer 和 extension 注入都进入 Session,但它们的语义不同:interrupt 取消现有 task,TurnInputMode::Steer 要求当前 Turn 可被引导,inject 只尝试把 ResponseItem 放入正在运行的 Turn,StartIfIdle 才可能创建新 Turn。

相关源码:

  • codex-rs/core/src/session/handlers.rs :: interrupt
  • codex-rs/core/src/session/turn_input.rs :: handle, start_or_steer, start_if_idle
  • codex-rs/core/src/session/inject.rs :: inject_if_running, inject_no_new_turn

普通输入没有“无条件注入当前 Turn”的保证。调用方必须先知道是用户发起的新 Turn,还是对当前 regular Turn 的 steer;而 extension 的 inject_if_running 在没有活动 Turn 时返回原始 items,不会偷偷创建 Turn。

2. interrupt取消 ​

源码位置:codex-rs/core/src/session/mod.rs :: interrupt_task。实现先记录是否有 active Turn,再调用 abort_all_tasks(TurnAbortReason::Interrupted);无 active Turn 时取消 MCP startup,最后尝试启动 pending work。

interrupt_task 先记录是否存在活动 Turn,再进入统一的 abort_all_tasks。有活动 task 时,task runner 观察 cancellation,发送中止生命周期并清理 pending input 与 approvals;没有活动 Turn 时,中断的目标是尚未完成的 MCP startup。中断结束后仍会尝试处理已经排队的工作,因此“中断”不是关闭 Session。

源码位置:codex-rs/core/src/tasks/mod.rs :: abort_all_tasks。实现先取出 active Turn/task,让 handle_task_abort 观察 cancellation 并发送中止生命周期,再清理 pending input 与 approvals。

这里的顺序很重要:handle_task_abort 先让正在运行的 task 看到取消,并由 runner 保证中断 marker 的持久化顺序;pending input 不能在 marker 之前被无条件丢弃。测试 abort_regular_task_emits_marker_before_turn_aborted、abort_gracefully_emits_marker_before_turn_aborted 和 turn_aborted_flushes_terminal_event_after_delivery 覆盖了这一事件顺序。

3. steer边界 ​

源码位置:codex-rs/core/src/session/mod.rs :: steer_input

rust
pub async fn steer_input(
    &self,
    input: Vec<UserInput>,
    additional_context: BTreeMap<String, AdditionalContextEntry>,
    expected_turn_id: Option<&str>,
    client_user_message_id: Option<String>,
    responsesapi_client_metadata: Option<HashMap<String, String>>,
) -> Result<String, SteerInputError> {
    let mut active = self.active_turn.lock().await;
    let Some(active_turn) = active.as_mut() else {
        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
    {
        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)
    };
    if let Some(metadata) = responsesapi_client_metadata {
        active_task
            .turn_context
            .turn_metadata_state
            .set_responsesapi_client_metadata(metadata);
    }
    let mut pending = additional_context_input
        .into_iter()
        .map(ResponseItem::from)
        .map(TurnInput::ResponseItem)
        .collect::<Vec<_>>();
    pending.push(TurnInput::UserInput {
        content: input,
        client_id: client_user_message_id,
    });
    self.input_queue
        .extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
            active_turn.turn_state.as_ref(),
            pending,
        )
        .await;
    Ok(active_turn_id.clone())
}

steer 的检查分别保护不同不变量:没有正在执行 task 的活动 Turn 时不能引导;expected_turn_id 防止客户端把旧请求写进新 Turn;Review/Compact 等非 Regular task 不接受 steer;空用户输入不制造无意义的下一次采样。通过后,Session 级 additional context 先合并并转换为 ResponseItem,用户输入再作为带 client ID 的 UserInput 写入 pending queue,并重新打开当前 Turn 的 mailbox delivery。responsesapi_client_metadata 也在入队前写入该 Turn 的 metadata。

4. inject ​

源码位置:codex-rs/core/src/session/inject.rs :: 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) => {
            let pending = input.into_iter().map(TurnInput::ResponseItem).collect();
            self.input_queue
                .extend_pending_input_and_accept_mailbox_delivery_for_turn_state(
                    active_turn.turn_state.as_ref(), pending,
                )
                .await;
            Ok(())
        }
        None => Err(input),
    }
}

inject_if_running 的返回值表达了“是否找到消费者”:Ok(()) 只表示 items 已进入活动 Turn 的 pending queue,不表示模型已经读到;Err(items) 把原 items 交还调用方。inject_no_new_turn 随后会在无活动 Turn 时记录 conversation items,但明确不创建新 Turn;这保证后台结果不会因为注入 API 而凭空产生一次模型请求。

5. idle gate ​

源码位置:codex-rs/core/src/session/turn_input.rs :: start_if_idle

rust
async fn start_if_idle(
    session: &Arc<Session>,
    request: TurnInputRequest,
    submission_id: String,
    is_recovery: bool,
) -> CodexResult<TurnInputSubmission> {
    let has_user_input = has_nonempty_user_input(&request.input);
    if session.input_queue.has_trigger_turn_mailbox_items().await {
        return Ok(TurnInputSubmission::NotSubmitted {
            reason: NotSubmittedReason::PendingTriggerTurn,
        });
    }
    if !has_user_input && !is_recovery
        && session.collaboration_mode().await.mode == ModeKind::Plan
    {
        return Ok(TurnInputSubmission::NotSubmitted {
            reason: NotSubmittedReason::PlanMode,
        });
    }
    let turn_state = {
        let mut active = session.active_turn.lock().await;
        if active.is_some() {
            return Ok(TurnInputSubmission::NotSubmitted {
                reason: NotSubmittedReason::NotIdle,
            });
        }
        Arc::clone(&active.get_or_insert_with(ActiveTurn::default).turn_state)
    };
    // reservation 后再次检查 mailbox、Plan mode 和 reservation 有效性,
    // 失效时清除 reservation 并把原 input 放回错误值。
    // 后续分支应用 start settings 并启动 RegularTask;竞态失败时清除 reservation。
}

这个 gate 处理的是“是否可以从 idle 创建 RegularTask”,不是“如何把内容追加到当前 Turn”。它拒绝已有活动 Turn、待处理 trigger-turn mailbox,以及 Plan mode 下没有用户输入的自动工作。预留 ActiveTurn 后还要再次检查条件;竞态期间如果 trigger-turn 到达,必须清除 reservation 并把原 input 返回给调用方。相关测试覆盖 Busy、PlanMode、PendingTriggerTurn、Plan mode 中真实用户输入和 active Review。

6. mailbox ​

源码位置:codex-rs/core/src/state/turn.rs :: MailboxDeliveryPhase

rust
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub enum MailboxDeliveryPhase {
    #[default]
    CurrentTurn,
    NextTurn,
}

CurrentTurn 表示 mailbox 可以在当前 Turn 的下一次模型请求前被 drain;NextTurn 表示已经发出用户可见终态,晚到消息不能污染刚展示的答案。工具调用和显式 steer 会重新接受当前 Turn mailbox。get_pending_input 只有在 CurrentTurn 时才消费 mailbox,因此“消息已经到达”与“本次 sampling 会看到消息”是两个不同事实。

源码位置:codex-rs/core/src/session/input_queue.rs :: get_pending_input。它先取出该 Turn 的 pending input;只有 MailboxDeliveryPhase::CurrentTurn 时才继续 drain mailbox,NextTurn 则保留晚到消息。

队列中的 TurnInput 至少区分 UserInput、ResponseItem 和 inter-agent communication。它们共享 pending queue,却在历史写入、模型可见性和触发下一 Turn 的条件上不同;不要把三者都简化成字符串。

7. Turn Loop终止 ​

Turn 主循环会在一次 sampling 返回后,根据 can_drain_pending_input 和 model_needs_follow_up 判断是否继续。pending input 先经过 hook 和 history 写入,再参与下一次 sampling;如果当前 Turn 已跨过 answer boundary,则 mailbox 可能只会在新 Turn 中出现。

源码位置:codex-rs/core/src/session/turn.rs :: can_drain_pending_input、model_needs_follow_up、pending input hook/history 路径。

理解这段循环时要分开三个时间点:API 调用完成、pending queue 被 drain、下一次模型请求开始。steer/inject 只保证第一个时间点;第二、第三个时间点取决于 Turn 当前是否仍允许继续,以及 cancellation 是否已经被观察。

8. 失败路径 ​

  • 没有 active Turn 的 steer 返回 NoActiveTurn;错误不会创建 reservation。
  • expected_turn_id 过期返回 ExpectedTurnMismatch,避免旧客户端覆盖新 Turn。
  • Review/Compact 返回 ActiveTurnNotSteerable;它们不是 regular 对话的可插入点。
  • interrupt 期间 task 取消、TurnAborted 生命周期和 pending 清理按固定顺序发生;空 active Turn 的测试确认 pending input 可保留。
  • answer boundary 后 queue-only child mail 和 trigger-turn mail 留到下一 Turn;工具调用与显式 steer 会重新打开当前 Turn mailbox。

相关测试分别覆盖身份检查、拒绝原因、pending input 保留及 mailbox 相位;它们没有证明所有调度交错都可复现,也没有证明跨平台宿主行为完全相同。

9. 注入边界验证 ​

  1. 复述 Op::Interrupt、typed steer、inject_if_running 和 StartIfIdle 各自是否取消、追加当前 Turn 或创建新 Turn。
  2. 出现“注入 API 返回成功但模型未看到内容”时,依次检查 active Turn、MailboxDeliveryPhase、get_pending_input 和下一次 sampling 是否已经开始。

可以用下面的只读搜索把本文的 interrupt and mailbox 主线落回源码:

bash
rg -n "interrupt_task|TurnInputMode|inject_if_running|MailboxDeliveryPhase|start_if_idle" codex-rs/core/src