Skip to content

Session依赖装配

展开 SessionSpawnArgs 到 SessionServices 的依赖解析、并行初始化、发布屏障和失败回滚。

基于rust-v0.150.0
CodexRustSessionDependency

Session依赖装配 ​

Session::spawn() 参数很多,不是因为缺少 builder,而是因为它位于多个生命周期的汇合点:ThreadManager 共享服务、调用方 Thread 选项、持久化历史、当前认证、模型目录、环境 registry 和宿主 extension 都要 在此收敛成一个 Thread 级运行时。

理解构造逻辑的关键不是记住参数顺序,而是分清三类依赖:构造前已经确定的输入、构造中产生的 thread-owned 资源,以及 SessionConfigured 之后才允许启动的后台能力。错置任一类,都会破坏首事件 顺序、持久化清理或 MCP 生命周期。

阅读本文前,应先理解 ThreadManager创建 中 manager registry 的提交位置,以及 Session核心数据结构 中 SessionState、 SessionServices 和 LiveThread 的所有权差异。本文把 Session::spawn() 看成一笔异步事务:先准备资源, 再建立客户端可见事件,最后提交持久化 owner;任何中途错误都必须沿 Result 返回并回收尚未提交的资源。

本文从 SessionSpawnArgs 开始,到 Session::spawn() 返回 (Arc<Session>, SessionIo) 为止。首事件读取、 HashMap vacant entry和manager registry发布属于RUN007,不在本文重复。

读完后,应能回答三个问题:一个依赖为什么属于某轮 tokio::join!,何时提交 LiveThreadInitGuard, 以及 shell runtime、AGENTS 刷新或 required MCP 失败后为什么不能留下可继续写入的 LiveThread。

1. 依赖入口 ​

下面的装配图展示五类输入如何汇入 SessionSpawnArgs,以及该对象随后分流到策略解析与资源初始化。

ThreadManager 负责把这些来源归一化为一个 request;Session 构造不再回到全局单例查找依赖。显式传参 使 subagent、测试和不同产品宿主能够替换 auth、store、extension 或 environment owner。

源码位置:codex-rs/core/src/session/mod.rs :: SessionSpawnArgs(节选)

rust
pub(crate) struct SessionSpawnArgs {
    pub(crate) config: Config,
    pub(crate) allow_provider_model_fallback: bool,
    pub(crate) user_instructions: LoadedUserInstructions,
    pub(crate) installation_id: String,
    // 以下manager级Arc跨Thread复用,Session只clone自己的强引用。
    pub(crate) auth_manager: Arc<AuthManager>,
    pub(crate) models_manager: SharedModelsManager,
    pub(crate) environment_manager: Arc<EnvironmentManager>,
    pub(crate) skills_service: Arc<HostSkillsService>,
    pub(crate) plugins_manager: Arc<PluginsManager>,
    pub(crate) mcp_manager: Arc<McpManager>,
    pub(crate) code_mode_session_provider: Arc<dyn CodeModeSessionProvider>,
    pub(crate) extensions: Arc<ExtensionRegistry<Config>>,
    pub(crate) conversation_history: InitialHistory,
    pub(crate) requested_history_mode: Option<ThreadHistoryMode>,
    pub(crate) fork_persistence: ForkPersistence,
    pub(crate) session_source: SessionSource,
    pub(crate) forked_from_thread_id: Option<ThreadId>,
    pub(crate) parent_thread_id: Option<ThreadId>,
    pub(crate) thread_source: Option<ThreadSource>,
    pub(crate) originator: String,
    pub(crate) agent_control: AgentControl,
    pub(crate) dynamic_tools: Vec<DynamicToolSpec>,
    pub(crate) inherited_exec_policy: Option<Arc<ExecPolicyManager>>,
    pub(crate) inherited_environments: Option<TurnEnvironmentSnapshot>,
    pub(crate) parent_rollout_thread_trace: ThreadTraceContext,
    pub(crate) thread_extension_init: ExtensionDataInit,
    pub(crate) client_mcp_extensions: ClientMcpExtensions,
    pub(crate) thread_store: Arc<dyn ThreadStore>,
    pub(crate) attestation_provider: Option<Arc<dyn AttestationProvider>>,
    pub(crate) external_time_provider: Option<Arc<dyn TimeProvider>>,
    pub(crate) inherited_multi_agent_version: Option<MultiAgentVersion>,
    // ... metrics、trace、shell、environment与平台策略字段。
}

