Skip to content

CodexThread公共API

按调用契约梳理 CodexThread 的提交、事件、Turn 控制、配置、持久化、MCP 与关闭接口。

基于rust-v0.150.0
CodexRustCodexThreadAPI

CodexThread公共API ​

CodexThread 是已加载 Thread 的进程内操作面。它把 Arc<Session> 与 SessionIo 放在一个 handle 中, 让宿主不必直接取得 Session lock、submission receiver 或后台 task handle。

这个 API 面并不等同于 App Server wire protocol,也不是第三方稳定 ABI。部分方法只服务仓库内部产品, 还有方法带 #[doc(hidden)] 或明确计划移除。正确使用它的前提,是先分清“提交成功”“Turn 已准入” “事件已消费”“状态快照”和“持久化完成”五种不同保证。

本文面向已理解 Thread/Session/Turn 层级,并读过 ThreadManager管理 的读者。本文是 API 契约参考:帮助宿主选择正确 方法并理解返回保证;它不展开 run_turn、MCP transport或ThreadStore内部算法。

读完后,应能为“发送普通Op、等待用户消息准入、向active Turn注入、只写history、读取状态、直接调用 MCP、等待关闭”分别选择API,并知道哪些返回值只是入队确认、哪些是业务完成确认。

1. Handle ​

源码位置:codex-rs/core/src/codex_thread.rs :: CodexThread, OutOfBandElicitations

rust
pub struct CodexThread {
    // Session 拥有运行状态和 thread-scoped services;handle 只共享 Arc。
    pub(crate) session: Arc<Session>,
    // SessionIo 保存 submission/event/status/termination 四类端点。
    pub(crate) io: SessionIo,
    pub(crate) session_source: SessionSource,
    session_configured: SessionConfiguredEvent,
    rollout_path: Option<PathBuf>,
    out_of_band_elicitations: Mutex<OutOfBandElicitations>,
    // 线程存活期间维持 live-thread 诊断计数。
    _diagnostics_guard: GaugeGuard,
}

#[derive(Default)]
struct OutOfBandElicitations {
    count: i64,
    registration: Option<ElicitationRegistration>,
}

session_configured 和 rollout_path 是创建时保留的快照;agent_status()、config_snapshot()、token usage 与 environment selections 则查询当前 Session。调用方不能看到一个 getter 就默认它是实时值。

2. API地图 ​

下面的 mindmap 按调用意图而不是源码出现顺序整理方法。它也解释了为什么本文不采用逐行罗列的 TOC。

这几组方法共享同一个 handle,却采用不同一致性模型:submission 用有界队列,event 用无界队列,status 用 watch,persistence 经 LiveThread,MCP 直接访问当前 runtime。

3. Turn输入路由 ​

方法接受内容返回时保证适合调用方
submit(Op)任意 Core Op已进入 submission queue,返回 submission ID普通控制操作;Interrupt、Shutdown 等
submit_with_traceOp + W3C trace同上,并传播显式 traceApp Server/宿主 trace 桥接
start_or_steer_turnTurnInputRequestCore 已返回 Started、Steered 或 NotSubmittedApp Server turn/start 与普通用户输入
start_turn_if_idle任意 TurnInputRequest仅在 idle 时启动,否则给出稳定拒绝原因Queue、Goal、Agent 等自动唤醒
steer_turn用户输入 + expected Turn ID只注入指定 active regular TurnApp Server turn/steer
recover_turn_if_idle中断 Turn ID + thread settingsidle 时用原 Turn ID 恢复 samplingworker handoff 与恢复
suspend_turn_and_shutdown无显式输入停止并落盘未完成 root Turn,再关闭 writerroot Turn 所有权转移

CodexThread 没有 interrupt() 方法。中断通过 submit(Op::Interrupt) 进入同一 Session control plane; 审批响应、配置更新和 shutdown 同样是不同 Op 变体。

