Skip to content

Session运行时处理

拆解 Submission 串行分派、运行时操作分类、响应关联、错误事件和显式或隐式关闭语义。

基于rust-v0.150.0
CodexRustSessionHandler

Session运行时处理 ​

Codex Core 对外接收的不是一组可以随意并发调用的 Session 方法,而是带 ID 和 trace 的 Submission。 每个 Thread 只有一个 submission loop,按 channel 接收顺序调用对应 handler。模型 Turn、Review、Compact 和部分 shell 工作会被 handler 启动为独立任务,但控制面上的配置变更、审批响应、rollback 与 shutdown 仍先经过这条串行入口。

理解 handlers 的重点不是背下 match Op 的每个分支,而是识别四种副作用:启动或控制 Task、修改 Session/持久化状态、唤醒当前 Turn 的 oneshot waiter,以及向独立的 Realtime/MCP/exec 子系统转发消息。 不同类别的完成信号和失败可见性并不相同。

阅读前先看 Core异步任务拓扑 理解单消费者控制面,再看 Session输入队列 区分 direct input、Turn pending 与 mailbox。本文按 Op 副作用分类,不重复每个 Task 的模型循环实现。

1. 分派入口 ​

SessionIo 将 Submission 写入 async channel,Session 构造完成后启动一个长期 submission_loop。Loop 每次只取一个 submission,并在继续 recv() 前 await 当前分支;因此 handler 本身的状态变更具有队列顺序, 但 handler 启动的后台 Task 可以与后续 submission 并发。

“串行”只约束 dispatch 阶段。例如 Compact handler 会 await spawn_task() 完成旧 Task 替换和新 Task 安装,然后 loop 可以继续接收审批或 interrupt;真正的 CompactTask 在 Tokio task 中运行。反过来, ThreadRollback 的 flush、load 和 replay 都在 handler 内 await,会推迟后续 Op 的分派。

源码位置:codex-rs/core/src/session/handlers.rs :: submission_loop骨架

rust
pub(super) async fn submission_loop(
    sess: Arc<Session>,
    config: Arc<Config>,
    rx_sub: Receiver<Submission>,
) {
    // To break out of this loop, send Op::Shutdown.
    let mut shutdown_received = false;
    while let Ok(sub) = rx_sub.recv().await {
        debug!(?sub, "Submission");
        let dispatch_span = submission_dispatch_span(&sub);
        let should_exit = async {
            match sub.op.clone() {
                // 每个handler在同一async分派体内await;分支返回bool只控制loop退出。
                Op::Interrupt => {
                    interrupt(&sess).await;
                    false
                }
                // 中间Op分支在后文按类别展开。
                Op::Shutdown => shutdown(&sess, sub.id.clone()).await,
                _ => false,
            }
        }
        .instrument(dispatch_span)
        .await;
        if should_exit {
            shutdown_received = true;
            break;
        }
    }
    // channel关闭也必须进入完整teardown,不能依赖显式Shutdown永远到达。
    if !shutdown_received {
        shutdown_session_runtime(&sess).await;
        emit_thread_stop_lifecycle(sess.as_ref()).await;
        if let Some(live_thread) = sess.live_thread()
            && let Err(err) = live_thread.shutdown().await
        {
            warn!("failed to shutdown thread persistence after submission channel closed: {err}");
        }
    }
}

2. Submission字段 ​

Submission 携带 submission ID、Op、W3C trace carrier、直接父 Turn ID 和因果 root Turn ID。 client message ID、thread settings、additional context 与 start options 已进入 TurnInputRequest,不再是 Submission 的平铺字段。parent/root lineage 仍用于 inter-agent communication;普通用户输入的 lineage 位于 TurnStartOptions。

调用方未提供 trace 时,submit_with_id() 捕获当前 span 的 W3C context。Loop 收到信封后再创建 op.dispatch.<kind> span,并把信封中的 carrier 设为 parent;这使 channel 排队后的异步 dispatch 仍能 接回原始请求 trace。

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

rust
pub(crate) async fn submit_with_id(&self, mut sub: Submission) -> CodexResult<()> {
    if sub.trace.is_none() {
        // 跨async channel前捕获调用方trace,避免接收端只能看到session loop上下文。
        sub.trace = current_span_w3c_trace_context();
    }
    self.tx_sub
        .send(sub)
        .await
        .map_err(|_| CodexErr::InternalAgentDied)?;
    Ok(())
}