SessionSpawnArgs 是一次性消费对象。spawn_internal() 解构后,各依赖分别进入 SessionConfiguration、SessionServices、持久化 future 或后台 worker,不会原样保存在 Session 中。

2. 策略与模型 ​

构造 Session 前必须先得到稳定的 SessionConfiguration。这一步包含 exec policy、模型刷新、模型 fallback、instructions、history mode、dynamic tools、service tier 和权限环境快照。

2.1 Exec策略来源 ​

源码位置:codex-rs/core/src/session/mod.rs :: Session::spawn_internal exec policy selection

rust
let exec_policy = if crate::guardian::is_basic_session_source(&session_source) {
    // Basic guardian只采用requirements中的managed policy,忽略用户/项目可塑造规则。
    let managed_policy = config
        .config_layer_stack
        .requirements()
        .exec_policy
        .as_deref()
        .map_or_else(codex_execpolicy::Policy::empty, |policy| {
            policy.as_ref().clone()
        });
    Arc::new(ExecPolicyManager::new(Arc::new(managed_policy)))
} else if let Some(exec_policy) = &inherited_exec_policy {
    Arc::clone(exec_policy)
} else {
    if !config
        .config_layer_stack
        .ignore_user_and_project_exec_policy_rules()
    {
        let codex_home = config.codex_home.clone();
        let policy_path = default_policy_path(codex_home.as_path());
        if let Err(err) = prefix_rule_migration(
            codex_home.as_path(),
            policy_path.as_path(),
            BANNED_PREFIX_SUGGESTIONS,
        )
        .await
        {
            tracing::warn!(error = %err, "failed to run prefix rule migration");
        }
    }
    Arc::new(
        // 普通root从resolved layer stack加载;解析失败属于Session启动失败。
        ExecPolicyManager::load(&config.config_layer_stack)
            .await
            .map_err(|err| CodexErr::Fatal(format!("failed to load rules: {err}")))?,
    )
};

Basic guardian 只接收 requirements 里的 managed policy;普通 child 可以继承 parent policy;普通 root 才 合并用户/项目规则。policy migration 失败只 warning,最终 policy load 失败则 fatal。二者的容错级别 不同:迁移是兼容辅助,没有可执行 policy 才破坏启动不变量。

2.2 模型与指令 ​

root 使用 OnlineIfUncached,non-root agent 使用 Offline refresh strategy,避免每个 child 重复联网刷新 catalog。只有调用方显式允许 provider fallback,requested model 不可用时才替换并记录日志。

源码位置:codex-rs/core/src/session/mod.rs :: model, history and SessionConfiguration resolution(节选)

rust
let refresh_strategy = if session_source.is_non_root_agent() {
    codex_models_manager::manager::RefreshStrategy::Offline
} else {
    codex_models_manager::manager::RefreshStrategy::OnlineIfUncached
};
let model = models_manager
    .get_default_model(
        &config.model,
        allow_provider_model_fallback,
        refresh_strategy,
        config.http_client_factory(),
    )
    .await;
let model_info = models_manager
    .get_model_info(model.as_str(), &config.to_models_manager_config())
    .await;
let history_mode = conversation_history.get_history_mode(
    requested_history_mode.unwrap_or_else(|| thread_store.default_history_mode()),
);
let base_instructions = config
    .base_instructions
    .clone()
    .or_else(|| conversation_history.get_base_instructions().map(|s| s.text))
    .unwrap_or_else(|| model_info.get_model_instructions(config.personality));
let dynamic_tools = if dynamic_tools.is_empty() {
    // resume/fork在调用方未覆盖时恢复持久化dynamic tools。
    conversation_history.get_dynamic_tools().unwrap_or_default()
} else {
    dynamic_tools
};