源码位置:codex-rs/core/src/codex_thread.rs :: submit, submit_with_trace, start_or_steer_turn, start_turn_if_idle, steer_turn

rust
pub async fn submit(&self, op: Op) -> CodexResult<String> {
    self.io.submit(op).await
}

pub async fn submit_with_trace(
    &self,
    op: Op,
    trace: Option<W3cTraceContext>,
) -> CodexResult<String> {
    self.io
        .submit_with_trace(
            op, trace, /*parent_turn_id*/ None, /*root_turn_id*/ None,
        )
        .await
}

pub async fn start_or_steer_turn(
    &self,
    request: TurnInputRequest,
) -> CodexResult<TurnInputSubmission> {
    self.submit_turn_input_with_mode(request, TurnInputMode::StartOrSteer)
        .await
}

pub async fn start_turn_if_idle(
    &self,
    request: TurnInputRequest,
) -> CodexResult<StartIfIdleSubmission> {
    match self
        .submit_turn_input_with_mode(request, TurnInputMode::StartIfIdle)
        .await?
    {
        TurnInputSubmission::Started { turn_id } => {
            Ok(StartIfIdleSubmission::Started { turn_id })
        }
        TurnInputSubmission::NotSubmitted { reason } => {
            Ok(StartIfIdleSubmission::NotSubmitted { reason })
        }
        TurnInputSubmission::Steered { .. } => unreachable!(),
    }
}

前两个方法不会检查 Op 是否为 UserInput,也不会等待 handler 采用它。queue sender 成功只说明 Session loop 将看到 submission;Turn input API 则会等待 Core 的路由决定,但仍不等待 prompt Hook、history 写入、rollout 持久化或模型 sampling 完成。

3.1 请求模型 ​

TurnInputRequest 把输入、持久 thread settings、仅在新 Turn 使用的 start options、附加上下文、 Responses metadata 和 trace 放进一个对象。调用方通过方法选择路由语义,而不是自己操作独立的 准入等待表:

源码位置:codex-rs/protocol/src/turn_input.rs :: TurnInputRequest, TurnInputMode, TurnInputSubmission, NotSubmittedReason

rust
pub struct TurnInputRequest {
    pub input: TurnInput,
    pub thread_settings: ThreadSettingsOverrides,
    pub start: TurnStartOptions,
    pub additional_context: BTreeMap<String, AdditionalContextEntry>,
    pub responsesapi_client_metadata: Option<HashMap<String, String>>,
    pub trace: Option<W3cTraceContext>,
}

pub enum TurnInputMode {
    StartOrSteer,
    StartIfIdle,
    Steer { expected_turn_id: String },
}

pub enum TurnInputSubmission {
    Started { turn_id: String },
    Steered { turn_id: String },
    NotSubmitted { reason: NotSubmittedReason },
}

Started 和 Steered 都只表示 Core 已接受路由。NotSubmittedReason 能区分 NotIdle、 PendingTriggerTurn、PlanMode、NoActiveTurn、expected Turn 不匹配、Review/Compact 不可 steer、 输出 schema 不匹配和空输入。App Server 将这些 reason 映射成具体 JSON-RPC 错误,而不是再等待一个 独立准入等待表。

4. ID关系 ​

四种 ID 的关系适合用实体图而不是再画调用流程:

start_or_steer_turn 启动新 Turn 时,submission ID 成为 Turn ID;steer 成功时返回的是当前 active Turn ID。控制类 Op 仍有 submission ID,却不产生 Turn。父 Turn、因果 root Turn 与结构化输出 schema 位于 TurnStartOptions,只在 Started 分支写入;Steered 分支不会重写 active Turn 的 lineage。 tool call ID 则用于配对模型调用与工具输出,不能拿 Event.id 代替。

5. 读取侧 ​

源码位置:codex-rs/core/src/codex_thread.rs :: next_event, agent_status, token_usage_info

rust
pub async fn next_event(&self) -> CodexResult<Event> {
    // 每次调用从同一个 Receiver 取走一个 Event;多个并发消费者会竞争分流。
    self.io.next_event().await
}