pub(super) fn submission_dispatch_span(sub: &Submission) -> tracing::Span {
    let op_name = sub.op.kind();
    let span_name = format!("op.dispatch.{op_name}");
    let dispatch_span = match &sub.op {
        Op::RealtimeConversationAudio(_) => {
            // 高频audio frame降为debug span,其他Op使用info span。
            debug_span!(
                "submission_dispatch",
                otel.name = span_name.as_str(),
                submission.id = sub.id.as_str(),
                codex.op = op_name,
            )
        }
        _ => info_span!(
            "submission_dispatch",
            otel.name = span_name.as_str(),
            submission.id = sub.id.as_str(),
            codex.op = op_name,
        ),
    };
    if let Some(trace) = sub.trace.as_ref()
        && !set_parent_from_w3c_trace_context(&dispatch_span, trace)
    {
        // 非法carrier只被忽略并warning,不拒绝业务Op。
        warn!(submission.id = sub.id.as_str(), "ignoring invalid submission trace carrier");
    }
    dispatch_span
}

3. Op副作用分类 ​

当前 Op 标记为 #[non_exhaustive],Core 的 match 保留 _ => false,使较新的协议变体不会让旧 Core 崩溃。代价是未知 Op 被忽略且没有统一 ack;跨版本调用方不能把“成功写入 channel”解释为“handler 已执行”。

类别Ophandler 的主要动作
Session 控制Interrupt、CleanBackgroundTerminals、Shutdown中断当前 Task、结束后台终端或关闭整个 Session
Turn 输入与移交TurnInput、RecoverTurn、SuspendTurnAndShutdown返回 typed routing result,恢复旧 Turn,或持久化后移交 owner
Agent 输入InterAgentCommunication写入 mailbox,必要时触发 pending-work scheduler
RealtimeStart、Audio、Text、Speech、Close、ListVoices转发 RealtimeConversationManager 或直接返回内置 voices
设置与刷新ThreadSettings、RefreshMcpServers、ReloadUserConfig、SetThreadMemoryMode更新 Session 配置、标脏 runtime、重载缓存或更新 Thread metadata
Task 操作Compact、Review、RunUserShellCommand启动/替换 SessionTask,或在 active Turn 上运行辅助 shell
Turn waiter 响应ExecApproval、PatchApproval、UserInputAnswer、RequestPermissionsResponse、DynamicToolResponse、ResolveElicitation依据业务 ID 唤醒 oneshot;elicitation 可回退给 MCP runtime
历史重写与注入ThreadRollback、ApproveGuardianDeniedActionreplay rollout,或记录一次精确动作的 developer approval context

下面是分派表中与普通运行时最相关的主体。大部分分支返回 false;Shutdown 返回 true,成功的 SuspendTurnAndShutdown 也会退出 loop。

源码位置:codex-rs/core/src/session/handlers.rs :: submission_loop Op分派(节选)

rust
match sub.op {
    Op::TurnInput { request, mode, reply } => {
        let result = turn_input::handle(&sess, *request, mode, sub.id.clone()).await;
        let _ = reply.send(result);
        false
    }
    Op::RecoverTurn { thread_settings, reply } => {
        let result = turn_input::handle_recovery(
            &sess,
            thread_settings,
            sub.id.clone(),
        ).await;
        let _ = reply.send(result);
        false
    }
    Op::SuspendTurnAndShutdown { reply } => {
        let result = super::turn_suspension::suspend_turn_and_shutdown(
            &sess,
            sub.id.clone(),
        ).await;
        let should_exit = matches!(
            &result,
            Ok(SuspendTurnOutcome::Suspended { .. })
        );
        let _ = reply.send(result);
        should_exit
    }
    Op::ThreadSettings { thread_settings } => {
        thread_settings::update(&sess, sub.id.clone(), thread_settings).await;
        false
    }
    Op::InterAgentCommunication { communication } => {
        inter_agent_communication(
            &sess,
            sub.id.clone(),
            communication,
            sub.parent_turn_id,
            sub.root_turn_id,
        )
        .await;
        false
    }
    Op::ExecApproval {
        id: approval_id,
        turn_id,
        decision,
    } => {
        exec_approval(&sess, approval_id, turn_id, decision).await;
        false
    }
    Op::UserInputAnswer { id, response } => {
        request_user_input_response(&sess, id, response).await;
        false
    }
    Op::RefreshMcpServers => {
        refresh_mcp_servers(&sess);
        false
    }
    Op::Compact => {
        compact(&sess, sub.id.clone()).await;
        false
    }
    Op::ThreadRollback { num_turns } => {
        thread_rollback(&sess, sub.id.clone(), num_turns).await;
        false
    }
    Op::RunUserShellCommand { command } => {
        run_user_shell_command(&sess, sub.id.clone(), command).await;
        false
    }
    Op::Shutdown => shutdown(&sess, sub.id.clone()).await,
    Op::Review { review_request } => {
        review(&sess, &config, sub.id.clone(), review_request).await;
        false
    }
    // 其余Realtime、approval response、reload、memory和Guardian分支遵循同一结构。
    _ => false,
}