let session_configuration = SessionConfiguration {
    provider: create_model_provider(
        config.model_provider.clone(),
        Some(Arc::clone(&auth_manager)),
    ),
    collaboration_mode,
    base_instructions,
    approval_policy: config.permissions.approval_policy.clone(),
    permission_profile_state: session_permission_profile_state_from_config(&config)?,
    environments: TurnEnvironmentSelections::new(config.cwd.clone(), environment_selections),
    history_mode,
    forked_from_thread_id,
    parent_thread_id,
    dynamic_tools,
    // ... 其余reasoning、identity、originator与platform字段。
};

到这里才具备继续执行 Session::spawn() 的逻辑快照。后续环境和 service 装配读取该 configuration,不再分别 从原 Config、history 和 manager options 重算优先级。

3. 第一轮并行 ​

Session identity 与 extension init 确定后,三条互不依赖的 future 同时启动:

源码位置:codex-rs/core/src/session/session.rs :: first initialization join(节选)

rust
let thread_persistence_fut = async {
    // persistence future只拥有创建/恢复职责,不等待State DB或MCP projection。
    if config.ephemeral {
        Ok::<_, anyhow::Error>(None)
    } else {
        let live_thread = match &initial_history {
            InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => {
                // ... CreateThreadParams在此按identity、lineage、history mode与metadata组装。
                if is_paginated_subagent
                    && matches!(&fork_persistence, ForkPersistence::Copied)
                    && let InitialHistory::Forked(items) = &initial_history
                {
                    LiveThread::create_with_inherited_model_context(
                        Arc::clone(&thread_store),
                        params,
                        items,
                    )
                    .await?
                } else {
                    LiveThread::create(Arc::clone(&thread_store), params).await?
                }
            }
            InitialHistory::Resumed(resumed_history) => {
                let params = ResumeThreadParams {
                    thread_id: resumed_history.conversation_id,
                    rollout_path: resumed_history.rollout_path.clone(),
                    history: Some(resumed_history.history.clone()),
                    include_archived: true,
                    metadata: ThreadPersistenceMetadata {
                        cwd: Some(config.cwd.to_path_buf()),
                        model_provider: config.model_provider_id.clone(),
                        memory_mode: if config.memories.generate_memories {
                            ThreadMemoryMode::Enabled
                        } else {
                            ThreadMemoryMode::Disabled
                        },
                    },
                };
                LiveThread::resume(
                    Arc::clone(&thread_store),
                    session_configuration.history_mode,
                    params,
                )
                .await?
            }
        };
        Ok(Some(live_thread))
    }
};

let state_db_fut = async {
    if config.ephemeral {
        None
    } else if let Some(local_store) = thread_store.as_any().downcast_ref::<LocalThreadStore>() {
        local_store.state_db().await
    } else {
        None
    }
};

let auth_and_mcp_fut = async move {
    let auth = auth_manager_clone.auth().await;
    let mcp_projection = mcp_manager_for_mcp
        .runtime_config_for_step(
            &config_for_mcp,
            mcp_thread_init_for_startup,
            thread_extension_data_for_mcp,
            McpThreadIdentity {
                session_source: &mcp_session_source,
                originator: &mcp_originator,
            },
            /*ready_selected_capability_roots*/ &[],
            /*executor_capability_discovery*/ None,
        )
        .await;
    (auth, mcp_projection)
};

let (thread_persistence_result, state_db_ctx, (auth, mcp_projection)) =
    tokio::join!(thread_persistence_fut, state_db_fut, auth_and_mcp_fut);
let mut live_thread_init =
    LiveThreadInitGuard::new(thread_persistence_result.map_err(|e| {
        error!("failed to initialize thread persistence: {e:#}");
        e
    })?);

上面是按真实分支压缩的摘录:CreateThreadParams、ResumeThreadParams 和 MCP 参数已在前文/相邻文章完整 展开。关键并发边界是三条 future 只共享不可变 snapshot/Arc,没有互相依赖的输出。

4. 第二轮并行依赖 ​

默认 shell、ShellSnapshot 和 ThreadEnvironments 必须先完成,因为 AGENTS、Plugin/Skill 都需要在正确 executor filesystem 上读取。环境 snapshot 形成后,三条工作再次并行:

  • AgentsMdManager::refresh();
  • Plugin 与 Skill warmup;
  • 从 ThreadStore 查询现有 thread name。

源码位置:codex-rs/core/src/session/session.rs :: environment and plugin-skill initialization