pub async fn agent_status(&self) -> AgentStatus {
    // watch 只返回最新状态,不重放中间变化。
    self.io.agent_status().await
}

pub(crate) fn subscribe_status(&self) -> watch::Receiver<AgentStatus> {
    self.io.agent_status.clone()
}

pub async fn token_usage_info(&self) -> Option<TokenUsageInfo> {
    self.session.token_usage_info().await
}

next_event() 不是 broadcast subscription。App Server 为一个 loaded Thread 建立单一 listener,再把事件 投影和 fan-out 到客户端;两个业务模块若直接并发调用 next_event(),每个只会拿到其中一部分事件。

agent_status() 是 cheap latest snapshot,适合列表和状态判断;需要观察变化的内部调用方使用 subscribe_status()。token usage 也是 Session 缓存的完整快照,不应从最后一个 delta event 自行拼接。

6. Turn控制 ​

Turn input API 已把“启动、steer、恢复和未准入”变成 typed result;inject_if_running 与 history-only 注入仍是不同的低层桥接:

APIActiveTurn要求无可用 Turn 时是否新建 Turn失败输入处理
start_or_steer_turn可有可无idle 时启动可能NotSubmittedReason
start_turn_if_idle必须没有 active TurnNotIdle 等 reason是不记录、不入队
steer_turn指定 ID 的 regular TurnNoActiveTurn 或 mismatch否不记录、不入队
recover_turn_if_idle必须 idleNotIdle恢复旧 Turn ID不应用 settings
inject_if_running任意 active TurnErr(original_items)否原样返还 items
inject_response_items不要求写 history/context否空列表报 InvalidRequest

start_turn_if_idle() 供 Queue、Goal、Agent 等扩展启动 idle work。它在持有 active-turn reservation 前后都检查 trigger-turn mailbox,避免自动工作抢在用户/agent 触发输入之前;没有显式用户输入时,Plan 模式也会得到 PlanMode。steer_turn() 则只接受匹配 expected Turn ID 的 regular Turn,Review 与 Compact 返回 ActiveTurnNotSteerable。

源码位置:codex-rs/core/src/codex_thread.rs :: start_turn_if_idle, steer_turn, recover_turn_if_idle, inject_if_running

rust
pub async fn inject_if_running(
    &self,
    items: Vec<ResponseItem>,
) -> Result<(), Vec<ResponseItem>> {
    self.session.inject_if_running(items).await
}

pub async fn steer_turn(
    &self,
    request: TurnInputRequest,
    expected_turn_id: String,
) -> CodexResult<SteerSubmission> {
    match self
        .submit_turn_input_with_mode(request, TurnInputMode::Steer { expected_turn_id })
        .await?
    {
        TurnInputSubmission::Steered { turn_id } => {
            Ok(SteerSubmission::Steered { turn_id })
        }
        TurnInputSubmission::NotSubmitted { reason } => {
            Ok(SteerSubmission::NotSubmitted { reason })
        }
        TurnInputSubmission::Started { .. } => unreachable!(),
    }
}

inject_response_items() 是 history-only API:输入不得为空;没有 reference context 时会先捕获一次 step context,再记录 items 并 flush rollout。它不会产生新的用户 Turn boundary,不能代替正常 UserInput。 inject_response_items_for_turn() 则用于调用方持有 thread-operation lock 时,把 raw response items 紧邻后续用户输入写入,不单独 flush rollout。

7. 配置API刷新 ​

类别方法语义
创建快照session_configured()、rollout_path()创建/恢复时固定,不随每轮变化
当前只读config_snapshot()、config()、environment_selections()从 Session 当前状态读取
可恢复设置thread_settings_snapshot()、restorable_thread_settings()捕获 thread-owned 设置供持久化或 runtime replacement
来源诊断instruction_sources()、legacy_instruction_sources()返回实际加载指令的文件来源
约束预览preview_thread_settings_overrides()计算候选 snapshot,不提交 mutation
恢复设置restore_thread_settings()resume 后恢复当前 loaded runtime 的 thread-owned mutable settings
定向刷新refresh_runtime_config()刷新 layer-backed runtime config,保留 session-static 设置
MCP 刷新refresh_mcp_config()刷新 MCP 与 managed requirements,不替换无关配置