4. 完成信号 ​

绝大多数 handler 返回 (),自行选择完成信号。Submission loop 不会自动生成“Op completed” Event, 调用方必须按 Op 语义等待正确的通道。

主要完成模式如下:

  • TurnInput 与 RecoverTurn 使用 Op 内自带的 reply oneshot 返回 typed result;Started/Steered 之后的 Turn 生命周期仍走 Event channel。
  • SuspendTurnAndShutdown 返回 SuspendTurnOutcome,只有 Suspended 才让 submission loop 退出;错误时 当前 worker 继续拥有 Thread。
  • approval、request_user_input、request_permissions 和 dynamic tool response 通过当前 TurnState 中的 oneshot sender 唤醒原请求。
  • settings、rollback、Review 解析失败和 shutdown 使用 Event 表达结果。
  • MCP refresh、config reload、background terminal clean 等没有统一完成 Event;失败可能只写日志。
  • submit() 返回的 submission ID 是关联键,不是通用 completion future。

5. Turn响应关联 ​

审批请求在发 Event 前把 oneshot sender 注册到当前 TurnState。响应 Op 进入 loop 后,再用 approval ID、 call ID 或 (server_name, request_id) 从当前 ActiveTurn 中移除并发送。迟到响应不会被缓存给下一 Turn, 找不到 entry 时通常只 warning。

以 UserInputAnswer 为例,handler 自身只是薄转发,原子关联发生在 Session 方法内:

源码位置:codex-rs/core/src/session/handlers.rs :: request_user_input_response

rust
pub async fn request_user_input_response(
    sess: &Arc<Session>,
    id: String,
    response: RequestUserInputResponse,
) {
    // id必须与当前TurnState注册时使用的sub_id一致。
    sess.notify_user_input_response(&id, response).await;
}

pub async fn notify_user_input_response(
    &self,
    sub_id: &str,
    response: RequestUserInputResponse,
) {
    let entry = {
        let mut active = self.active_turn.lock().await;
        match active.as_mut() {
            Some(at) => {
                let mut ts = at.turn_state.lock().await;
                // remove保证同一个请求最多被成功响应一次。
                ts.remove_pending_user_input(sub_id)
            }
            None => None,
        }
    };
    match entry {
        Some(tx_response) => {
            tx_response.send(response).ok();
        }
        None => {
            // 迟到、重复或跨Turn响应不会污染后续请求。
            warn!("No pending user input found for sub_id: {sub_id}");
        }
    }
}

MCP elicitation 多一层兼容:先查当前 TurnState,未找到时交给 McpRuntime::resolve_elicitation(),因为 server-level elicitation 不一定由某个 active model Turn 发起。handler 还负责把协议层 string/integer ID 转换成 rmcp 的 NumberOrString,并把 Accept 缺少 content 的旧客户端请求补成空对象。

5.1 审批中断 ​

Exec/Patch approval 的 ReviewDecision::Abort 会调用 interrupt_task(),不是把 Abort 值发送给单个审批 waiter。其他 decision 才走 notify_approval()。Exec approval 若包含 execpolicy amendment,会先尝试持久化; 持久化失败发 Warning,但仍继续处理原 decision。

源码位置:codex-rs/core/src/session/handlers.rs :: exec_approval

