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(节选)
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
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(节选)
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(节选)
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
// ... 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
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(节选)
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(节选)
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
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 policy | layer rule解析失败 | 否 | 映射 Fatal,停止构造 |
| Model | catalog/provider解析异常 | 否 | fallback仅在显式允许时生效 |
| Thread persistence | create/resume writer失败 | 可能部分打开 | guard创建前错误直接返回;store负责局部清理 |
| Environment | packaged zsh不可用、selection无效 | guard持有 | Session transaction返回错误并discard |
| AGENTS刷新 | environment filesystem读取失败 | guard持有 | 返回错误并discard live writer |
| Network proxy | 受管网络无法启动 | guard持有 | fail closed,不发布Session |
| Hook | startup warning | guard持有 | warning排在SessionConfigured之后 |
| Extension callback | callback自身处理结果 | guard持有 | registry contributor决定内部容错 |
| MCP install | projection/connection失败 | 首事件可能已排队 | transaction失败,manager不注册handle |
| Referenced fork | child 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
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
// :: 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
// :: 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. 构造屏障验证
可以先用只读搜索把事务的三个提交点定位出来:
# 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完成搜索后,尝试回答:
- 为什么 auth/MCP projection 可以和 persistence 并行,而 environment-dependent skill warmup 不可以?
- 为什么
SessionConfigured已写入内部 event queue,仍不代表 Thread 已经对客户端发布? - required MCP 与 optional MCP 的启动失败为什么不能采用相同返回策略?
commit()为什么只清空 guard 的 owner,而discard()必须调用 ThreadStore?
构造提交后的 prewarm 调度见 Session启动预热;资源释放的逆向过程见 Session关闭流程。