源码位置:codex-rs/core/src/codex_thread.rs :: preview_thread_settings_overrides, config_snapshot, refresh_runtime_config

rust
pub async fn preview_thread_settings_overrides(
    &self,
    overrides: CodexThreadSettingsOverrides,
) -> ConstraintResult<ThreadConfigSnapshot> {
    let updates = self.thread_settings_update(overrides).await;
    // preview 只运行约束合并,不写 SessionState。
    self.session.preview_settings(&updates).await
}

pub async fn config_snapshot(&self) -> ThreadConfigSnapshot {
    self.session.thread_config_snapshot().await
}

pub async fn refresh_runtime_config(&self, next_config: crate::config::Config) {
    // 该方法不是“替换整个 Session Config”,静态 feature/identity 等仍保持原值。
    self.session.refresh_runtime_config(next_config).await;
}

调用方要验证一次 Turn override 是否被 managed constraint 接受,应先 preview;要永久改变设置则走相应 Session Op/产品 API。直接读取 config() 后修改 clone 不会自动提交回 Session。

8. 持久化 API ​

源码位置:codex-rs/core/src/codex_thread.rs :: load_history, read_thread, update_thread_metadata, append_rollout_items

rust
pub async fn load_history(
    &self,
    include_archived: bool,
) -> ThreadStoreResult<StoredThreadHistory> {
    let live_thread = self
        .session
        .live_thread_for_persistence("load history")
        .map_err(|err| ThreadStoreError::Internal {
            message: err.to_string(),
        })?;
    live_thread.load_history(include_archived).await
}

pub async fn update_thread_metadata(
    &self,
    patch: ThreadMetadataPatch,
    include_archived: bool,
) -> ThreadStoreResult<StoredThread> {
    let live_thread = self
        .session
        .live_thread_for_persistence("update thread metadata")
        .map_err(|err| ThreadStoreError::Internal {
            message: err.to_string(),
        })?;
    // metadata 经 LiveThread 排序,避免绕过正在追加的 rollout writer。
    live_thread.update_metadata(patch, include_archived).await
}

pub async fn append_rollout_items(&self, items: &[RolloutItem]) -> ThreadStoreResult<()> {
    let live_thread = self
        .session
        .live_thread_for_persistence("append rollout items")
        .map_err(|err| ThreadStoreError::Internal {
            message: err.to_string(),
        })?;
    live_thread.append_items(items).await
}

ephemeral Thread 没有 LiveThread,这些方法不能承诺 durable 行为。ensure_rollout_materialized() 与 flush_rollout() 虽是 public,但标记 doc(hidden);它们主要服务测试、fork/resume 和产品生命周期, 不应被当成普通业务写入 API。

9. 宿主直连 ​

read_mcp_resource() 和 call_mcp_tool() 会先按 dirty state 刷新 MCP runtime,再直接调用最新 connection。它们不是模型产生 ToolCall 后经过 ToolRouter/ToolOrchestrator 的同一入口,宿主必须自己 承担参数、授权和错误投影。

当前 API 还增加了两个更窄的服务入口:read_mcp_resource_for_call(call_id, uri) 会沿原始 tool call authority 读取 app resource,避免调用方只凭 server 名重新选择身份;start_mcp_event_stream() 打开 MCP event stream,并要求 _meta 必须是 JSON object。它们仍是宿主直连,不自动继承模型工具审批链。

list_background_terminals() 返回 item ID、process ID、command 与 cwd; terminate_background_terminal(process_id) 直接委托 Session process manager。返回 bool 表示是否找到 并终止目标,不提供完整退出状态。