rust
pub async fn exec_approval(
    sess: &Arc<Session>,
    approval_id: String,
    turn_id: Option<String>,
    decision: ReviewDecision,
) {
    let event_turn_id = turn_id.unwrap_or_else(|| approval_id.clone());
    if let ReviewDecision::ApprovedExecpolicyAmendment {
        proposed_execpolicy_amendment,
    } = &decision
        && let Err(err) = sess
            .persist_execpolicy_amendment(proposed_execpolicy_amendment)
            .await
    {
        // amendment写失败对客户端可见,但不吞掉原审批decision。
        let message = format!("Failed to apply execpolicy amendment: {err}");
        sess.send_event_raw(Event {
            id: event_turn_id.clone(),
            msg: EventMsg::Warning(WarningEvent { message }),
        })
        .await;
    }
    match decision {
        ReviewDecision::Abort => {
            // Abort终止整个当前Task,而不是只解除某个approval wait。
            sess.interrupt_task().await;
        }
        other => sess.notify_approval(&approval_id, other).await,
    }
}

6. Settings处理 ​

Op::ThreadSettings 与 Op::TurnInput.request.thread_settings 使用同一 submission queue,因此独立设置和 随后输入的调用顺序可被 Core 保留。standalone settings 由 session/thread_settings.rs 处理;TurnInput 在 Started 或 Steered 被接受后复用同一 prepare_update / apply_update 语义。model 和 reasoning effort 当前嵌在 CollaborationMode.settings 中;若请求没有给完整 collaboration mode,handler 会在锁内读取当前 mode,再只更新 model/effort。

源码位置:codex-rs/core/src/session/thread_settings.rs :: prepare_update

rust
let collaboration_mode = match collaboration_mode {
    Some(collaboration_mode) => collaboration_mode,
    None => {
        let state = sess.state.lock().await;
        // 部分更新继承当前mode与developer instructions,只替换显式model/effort。
        state
            .session_configuration
            .collaboration_mode
            .with_updates(model, effort, /*developer_instructions*/ None)
    }
};
SessionSettingsUpdate {
    environments,
    profile_workspace_roots,
    approval_policy,
    approvals_reviewer,
    sandbox_policy,
    permission_profile,
    active_permission_profile,
    windows_sandbox_level,
    collaboration_mode: Some(collaboration_mode),
    reasoning_summary: summary,
    service_tier,
    personality,
    ..Default::default()
}

应用成功后发送当前完整 snapshot,而不是回显调用方 patch;失败则发送 BadRequest Error。成功 Event 使用 send_event_raw_without_materializing_rollout():若 Thread 的本地 rollout 尚未真正创建,仅改设置不会为了 记录一个 Applied Event 强制物化文件;已有 rollout 时仍持久化事件。

源码位置:codex-rs/core/src/session/thread_settings.rs :: update, apply_update

rust
pub(super) async fn apply_update(
    session: &Session,
    submission_id: String,
    updates: SessionSettingsUpdate,
) -> ConstraintResult<()> {
    session.update_settings(updates).await?;
    emit_applied(session, submission_id).await;
    Ok(())
}

pub(super) async fn update(
    session: &Arc<Session>,
    submission_id: String,
    overrides: ThreadSettingsOverrides,
) {
    let updates = prepare_update(session, overrides).await;
    if let Err(error) = apply_update(session, submission_id.clone(), updates).await {
        session.send_event_raw(Event {
            id: submission_id,
            msg: EventMsg::Error(ErrorEvent {
                message: format!("invalid thread settings override: {error}"),
                codex_error_info: Some(CodexErrorInfo::BadRequest),
            }),
        }).await;
    }
}

RefreshMcpServers 则是非阻塞标记:先要求 runtime 在下一次 refresh 重连,再把 MCP state 标脏并调度 prewarm。ReloadUserConfig 会读取 user layer、清 Plugin/Skill cache 并刷新 runtime config,但读取、解析或 校验失败只 warning 后返回,当前 Op 没有专用 Error Event。

7. Task并发 ​

Compact、Review 和 user shell 都可能产生 TurnComplete,但 handler 的启动方式不同:

  • compact() 构造默认 TurnContext 后调用 spawn_task(CompactTask);spawn_task() 会先以 Replaced 中止 现有 Task。
  • review() 先同步解析 review target,再构造专用 review context 并通过 spawn_task(ReviewTask) 替换当前 Task;解析失败发 Error。
  • user shell 在已有 active Turn 时 tokio::spawn() 一个 auxiliary execution,共用当前 TurnContext 和 cancellation token;idle 时才启动独立 UserShellCommandTask。

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