rust
// ... default_shell与shell_snapshot已按feature和override解析。
let turn_environments = Arc::new(ThreadEnvironments::new(
    environment_manager,
    default_shell.clone(),
    // Environment使用SessionConfiguration推导后的环境配置。
    session_configuration.inferred_environment_config(),
    shell_snapshot,
    inherited_environments.unwrap_or_default(),
    config.features.enabled(Feature::DeferredExecutor),
));
turn_environments.update_selections(
    environment_selections,
    &session_configuration.inferred_environment_config(),
);
let resolved_environments = turn_environments.snapshot().await;

let (agents_md_result, plugin_skill_errors, thread_name) = tokio::join!(
    agents_md_manager.refresh(config.as_ref(), &resolved_environments),
    plugin_skill_warmup,
    thread_name_lookup,
);
agents_md_result?;
// Skill错误记录日志后继续;AGENTS刷新错误直接终止事务。
session_configuration.thread_name = thread_name.clone();

ShellZshFork 被启用却没有可用 packaged zsh 时不能悄悄回退默认 shell,因为 feature 已承诺具体执行路径。 普通 Skill 加载错误则记录 error 日志并允许 Session 继续;当前这条路径不会自动生成客户端 warning。 当前版本没有旧版 config_lock.rs 的独立 Session lock 校验/导出步骤;thread settings 的约束与更新由 CodexThread::preview_thread_settings_overrides、restore_thread_settings 和 Turn input request 处理。

5. SessionServices ​

5.1 Managed失败 ​

SessionState 建立后,Core 根据 network requirements 构造 approval service、blocked-request observer 与 Weak Session decider。真正的 Session Arc 此时尚不存在,所以 decider 先保存空 Weak,Session 创建后 再回填。

proxy start 会读取当前 exec policy 和 permission profile。若配置要求受管网络却无法启动 proxy,构造 直接失败;不能创建一个看似受管、实际可直连的 Session。

5.2 Hook回调 ​

build_hooks_for_config() 依赖 resolved local environment 和 plugin manager。Hook startup warnings 进入 post_session_configured_events,保证客户端先拿到 canonical Session snapshot,再接收 warning。

5.3 Extension ​

thread lifecycle contributor 可能在 on_thread_start 中读取 MCP resources,但正式 connection set 只有 在 SessionConfigured 后才能发布。Core 先创建空 McpRuntime 和稳定的 McpResourceClient,执行 extension callback,稍后再安装 projection。

源码位置:codex-rs/core/src/session/session.rs :: empty MCP runtime and extension startup

rust
let mcp_runtime = Arc::new(McpRuntime::empty(
    mcp_projection.config.prefix_mcp_tool_names,
));
let session_extension_data = ExtensionData::new(session_id.to_string());
let mcp_resource_client = Arc::new(McpResourceClient::new(Arc::clone(&mcp_runtime)));
for contributor in extensions.thread_lifecycle_contributors() {
    // callback获得稳定resource client,但不会提前看到未发布的正式connections。
    contributor
        .on_thread_start(ThreadStartInput {
            config: config.as_ref(),
            session_source: &session_configuration.session_source,
            persistent_thread_state_available: state_db_ctx.is_some(),
            environments: session_configuration.environment_selections(),
            mcp_resource_client: Some(Arc::clone(&mcp_resource_client)),
            session_store: &session_extension_data,
            thread_store: &thread_extension_data,
            // ... metrics字段。
        })
        .await;
}

这个顺序解决了循环依赖:Extension 需要 thread-owned MCP handle,正式 MCP 初始化又依赖完整 Session identity、environment 和 auth。

6. 服务装配 ​

下面的 ER 图把 configuration、state、services 与主要 thread-owned resource 的一对一/可选关系放在 一起。它描述生命周期基数,不表示关系型数据库。

LiveThread 与 State DB 在 ephemeral/remote store 场景可以缺失;model client、MCP runtime 和 exec manager 则是每个正常 SessionServices 的明确字段。network proxy 使用 ArcSwapOption,允许运行中重建替换。

源码位置:codex-rs/core/src/session/session.rs :: SessionServices assembly(节选)