10. Elicitation计数 ​

Code Mode 等宿主能力可在模型工具结果回灌期间打开带外 elicitation。计数从 0 变为 1 时注册 pause, 降回 0 时释放;嵌套 elicitation 共享同一 registration。

源码位置:codex-rs/core/src/codex_thread.rs :: increment_out_of_band_elicitation_count, decrement_out_of_band_elicitation_count

rust
pub async fn increment_out_of_band_elicitation_count(&self) -> CodexResult<i64> {
    let mut elicitations = self.out_of_band_elicitations.lock().await;
    let incremented = elicitations.count.checked_add(1).ok_or_else(|| {
        CodexErr::Fatal("out-of-band elicitation count overflowed".to_string())
    })?;
    if elicitations.count == 0 {
        // 只在 0→1 时注册一次 pause,嵌套调用不重复注册。
        elicitations.registration = Some(self.session.services.elicitations.register());
    }
    elicitations.count = incremented;
    Ok(incremented)
}

pub async fn decrement_out_of_band_elicitation_count(&self) -> CodexResult<i64> {
    let mut elicitations = self.out_of_band_elicitations.lock().await;
    if elicitations.count == 0 {
        // underflow 是调用协议错误,不能静默保持零。
        return Err(CodexErr::InvalidRequest(
            "out-of-band elicitation count is already zero".to_string(),
        ));
    }
    elicitations.count -= 1;
    if elicitations.count == 0 {
        elicitations.registration = None;
    }
    Ok(elicitations.count)
}

测试验证 count 为 2 时释放一层仍保持 pause,只有第二次 decrement 到 0 后 captured exec result 才能 回到模型。计数不是 UI 装饰字段,而是资源注册的引用计数。

11. 关闭与移交 ​

shutdown_and_wait() 提交 Shutdown,并等待 Session loop termination;wait_until_terminated() 只等待 已经发生的终止,不发送控制消息。

源码位置:codex-rs/core/src/codex_thread.rs :: shutdown_and_wait, wait_until_terminated

rust
pub async fn shutdown_and_wait(&self) -> CodexResult<()> {
    self.io.shutdown_and_wait().await
}

pub async fn wait_until_terminated(&self) {
    // clone shared future 允许多个 owner 等同一次 loop 终止。
    self.io.session_loop_termination.clone().await;
}

若宿主只观察另一个 owner 发起的关闭,用 wait_until_terminated();若当前 owner 负责触发关闭,用 shutdown_and_wait()。两者都不从 ThreadManager map 移除 handle,也不删除持久化 Thread。

suspend_turn_and_shutdown() 的语义更强:它只允许 owning root thread 使用,拒绝仍有 loaded descendant 或不支持的 task;成功时停止当前未完成 Turn,但不记录 TurnAborted/TurnComplete,随后 flush history、 关闭 writer,并等待 Session loop 终止。这样另一个 worker 才能用 recover_turn_if_idle() 在原 Turn ID 上继续 sampling。调用方在 Suspended 返回前不能转移所有权,因为 accepted suspension 即使调用方断连 也会由 Session 继续完成。

源码位置:codex-rs/core/src/codex_thread.rs :: suspend_turn_and_shutdown, recover_turn_if_idle

rust
pub async fn recover_turn_if_idle(
    &self,
    request: RecoverTurnRequest,
) -> CodexResult<StartIfIdleSubmission> {
    self.session
        .services
        .agent_control
        .ensure_execution_capacity_for_turn_start(self)
        .await?;
    let RecoverTurnRequest { turn_id, thread_settings, trace } = request;
    match self.io.submit_recover_turn(thread_settings, trace, turn_id).await? {
        TurnInputSubmission::Started { turn_id } => {
            Ok(StartIfIdleSubmission::Started { turn_id })
        }
        TurnInputSubmission::NotSubmitted { reason } => {
            Ok(StartIfIdleSubmission::NotSubmitted { reason })
        }
        TurnInputSubmission::Steered { .. } => unreachable!(),
    }
}