rust
// :: compact、review与run_user_shell_command入口
pub async fn compact(sess: &Arc<Session>, sub_id: String) {
    let turn_context = sess.new_default_turn_with_sub_id(sub_id).await;
    // spawn_task内部先abort_all_tasks(Replaced),因此Compact是替换语义。
    sess.spawn_task(Arc::clone(&turn_context), Vec::new(), CompactTask)
        .await;
}

pub async fn run_user_shell_command(
    sess: &Arc<Session>,
    sub_id: String,
    command: String,
) {
    if let Some((turn_context, cancellation_token)) =
        sess.active_turn_context_and_cancellation_token().await
    {
        let session = Arc::clone(sess);
        // active路径不是第二个SessionTask,而是当前Turn的辅助异步工作。
        tokio::spawn(async move {
            execute_user_shell_command(
                session,
                turn_context,
                command,
                cancellation_token,
                UserShellCommandMode::ActiveTurnAuxiliary,
            )
            .await;
        });
        return;
    }

    let turn_context = sess.new_default_turn_with_sub_id(sub_id).await;
    // idle路径拥有独立Task生命周期和TurnComplete。
    sess.spawn_task(
        Arc::clone(&turn_context),
        Vec::new(),
        UserShellCommandTask::new(command),
    )
    .await;
}

因此“handler 已返回”只表示 Task 已安装或 auxiliary task 已调度,不代表模型、压缩或 shell 命令已经完成。 这些工作的最终状态仍由 Turn/exec Event 表达。

8. 历史回滚 ​

ThreadRollback 不直接从内存 Vec 尾部删除 N 条消息。它要求 num_turns >= 1、当前没有 active Turn、 Thread 已持久化;然后 flush writer、加载 durable history、追加 rollback marker 做一次 reconstruction,最后 把 marker 持久化并向客户端交付。

这里不是数据库式原子事务:replay 已更新内存后,marker flush 仍可能失败。实现选择继续交付 ThreadRolledBack,同时发 Warning 并依赖 writer 后续重试。前置 flush/load 失败则完全不 replay。

源码位置:codex-rs/core/src/session/handlers.rs :: thread_rollback(replay与提交节选)

rust
let rollback_event = ThreadRolledBackEvent { num_turns };
let rollback_msg = EventMsg::ThreadRolledBack(rollback_event.clone());
let replay_items = stored_history
    .items
    .into_iter()
    // reconstruction看到的是durable历史加上本次新marker。
    .chain(std::iter::once(RolloutItem::EventMsg(rollback_msg.clone())))
    .collect::<Vec<_>>();
sess.apply_rollout_reconstruction(
    turn_context.as_ref(),
    replay_items.as_slice(),
)
.await;
sess.services
    .agent_control
    .rollout_budget()
    .rearm_reminder(sess.thread_id());
sess.recompute_token_usage(turn_context.as_ref()).await;

// 先持久化marker;deliver_event_raw避免同一事件被send_event_raw再次追加。
sess.persist_rollout_items(&[RolloutItem::EventMsg(rollback_msg.clone())])
    .await;
if let Err(err) = sess.flush_rollout().await {
    sess.send_event(
        turn_context.as_ref(),
        EventMsg::Warning(WarningEvent {
            message: format!(
                "Rolled the thread back, but failed to save the rollback marker. Codex will continue retrying. Error: {err}"
            ),
        }),
    )
    .await;
}
sess.deliver_event_raw(Event {
    id: turn_context.sub_id.clone(),
    msg: rollback_msg,
})
.await;

9. Event发送 ​

Handlers 中看似相似的 send_event_raw()、send_event_raw_without_materializing_rollout() 和 deliver_event_raw() 有不同保证。错误使用会导致事件重复写入、在只改配置时意外创建 rollout,或绕过 trace/status 更新。

源码位置:codex-rs/core/src/session/mod.rs :: raw event发送路径

rust
pub(crate) async fn send_event_raw(&self, event: Event) {
    self.send_event_raw_with_persistence(event, /*persist*/ true)
        .await;
}