rust
let services = SessionServices {
    // 初始为空connections,正式projection在SessionConfigured后安装。
    mcp_runtime,
    mcp_handler_cache: Default::default(),
    unified_exec_manager: UnifiedExecProcessManager::new(
        config.background_terminal_max_timeout,
    ),
    elicitations: ElicitationService::new(),
    analytics_events_client,
    hooks: ArcSwap::from_pointee(hooks),
    rollout_thread_trace,
    user_shell: Arc::new(default_shell),
    exec_policy,
    auth_manager: Arc::clone(&auth_manager),
    models_manager: Arc::clone(&models_manager),
    tool_approvals: Mutex::new(ApprovalStore::default()),
    guardian_rejection_circuit_breaker: Mutex::new(Default::default()),
    runtime_handle: tokio::runtime::Handle::current(),
    skills_service,
    agents_md_manager,
    plugins_manager: Arc::clone(&plugins_manager),
    mcp_manager: Arc::clone(&mcp_manager),
    extensions,
    session_extension_data,
    thread_extension_data,
    selected_capability_roots,
    mcp_thread_init,
    client_mcp_extensions,
    agent_control,
    network_proxy: ArcSwapOption::from(network_proxy.map(Arc::new)),
    state_db: state_db_ctx.clone(),
    live_thread: live_thread_init.as_ref().cloned(),
    thread_store: Arc::clone(&thread_store),
    // ... ModelClient按auth、provider、identity、features与http factory完整构造。
    code_mode_service: CodeModeService::new(
        Arc::clone(&code_mode_session_provider),
        &config.code_mode,
    ),
    tool_search_handler_cache: Default::default(),
    turn_environments: Arc::clone(&turn_environments),
    // ... audit、capability、time与工具缓存字段。
};

这里没有 lazy get_or_create service map。一个必要服务构造失败,Session 不会以部分可用状态发布。

7. 发布后收尾 ​

Arc<Session> 构造完成后,Core 仍未完成事务。它必须按顺序发送首事件、安装 MCP、启动 worker、记录 历史并提交 LiveThread guard。

源码位置:codex-rs/core/src/session/session.rs :: publication barrier and commit(节选)

rust
let events = std::iter::once(Event {
    id: INITIAL_SUBMIT_ID.to_owned(),
    msg: EventMsg::SessionConfigured(SessionConfiguredEvent {
        session_id,
        thread_id,
        cwd: thread_config.cwd().clone(),
        forked_from_id: thread_config.forked_from_thread_id,
        parent_thread_id: thread_config.parent_thread_id,
        thread_source: thread_config.thread_source,
        thread_name: session_configuration.thread_name.clone(),
        model: thread_config.model,
        permission_profile: thread_config.permission_profile,
        active_permission_profile: thread_config.active_permission_profile,
        rollout_path,
        // ... client bootstrap字段。
    }),
})
.chain(post_session_configured_events.into_iter());
for event in events {
    // SessionConfigured严格排在warning和MCP startup event之前。
    sess.send_event_raw(event).await;
}
turn_environments.start_connection_event_forwarding(tx_event.clone());
sess.install_initial_mcp_runtime(
    &session_configuration,
    latest_auth,
    mcp_projection,
    &resolved_environments,
    mcp_runtime_cwd,
)
.await?;
sess.start_mcp_prewarm_worker(mcp_prewarm_rx, mcp_auth_changes);
sess.schedule_startup_prewarm(session_configuration.base_instructions.clone())
    .await;
Box::pin(sess.record_initial_history(initial_history)).await;

match session_result {
    Ok(sess) => {
        // 直到所有fallible装配完成,LiveThread才脱离失败清理guard。
        live_thread_init.commit();
        Ok(sess)
    }
    Err(err) => {
        live_thread_init.discard().await;
        Err(err)
    }
}

SessionConfigured 在内部 queue 排队后仍可能因 MCP 安装或 referenced fork materialization 失败而回滚。 此时 Session::spawn() 不会返回 SessionIo,调用方拿不到该 receiver;manager 也不会注册 CodexThread,因此内部排队过首事件不等于客户端观察到已发布 Session。

8. Submission Loop ​

spawn_internal() 完成 Session 装配后,才启动 submission loop,并把 JoinHandle 转成 shared termination future。这样第一条 Op 不会与 SessionServices、history 或 MCP 安装竞争初始化。