12. 调用与误用 ​

目标推荐 API常见误用
发普通控制 Opsubmit / submit_with_trace认为返回 ID 表示 handler 已成功
确认输入已启动或 steerstart_or_steer_turn轮询 agent_status 猜测路由结果
单一事件消费循环next_event多组件并发 recv,误以为每个都收到全量事件
读取最新轻量状态agent_status把 watch snapshot 当事件历史
向当前 Turn 注入且不创建新 Turninject_if_running无 active 时丢弃返回的原 items
idle extension 自动开 Turnstart_turn_if_idle忽略 PendingTriggerTurn/PlanMode/NotIdle reason
只写历史、不形成用户 Turninject_response_items用它模拟正常 UserInput
变更 loaded metadataupdate_thread_metadata直接写 store 绕过 LiveThread ordering
直接 MCP 宿主调用read_mcp_resource / call_mcp_tool假设自动经过模型工具审批链
触发并等待关闭shutdown_and_wait只等 termination,却从未发送 Shutdown

审阅调用点时,最重要的问题不是“这个方法是否 public”,而是“它返回时究竟确认了哪一层”。 CodexThread 的价值正是把底层 owner 收口成窄方法;错误地扩大返回保证,反而会重新引入 Session 并发 和持久化竞态。

13. 返回值验证 ​

13.1 路由拒绝 ​

turn_input_tests 直接验证三个 API 最容易被混淆的拒绝边界:start-only 遇到 active Turn 返回 NotIdle 且不注入输入;自动空输入在 Plan mode 返回 PlanMode;steer-only 没有 active Turn 时返回 NoActiveTurn。这些结果都不需要等待模型 sampling。

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

rust
let submission = submit_start_only(&session, synthetic_input).await;
assert_eq!(
    TurnInputSubmission::NotSubmitted {
        reason: NotSubmittedReason::NotIdle,
    },
    submission
);
assert_eq!(
    (Vec::<TurnInput>::new(), None, None),
    session.input_queue.get_pending_input(&session.active_turn).await
);

let steer = submit_steer_only(&idle_session, input, "missing-turn-id").await;
assert_eq!(
    TurnInputSubmission::NotSubmitted {
        reason: NotSubmittedReason::NoActiveTurn,
    },
    steer
);

相邻测试还覆盖 expected Turn ID mismatch、Review/Compact 不可 steer、Plan mode 和 trigger-turn mailbox。 关键不变量是 NotSubmitted 不应用 thread settings、start options,也不把输入写入 pending queue。

13.2 Active边界 ​

pending-input集成测试让第一响应已经产生final answer,但继续用reasoning item保持stream打开;此时调用 inject_if_running(),释放response gate后等待TurnComplete。最终server必须收到第二次请求,且其中真实用户 文本顺序为初始prompt后跟injected context。

源码位置:codex-rs/core/tests/suite/pending_input.rs

rust
// :: injected_response_item_reopens_turn_after_final_answer(核心断言节选)
assert!(
    codex
        .inject_if_running(vec![responses::user_message_item(INJECTED_CONTEXT)])
        .await
        .is_ok()
);
let _ = gate_completed_tx.send(());
wait_for_turn_complete(&codex).await;

let requests = server.requests().await;
assert_eq!(requests.len(), 2);
let second: Value = from_slice(&requests[1]).expect("parse second request");
let relevant_user_input = message_input_texts(&second, "user")
    .into_iter()
    .filter(|text| text == INITIAL_PROMPT || text == INJECTED_CONTEXT)
    .collect::<Vec<_>>();
// 注入不抢占正在解析的item,而是在follow-up request中按history顺序出现。
assert_eq!(
    relevant_user_input,
    vec![INITIAL_PROMPT, INJECTED_CONTEXT],
);