pub(crate) async fn send_event_raw_without_materializing_rollout(&self, event: Event) {
    let persist = match self.current_rollout_path().await {
        Ok(Some(path)) => codex_rollout::existing_rollout_path(&path).await.is_some(),
        Ok(None) => true,
        Err(err) => {
            warn!("failed to check whether thread persistence is materialized: {err}");
            true
        }
    };
    // 已有rollout仍持久化;只有尚未物化的本地路径跳过。
    self.send_event_raw_with_persistence(event, persist).await;
}

async fn send_event_raw_with_persistence(&self, event: Event, persist: bool) {
    // MCP resource/call authority也观察Core事件,用于更新调用来源状态。
    self.services.mcp_runtime.observe_event(&event.msg);
    if persist {
        self.persist_rollout_items(&[
            RolloutItem::EventMsg(event.msg.clone()),
        ])
        .await;
    }
    // 无论是否落盘,正常raw路径都记录rollout trace。
    self.services
        .rollout_thread_trace
        .record_protocol_event(&event.msg);
    self.deliver_event_raw(event).await;
}

async fn deliver_event_raw(&self, event: Event) {
    if let Some(status) = agent_status_from_event(&event.msg) {
        self.agent_status.send_replace(status);
    }
    // 客户端断开只debug记录,不能让Session业务流程因event receiver关闭而panic。
    if let Err(e) = self.tx_event.send(event).await {
        debug!("dropping event because channel is closed: {e}");
    }
}

Turn-scoped send_event() 还会在 CodexErrorInfo::affects_turn_status 为真时写入 TurnContext 的 terminal_error,并记录 turn/tool trace;handler 直接用 raw event 时没有这层 Turn 状态归约。选择哪条 API 取决于事件是否属于某个真实 Turn,而不是仅看有没有 sub_id。

10. 关闭处理 ​

显式 Op::Shutdown 会让 match 返回 true。SuspendTurnAndShutdown 在返回 SuspendTurnOutcome::Suspended 时也退出 loop,但它用于保留未完成 Turn ID 的 worker handoff;它在关闭 writer 和发送 thread-stop lifecycle 后直接交付 ShutdownComplete,不记录普通 Turn 的 aborted/complete 终态。普通 Shutdown 执行 runtime teardown、统计真实用户 Turn 数、运行 thread-stop extension、关闭持久化,然后直接交付 ShutdownComplete 并结束 rollout trace。由于 LiveThread 已 shutdown,完成事件不再经普通持久化路径。

channel 被所有 sender 关闭时,while let Ok 自然退出。隐式路径仍执行 runtime teardown、thread-stop 和 LiveThread shutdown,但没有可用请求方等待的 ShutdownComplete,也不会走显式 handler 的会话计数逻辑。

两条路径共享的 shutdown_session_runtime() 顺序是:abort 未消费 prewarm、关闭 realtime、interrupt active Task、终止 exec/code mode、停止 MCP prewarm、关闭 MCP runtime、关闭 Guardian、运行 session-end hooks。 thread-stop extension 被放在共享 runtime teardown 之后。

源码位置:codex-rs/core/src/session/handlers.rs :: shutdown尾部

rust
let event = Event {
    id: sub_id,
    msg: EventMsg::ShutdownComplete,
};
// LiveThread已关闭,完成事件只记trace并直接deliver,不再写rollout。
sess.services
    .rollout_thread_trace
    .record_protocol_event(&event.msg);
sess.deliver_event_raw(event).await;
sess.services
    .rollout_thread_trace
    .record_ended(codex_rollout_trace::RolloutStatus::Completed);
true

11. 错误返回 ​

Handlers 的错误策略取决于是否存在可关联的客户端操作和是否还能安全继续:

场景可见结果Session 是否继续
TurnInput settings/steer 无效reply oneshot 返回 Error 或 NotSubmittedReason是
RecoverTurn 遇到 active TurnNotSubmitted(NotIdle)是,不应用 settings
SuspendTurnAndShutdown 拒绝reply 返回 NotActive/HasLiveDescendants/UnsupportedTask 或 Error是,当前 worker保留owner
Standalone ThreadSettings 无效BadRequest Error Event是
Review request 无法解析Other Error Event是
Rollback 前置条件/读取失败ThreadRollbackFailed Error Event是,历史不变
Rollback marker flush 失败Warning,rollback 已交付是,writer 后续重试
Execpolicy amendment 写失败Warning,仍处理审批 decision是
找不到 Turn waiterwarning log是
Reload user config 读取/解析失败warning log是,保留旧配置
Realtime start 失败Other Error Event是
未知未来 Op无事件,忽略是
显式 Shutdown persistence 关闭失败Error 后仍发 ShutdownComplete否