源码位置:codex-rs/core/src/session/mod.rs :: submission loop final assembly

rust
let thread_id = session.thread_id;
let session_for_loop = Arc::clone(&session);
let session_loop_handle = tokio::spawn(async move {
    submission_loop(session_for_loop, configured_config, rx_sub)
        .instrument(info_span!("session_loop", thread_id = %thread_id))
        .await;
});
let io = SessionIo {
    tx_sub,
    rx_event,
    agent_status: agent_status_rx,
    // 多个CodexThread owner可等待同一个loop终止,不复制后台task。
    session_loop_termination: session_loop_termination_from_handle(session_loop_handle),
};
Ok((session, io))

这一步之后 manager 才能读取已排队的 SessionConfigured 并发布 CodexThread。Session::spawn、submission loop 与 manager registry 是三个连续提交点,不是一个构造函数的同义表达。

9. 装配失败矩阵 ​

阶段典型失败是否已有持久化清理/结果
Exec policylayer rule解析失败否映射 Fatal,停止构造
Modelcatalog/provider解析异常否fallback仅在显式允许时生效
Thread persistencecreate/resume writer失败可能部分打开guard创建前错误直接返回;store负责局部清理
Environmentpackaged zsh不可用、selection无效guard持有Session transaction返回错误并discard
AGENTS刷新environment filesystem读取失败guard持有返回错误并discard live writer
Network proxy受管网络无法启动guard持有fail closed,不发布Session
Hookstartup warningguard持有warning排在SessionConfigured之后
Extension callbackcallback自身处理结果guard持有registry contributor决定内部容错
MCP installprojection/connection失败首事件可能已排队transaction失败,manager不注册handle
Referenced forkchild reference无法materialize首事件可能已排队reservation和child persistence清理

设计新依赖时应明确它属于哪个阶段:若客户端在首事件前必须知道结果,应放在 barrier 之前;若它会 产生 startup event,则应在 SessionConfigured 之后启动;若失败后必须释放资源,则必须进入 LiveThreadInitGuard 覆盖的 fallible transaction,而不是在返回 Session 后裸 spawn。

10. 失败路径测试 ​

没有一组测试可以单独覆盖整个构造事务。Session 单测覆盖 packaged zsh 缺失的构造失败,App Server 集成测试覆盖 required MCP 错误向调用方传播,ThreadStore 测试覆盖 discard 撤销 writer ownership。Manager registry屏障测试已归入RUN007。

10.1 Shell依赖失败 ​

ShellZshFork feature 被启用时,Session 构造不能静默退回普通 shell。测试启用 feature 但不提供 packaged zsh,直接调用 Session::new(),断言启动失败并保留明确原因。

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

rust
let result = Session::new(
    session_configuration,
    /*environment_selections*/ &[],
    Arc::clone(&config),
    // 其余依赖由测试正常构造,此处省略。
).await;

let err = match result {
    Ok(_) => panic!("expected startup to fail"),
    Err(err) => err,
};
let msg = format!("{err:#}");
assert!(msg.contains(
    "zsh fork feature enabled, but no packaged zsh fork is available"
));

测试直接证明缺少必要 shell runtime 时 Session::new() 不会返回部分可用 Session。它没有覆盖 manager registry;外层 spawn_internal() 的 LiveThread discard 仍由构造控制流和 ThreadStore 测试共同证明。

10.2 Required MCP失败 ​

MCP server 分为 optional 与 required。install_initial_mcp_runtime() 先把 connection set 安装进稳定 runtime, 再等待 validate_required_servers();只有 required server 的启动失败才会让 Session 构造返回 Err。 App Server 集成测试使用不存在的命令创建 required server,断言 thread/start 收到 JSON-RPC error。

源码位置:codex-rs/app-server/tests/suite/v2/thread_start.rs

rust
// :: thread_start_fails_when_required_mcp_server_fails_to_initialize(关键断言)
let req_id = mcp
    .send_thread_start_request_with_auto_env(ThreadStartParams::default())
    .await?;

let err: JSONRPCError = timeout(
    DEFAULT_READ_TIMEOUT,
    mcp.read_stream_until_error_message(RequestId::Integer(req_id)),
)
.await??;