该测试证明active路径,不证明 idle语义。inject_response_items() 有 agent-control测试覆盖持久化调用和history mode,但没有直接断言“idle时未创建Turn、items已进入模型history、空输入被拒绝”三项完整契约;这些仍 主要由 CodexThread::inject_response_items 源码保证。

13.3 Elicitation计数 ​

Code Mode测试先把counter加到2,再让第一个exec result等待。decrement到1后第二个模型请求必须仍为空; 只有归零后结果才允许回灌并产生下一请求。

源码位置:codex-rs/core/tests/suite/code_mode.rs

rust
// :: code_mode_exec_holds_captured_result_during_elicitation(核心断言)
assert_eq!(
    test.codex.increment_out_of_band_elicitation_count().await?,
    1,
);
assert_eq!(
    test.codex.increment_out_of_band_elicitation_count().await?,
    2,
);
let release_elicitation = async {
    tokio::time::timeout(Duration::from_secs(5), async {
        while first_mock.requests().is_empty() {
            tokio::time::sleep(Duration::from_millis(10)).await;
        }
    })
    .await
    .expect("initial response request should arrive");
    tokio::time::sleep(Duration::from_millis(250)).await;
    // count大于0时,captured exec result不能进入模型follow-up。
    assert!(
        second_mock.requests().is_empty(),
        "captured exec result should not return during an elicitation"
    );
    assert_eq!(
        test.codex.decrement_out_of_band_elicitation_count().await?,
        1,
    );
    tokio::time::sleep(Duration::from_millis(100)).await;
    assert!(
        second_mock.requests().is_empty(),
        "captured exec result should wait for every elicitation"
    );
    assert_eq!(
        test.codex.decrement_out_of_band_elicitation_count().await?,
        0,
    );
    Ok::<(), anyhow::Error>(())
};

它验证嵌套释放行为,不覆盖overflow/underflow;这两个错误分支由 checked arithmetic源码定义。

13.4 多个关闭者 ​

测试构造一个只接收一条Shutdown后延迟结束的session loop,同时启动两个 shutdown_and_wait()。两位waiter都必须成功完成;第二次submit即使遇到receiver关闭,也不能导致它跳过 termination等待。

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

rust
// :: shutdown_and_wait_allows_multiple_waiters(断言节选)
let waiter_1 = {
    let io = Arc::clone(&io);
    tokio::spawn(async move { io.shutdown_and_wait().await })
};
let waiter_2 = {
    let io = Arc::clone(&io);
    tokio::spawn(async move { io.shutdown_and_wait().await })
};

// 两个调用者等待同一个Shared<BoxFuture>,不竞争消费一次性JoinHandle。
waiter_1
    .await
    .expect("first shutdown waiter join")
    .expect("first shutdown waiter");
waiter_2
    .await
    .expect("second shutdown waiter join")
    .expect("second shutdown waiter");

这证明多waiter,不证明 ShutdownComplete 会被每个消费者广播;next_event() 仍是单receiver队列语义。

14. API契约 ​

  1. 对 submit()、start_or_steer_turn() 和 shutdown_and_wait() 分别说明“返回时保证到哪一层”。
  2. 对比 steer_turn、inject_if_running、start_turn_if_idle 和 inject_response_items:无active Turn时 各自返回什么、是否创建Turn、是否保留原输入。
  3. 解释为什么 agent_status() 不能替代Event消费,为什么 wait_until_terminated() 不能触发关闭。
  4. 使用只读命令定位四组API和测试:
bash
rg -n "start_or_steer_turn|start_turn_if_idle|steer_turn|inject_if_running|shutdown_and_wait" \
  codex-rs/core/src/codex_thread.rs codex-rs/core/src/session/mod.rs
rg -n "start_only_rejects|steer_only_|injected_response_item|captured_result_during|multiple_waiters" \
  codex-rs/core/src/session codex-rs/core/tests

继续学习 CodexThread背压,理解本文的 next_event() 为什么不 提供fan-out;需要深入某个 Op如何执行时,回看 Session运行时处理。