这意味着调用方不能只等待“同 ID 的任意 Event”。有些 Op 通过 reply oneshot 完成,有些是 fire-and-forget,有些即使发 Warning 也已经部分生效。技术上最危险的是把 enqueue 成功、handler 返回和 业务完成混为一个时刻。

12. 分派测试 ​

第一组测试验证 trace 不是被 session loop 截断:submit_with_id() 会捕获当前 span,dispatch span 优先 采用 submission carrier;Realtime audio 使用 Debug level,避免高频帧污染 Info trace。

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

rust
// :: submission_dispatch_span_prefers_submission_trace_context(断言节选)
let dispatch_span = ambient_span.in_scope(|| {
    submission_dispatch_span(&Submission {
        id: "sub-1".into(),
        op: Op::Interrupt,
        trace: Some(submission_trace),
        parent_turn_id: None,
        root_turn_id: None,
        // 显式submission trace应覆盖接收端ambient span作为parent。
    })
});
let trace_id = dispatch_span.context().span().span_context().trace_id();
assert_eq!(
    trace_id,
    TraceId::from_hex("00000000000000000000000000000055")
        .expect("trace id"),
);

第二组测试验证 submission sender 异常消失也不会跳过清理,并锁定 turn-abort 必须先于 thread-stop extension:

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

rust
// :: submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle(节选)
session
    .spawn_task(
        Arc::new(turn_context),
        Vec::new(),
        NeverEndingTask {
            kind: TaskKind::Regular,
            listen_to_cancellation_token: true,
        },
    )
    .await;

let (tx_sub, rx_sub) = async_channel::bounded(1);
drop(tx_sub);
// 没有发送Op::Shutdown,recv因channel关闭退出。
submission_loop(
    Arc::clone(&session),
    session.get_config().await,
    rx_sub,
)
.await;

assert_eq!(
    vec!["turn_abort", "thread_stop"],
    *calls.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
);

排查运行时 Op 时,建议先确定它属于哪类完成路径:

  • Op 已成功 submit 但没有结果:确认该 Op 是否本来就是 fire-and-forget,以及是否应该等待 admission、 Turn Event、exec Event 或 oneshot waiter。
  • 后续 Op 长时间不分派:检查前一个 inline handler 是否在做 rollback I/O、Review setup 或配置 reload; submission loop 本身不会并发 await 两个 handler。
  • approval 响应“找不到请求”:检查当前 ActiveTurn、业务 ID 和是否重复响应;waiter 不跨 Turn 缓存。
  • 设置 Event 与预期 patch 不同:ThreadSettingsApplied 返回的是更新后的完整 snapshot。
  • rollback 已生效但保存告警:这是 replay 成功、marker flush 失败的部分成功路径,不应自动重复 rollback。
  • channel 关闭后未收到 ShutdownComplete:隐式 teardown 不承诺该 Event,应等待 session-loop termination。
  • trace 断链:检查 Submission 是否携带合法 W3C carrier,而不是只看接收端 ambient span。
  • 新协议 Op 被旧 Core 忽略:Op 是 non-exhaustive,wildcard 分支不生成通用错误回执。

Session handlers 的职责是把一个有序控制流翻译成不同生命周期对象的动作。单消费者 loop 提供配置与控制 顺序,Task 和子系统提供并发执行,业务 ID/oneshot 提供精确响应关联,而 Event API 决定结果是否持久化。 只有同时看这四层,才能准确判断一个 Op 到底是“已入队”“已分派”“已启动”还是“真正完成”。

这些测试不证明所有 handler 都具有相同的完成语义,也不证明异步调度中的每一种交错都能由当前 fixture 复现;它们只锁定本文列出的 dispatch、rollback、waiter 和 teardown 顺序。

显式 Shutdown、channel close 与资源释放顺序见 Session关闭流程;字段锁与 waiter owner 见 Session核心数据结构。

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

bash
rg -n "Submission|dispatch|rollback|waiter|record_response" codex-rs/core/src