assert!(
    err.error
        .message
        .contains("required MCP servers failed to initialize"),
    "unexpected error message: {}",
    err.error.message
);
// 聚合错误保留具体失败server,便于用户定位配置项。
assert!(
    err.error.message.contains("required_broken"),
    "unexpected error message: {}",
    err.error.message
);

这说明错误穿过 McpConnectionSet → McpRuntime → Session → ThreadManager → App Server 到达调用方。测试没有 返回 ThreadStartResponse,因此客户端拿不到半初始化 Thread;optional MCP 的失败则通过 startup status event 报告,不适用这条 fatal 语义。

10.3 Discard ​

LiveThreadInitGuard::discard() 最终调用 ThreadStore::discard_thread()。LocalThreadStore 的测试覆盖 Legacy 和 Paginated 两种 history mode:创建后 rollout 尚未物化,但 writer lock 已存在;discard 后两者都不存在, 再 append 必须得到 ThreadNotFound。

源码位置:codex-rs/thread-store/src/local/mod.rs

rust
// :: discard_thread_drops_unmaterialized_live_writer(核心路径)
for history_mode in [ThreadHistoryMode::Legacy, ThreadHistoryMode::Paginated] {
    let thread_id = ThreadId::default();
    let mut params = create_thread_params(thread_id);
    params.history_mode = history_mode;

    store.create_thread(params).await.expect("create live thread");
    let rollout_path = store
        .live_rollout_path(thread_id)
        .await
        .expect("load rollout path");
    assert!(!rollout_path.exists());

    let lock_path = home
        .path()
        .join("thread-writer-locks")
        .join(format!("{thread_id}.lock"));
    assert!(lock_path.exists());

    // 回滚未提交Session:释放writer lock,但不强制生成空rollout。
    store.discard_thread(thread_id).await.expect("discard live thread");
    assert!(!rollout_path.exists());
    assert!(!lock_path.exists());

    let err = store
        .append_items(AppendThreadItemsParams {
            thread_id,
            items: vec![user_message_item("write after discard")],
        })
        .await
        .expect_err("discard should remove the live thread writer");
    assert!(matches!(
        err,
        ThreadStoreError::ThreadNotFound { thread_id: missing } if missing == thread_id
    ));
}

如果 rollback 只删除 rollout 文件而保留 lock,下一次启动会被误判为重复 live writer;如果只释放 lock 却保留可写 recorder,迟到的后台写入又可能复活失败 Session。测试同时封住这两个方向。

10.4 尚缺的故障注入 ​

当前测试分别覆盖 shell/required MCP 的失败判定、错误传播、registry 屏障和 discard 语义,但没有把 Session::spawn() 的这些失败注入与 InMemoryThreadStore::discard_thread、manager registry 为空组合成一条 端到端断言。因此可以依据构造函数的统一错误分支复述控制流,却不能把它解释成一个测试已经覆盖所有回滚动作; 这仍是适合补充的故障注入场景。

11. 构造屏障验证 ​

可以先用只读搜索把事务的三个提交点定位出来:

bash
# 1. 哪些初始化任务并行,哪些必须等待上一轮输出?
rg -n "tokio::join!|thread_persistence_fut|plugin_skill_warmup" \
  codex-rs/core/src/session/session.rs

# 2. Session何时提交LiveThread,错误时从哪里统一discard?
rg -n "live_thread_init\.(commit|discard)" codex-rs/core/src/session/session.rs

# 3. manager何时才把CodexThread写入registry?
rg -n "finalize_thread_spawn|threads\.entry" codex-rs/core/src/thread_manager.rs

# 4. required MCP失败如何从connection set进入Session Result?
rg -n "validate_required_servers" codex-rs/codex-mcp/src codex-rs/core/src/session

完成搜索后,尝试回答:

  1. 为什么 auth/MCP projection 可以和 persistence 并行,而 environment-dependent skill warmup 不可以?
  2. 为什么 SessionConfigured 已写入内部 event queue,仍不代表 Thread 已经对客户端发布?
  3. required MCP 与 optional MCP 的启动失败为什么不能采用相同返回策略?
  4. commit() 为什么只清空 guard 的 owner,而 discard() 必须调用 ThreadStore?

构造提交后的 prewarm 调度见 Session启动预热;资源释放的逆向过程见 Session关闭流程。