ThreadManager创建
ThreadManager::start_thread() 看起来只接收 StartThreadOptions 并返回 NewThread,实际需要完成一条 跨越配置、模型、持久化、Extension、MCP、环境和异步消息的启动事务。只有当 Session 已经构造成功, 并且第一个事件严格是 SessionConfigured 时,manager 才把 Arc<CodexThread> 注册进 active map。
这条边界很重要:ThreadId 分配不代表 thread 已对外可用,rollout 创建不代表 Session 初始化已提交, Session 对象存在也不代表 manager registry 已发布。创建流程用首事件握手和 LiveThreadInitGuard 把这些阶段连接起来。
本文面向已经读过 ThreadManager依赖 的读者,默认理解 Thread、Session和Turn的层级。本文只回答 fresh Thread 怎样从 StartThreadOptions 变成已注册 CodexThread;Resume和Fork的历史选择分别留给 ThreadManager恢复、ThreadManager分叉,Session内部每项依赖的细节留给 Session依赖装配。
读完后,应能从 start_thread() 复述到 finalize_thread_spawn(),说明 SessionConfigured 为什么必须 是首事件、HashMap vacant entry为什么是最终原子发布点,并知道 Session内部资源提交应转入RUN014。
1. start_thread参数
ThreadManager 不读取 config.toml。调用方在进入 start_thread() 之前已经完成 config layer、profile、 requirements、feature 和路径解析,StartThreadOptions 持有完整 Config:
源码位置:codex-rs/core/src/thread_manager.rs :: StartThreadOptions
pub struct StartThreadOptions {
// Config 已在调用方完成分层与约束解析,manager 不重新读取配置文件。
pub config: Config,
pub allow_provider_model_fallback: bool,
pub initial_history: InitialHistory,
// source、环境与扩展输入允许宿主覆盖;None 才触发 manager 默认值。
pub history_mode: Option<ThreadHistoryMode>,
pub session_source: Option<SessionSource>,
pub thread_source: Option<ThreadSource>,
pub dynamic_tools: Vec<DynamicToolSpec>,
pub metrics_service_name: Option<String>,
pub parent_trace: Option<W3cTraceContext>,
pub environments: Option<Vec<TurnEnvironmentSelection>>,
pub thread_extension_init: ExtensionDataInit,
pub client_mcp_extensions: ClientMcpExtensions,
pub reserved_thread_id: Option<ThreadId>,
}构造器把缺省状态收敛为 fresh root 语义,避免不同宿主各自拼一套默认值:
源码位置:codex-rs/core/src/thread_manager.rs :: StartThreadOptions::new
pub fn new(config: Config) -> Self {
Self {
config,
allow_provider_model_fallback: false,
// New 明确要求分配 fresh 身份,而不是恢复或复制历史。
initial_history: InitialHistory::New,
history_mode: None,
session_source: None,
thread_source: None,
dynamic_tools: Vec::new(),
metrics_service_name: None,
parent_trace: None,
// 空环境留给 manager 根据 cwd/workspace roots 统一派生。
environments: None,
thread_extension_init: ExtensionDataInit::default(),
client_mcp_extensions: ClientMcpExtensions::default(),
}
}StartThreadOptions::new(config) 的默认值专门表示 fresh root thread:
| 字段 | 默认值 | 新 Thread 语义 |
|---|---|---|
allow_provider_model_fallback | false | 请求模型不可用时不静默替换 |
initial_history | InitialHistory::New | 分配新 ThreadId,无历史重建 |
history_mode | None | 使用 ThreadStore 默认值 |
| session/thread source | None | 从 manager 默认 source 推导 |
dynamic_tools | 空 | 没有调用方动态工具 |
| metrics/trace | None | 使用默认 originator/新 trace |
environments | None | 从 cwd/workspace roots 生成默认 selection |
| extension/MCP init | 空 | 只使用宿主 registry 与 config |
App Server 可以显式填 environment、dynamic tools、parent trace 或 extension init;这些是同一个创建 API 的增强输入,不是另一条 Session 构造实现。
2. Manager
图中 Session::spawn 被视为一个下游事务节点。Manager 只有在它整体成功后才拿到 SessionIo,再从 queue 读取首事件;内部依赖装配和资源 rollback 不在本篇展开。
3. 入口归一化
start_thread() 只把调用委托给 start_thread_inner()。后者完成三件事:
- 根据 Config 创建带 rollout budget 的
AgentControl; - 从 initial history 解析 resumed source;fresh history 没有时使用 manager 默认
session_source; - 生成内部
ThreadSpawnRequest,注入 manager 的 AuthManager,并记录可选 fork source。
ThreadSpawnRequest 比 public options 多出 lineage 与继承字段:parent_thread_id、 forked_from_thread_id、ForkPersistence、inherited environments、inherited exec policy 和测试 shell override。当前版本的 public options 还允许宿主先提供 reserved_thread_id;它只能用于 fresh start, resume 时会被明确拒绝。fresh root 未指定时由 manager 的 ThreadIdGenerator 生成,persistence 默认为 Copied。
所有 new/resume/fork/subagent 入口最终都转换成该内部 request,保证 spawn_thread() 只有一条 Session 构造路径。后续文章会分别展开 history 和 lineage 差异;本文只跟踪 InitialHistory::New。
源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManager::start_thread_inner
pub async fn start_thread(&self, options: StartThreadOptions) -> CodexResult<NewThread> {
Box::pin(self.start_thread_inner(options, /*forked_from_thread_id*/ None)).await
}
async fn start_thread_inner(
&self,
mut options: StartThreadOptions,
forked_from_thread_id: Option<ThreadId>,
) -> CodexResult<NewThread> {
// AgentControl 在归一化阶段创建,使后续 request 已携带 agent tree 控制面。
let agent_control = self.agent_control_for_config(&options.config);
let (resumed_session_source, resumed_thread_source) = options
.initial_history
.get_resumed_session_sources()
.unwrap_or_else(|| (self.state.session_source.clone(), None));
options.session_source = Some(
// 显式 source 优先,否则恢复 history source,再退到 manager source。
options
.session_source
.take()
.unwrap_or(resumed_session_source),
);
options.thread_source = options.thread_source.take().or(resumed_thread_source);
let mut request =
ThreadSpawnRequest::new(options, Arc::clone(&self.state.auth_manager), agent_control);
request.forked_from_thread_id = forked_from_thread_id;
Box::pin(self.state.spawn_thread(request)).await
}start_thread_inner() 的关键不是转发,而是先恢复/补齐 source,再把公开 options 与 manager 级依赖 合并成内部 request。显式 session_source 使用 take() 取得优先权,只有缺省时才采用 history 或 manager 默认值。
创建参数会经过三层结构逐步补齐,最终才形成可发布的 NewThread。下面的类图突出 public options、 manager 内部 request、Session 构造参数和返回 handle 的边界。
StartThreadOptions 不携带 manager 的共享服务;ThreadSpawnRequest 才加入 auth、agent control 与继承 状态;NewThread 则只在 Session 和首事件都成功后出现。
4. spawn_thread
fresh request 进入 ThreadManagerState::spawn_thread() 后,首先消费 options:
4.1 Source
未设置 session_source 时使用 manager source;thread_source 保持可选。session_source 影响:
- restriction product;
- root/non-root 判断;
- originator 与 user instructions 加载;
- multi-agent/session identity;
- telemetry 和 rollout trace。
fresh root 通常加载 UserInstructionsProvider 的新 snapshot;non-root 则可能继承 live parent,避免 child 独立加载而产生不一致。
4.2 Environment
未提供 environment selections 时,manager 用 EnvironmentManager + config.cwd + workspace_roots 生成默认选择。这里形成 thread-level selection,真正的 environment connection/readiness snapshot 在 Session 和 Step 中捕获。
4.3 Agent版本与来源
multi-agent version 从 history、parent 和 Config feature/override 解析。fresh root 没有 history/parent, 最终使用 Config。
originator 的优先级为:受认可的 metrics service name → persisted originator → inherited parent → 环境 override → process default。fresh root 通常落到 service name、环境或默认值。
4.4 MCP失效标记
manager 在调用 Session::spawn 前把 Weak<AtomicBool> 加入 starting_mcp_runtimes。若 Plugin/Auth/MCP 源在 Session 尚未注册期间变化,全局 invalidation 会将 marker 设为 true;注册后立即补一次 refresh。
源码位置:codex-rs/core/src/thread_manager.rs :: spawn_thread startup marker
let source_changed_during_startup = Arc::new(AtomicBool::new(false));
{
// 同步锁内只清理 weak 和插入 marker,不跨 Session::spawn 的 await。
let mut starting = self
.starting_mcp_runtimes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
starting.retain(|runtime| runtime.strong_count() != 0);
// manager 只保存 Weak,启动失败不会因辅助 registry 泄漏 marker。
starting.push(Arc::downgrade(&source_changed_during_startup));
}marker 的强引用留在当前 spawn future;它覆盖的正是 Session 创建到 active map 发布之间的窗口。
spawn_thread() 对几组缺省值分别解析,再汇合为 SessionSpawnArgs。下面的流程图展示 source、环境、 history mode 和 originator 不共享一条简单的“Config 默认值”规则。
恢复路径还有一个短路:如果 InitialHistory::Resumed 的 conversation ID 已在 active map 且 thread 仍在运行,manager 会直接返回同一个 Arc<CodexThread>;如果请求携带了不同 rollout path,则返回 InvalidRequest。只有 stopped handle 才会先从 map 移除,再创建新的 Session。这个判断发生在 Session::spawn 之前,所以不会重复启动一个已经运行的 thread。
不同分支最后只生成一份 Session 参数快照。Session 初始化之后不会再回到 ThreadManager 重新询问 “当时为什么选择这个 source 或 environment”。
5. Session创建
ThreadManagerState::spawn_thread() 完成 source、environment、lineage、originator 和启动期 MCP marker 归一化后,把所有输入收敛为 SessionSpawnArgs。从 manager 视角,Session::spawn() 是一个可失败的 下游事务,只暴露两类结果:
| 结果 | Manager可以依赖的契约 | 内部实现归属 |
|---|---|---|
Ok((Arc<Session>, SessionIo)) | Session初始化完成,首事件已排队,submission loop可被等待 | Session依赖装配 |
Err(CodexErr) | 不得进入active map,Session侧负责丢弃未提交资源 | Session依赖装配 |
策略/model解析、三路 tokio::join!、LiveThreadInitGuard、Plugin/Skill/AGENTS、Hook、network proxy、空 MCP runtime、SessionConfigured 生成和 submission loop 装配都属于 Session依赖装配。ThreadManager创建 不再重复这些内部阶段。
Manager 只依赖 Session::spawn(args: SessionSpawnArgs) -> CodexResult<(Arc<Session>, SessionIo)> 这一返回契约,不依赖函数内部阶段。
Session::spawn(...).await? 成功后,源码紧接着进入 manager 自己拥有的提交阶段:
源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManagerState::spawn_thread(后续片段)
// Session返回只完成下游事务;manager仍需执行自己的首事件与registry提交。
let new_thread = self
.finalize_thread_spawn(session, io, tracked_session_source)
.await?;因此 ThreadManager创建 的下一提交点是 manager 读取首事件,而不是继续进入 SessionServices。读者若开始追问某项 service 为什么先构造、某个 startup failure 如何 discard,应立即转入 Session依赖装配。
6. Manager
从 manager 的 registry authority 看,创建过程只有一个成功终态:Registered。Session 已构造或首事件已 到达都只是中间状态,不能被 get_thread() 查到。
下面的时序图再展开首事件和 map 临界区;它不进入 Session 内部装配。
finalize_thread_spawn() 的顺序不能调换:
manager 验证两个条件:
event.id == INITIAL_SUBMIT_ID;event.msg是SessionConfigured。
否则返回 SessionConfiguredNotFirstEvent。验证成功后获取 threads write lock,用 HashMap Entry 检查 ThreadId 仍 vacant,再构造 Arc<CodexThread> 并插入。返回 NewThread 包含 ID、handle 和同一份 SessionConfigured snapshot。
如果并发路径已经注册相同 ID,当前 Session 会 shutdown_and_wait(),避免产生未追踪的后台 loop, 然后返回 invalid request。fresh UUIDv7 几乎不会碰撞,但 resume/fork 和测试可触发该保护。
源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManagerState::finalize_thread_spawn
async fn finalize_thread_spawn(
&self,
session: Arc<Session>,
io: SessionIo,
session_source: SessionSource,
) -> CodexResult<NewThread> {
let thread_id = session.thread_id();
// 等首事件时尚未持有 map 锁,Session 初始化不会阻塞其他 registry 操作。
let event = io.next_event().await?;
let session_configured = match event {
Event {
id,
msg: EventMsg::SessionConfigured(session_configured),
} if id == INITIAL_SUBMIT_ID => session_configured,
// 任何 warning/error 抢在配置事件前都视为启动协议破坏。
_ => return Err(CodexErr::SessionConfiguredNotFirstEvent),
};
{
// Entry 在单个写锁临界区内完成查重与发布,避免 check-then-insert 竞态。
let mut threads = self.threads.write().await;
if let std::collections::hash_map::Entry::Vacant(e) = threads.entry(thread_id) {
let thread = Arc::new(CodexThread::new(
session,
io,
session_configured.clone(),
session_configured.rollout_path.clone(),
session_source,
));
e.insert(thread.clone());
return Ok(NewThread {
thread_id,
thread,
session_configured,
});
}
}
// 冲突 Session 未进入 registry,必须主动关闭后台 loop,避免孤儿任务。
if let Err(err) = io.shutdown_and_wait().await {
warn!("failed to shut down duplicate thread {thread_id}: {err}");
}
Err(CodexErr::InvalidRequest(format!(
"thread {thread_id} is already running"
)))
}
pub struct NewThread {
pub thread_id: ThreadId,
pub thread: Arc<CodexThread>,
pub session_configured: SessionConfiguredEvent,
}首事件检查发生在拿 map 写锁之前,避免等待 Session 启动事件时阻塞其他注册;冲突分支又在释放写锁后 等待 shutdown,避免把异步清理放进 registry 临界区。
7. 注册后的补偿动作
spawn_thread() 拿到 NewThread 后还有两个条件动作:
- startup invalidation marker 为 true:请求 MCP runtime refresh;
- resumed thread:发送 thread-resume extension lifecycle。
fresh new thread 不执行 resume lifecycle。普通 start_thread() 本身不自动发送 thread-created broadcast;目前主要由 AgentControl 的 spawn/reload/resume 上层路径在创建成功后通知,因此监听者 不会看到尚未完成首事件握手的 agent thread。
源码位置:codex-rs/core/src/thread_manager.rs :: spawn_thread post-registration actions
let new_thread = self
.finalize_thread_spawn(session, io, tracked_session_source)
.await?;
// Acquire 与 invalidation 侧的 Release 配对,读取启动期间发生的配置变化。
if source_changed_during_startup.load(Ordering::Acquire) {
new_thread.thread.session.request_mcp_runtime_refresh();
}
// resume lifecycle 只补给恢复的 Thread,fresh Thread 已走正常 start lifecycle。
if is_resumed_thread {
new_thread.thread.emit_thread_resume_lifecycle().await;
}
Ok(new_thread)补偿动作在 finalize_thread_spawn() 成功后执行,因此 extension 与 MCP refresh 看到的是已经可查找的 Thread;任一前置启动错误都不会把半成品发布出去。
8. Ephemeral
config.ephemeral = true 不跳过 Session 创建,只改变持久化相关分支:
- 不创建 LiveThread;
- 不读取 local state DB;
- SessionConfigured 的 rollout path 为
None; - model、MCP、environment、tool 和 event runtime 仍完整存在;
- metadata update、section move 等需要 durable store 的操作会被拒绝。
ephemeral 是“运行但不 materialize thread persistence”,不是“轻量 Session”或“无历史”。SessionState 仍保存本进程内 model history,Turn 可以正常多轮执行,只是进程退出后没有可恢复记录。
9. 失败与清理矩阵
| Manager观察点 | 失败或变化 | Manager拥有的动作 | 下游资源责任 |
|---|---|---|---|
| 参数归一化 | source/environment/lineage输入无效 | 不调用或停止进入Session边界 | 尚无Session资源 |
Session::spawn | 返回任意 CodexErr | ?向调用方传播,不写active map | Session/RUN014负责未提交资源回滚 |
| 首事件握手 | first event非SessionConfigured | 返回协议错误,不发布handle | SessionIo drop触发控制面结束 |
| map注册 | ThreadId已存在 | 释放map锁后显式shutdown输掉竞态的Session | Session正常teardown |
| startup source变化 | marker为true | 注册成功后请求MCP refresh | 不回滚已发布Thread |
RUN007只判断“是否发布”和“冲突Session是否关闭”。Config lock、network proxy、MCP install等具体失败如何 释放 LiveThread 和 worker 的失败路径,见 Session依赖装配 的故障测试。
10. 观测与测试入口
启动关键阶段都有 span:thread_spawn、session_init.thread_persistence、session_init.state_db、 session_init.auth_mcp、Plugin/Skill warmup、network proxy、startup prewarm。排查慢启动应比较 span, 不要只测 start_thread() 总耗时。
测试按边界分层:
thread_manager_tests.rs:多 Thread 创建、注册、shutdown、internal source、MCP startup invalidation;session/tests.rs:身份、persistence、SessionConfigured 顺序、ephemeral 与初始化失败;thread-storetests:CreateThreadParams、LiveThread guard 和 discard;- Core integration suite:model/MCP/Skill/Hook/环境对首 Turn 的实际影响;
- App Server thread/start tests:JSON-RPC response、notification 和 live state attach。
一个完整的新 Thread 测试至少应断言:返回的 ThreadId、首个 SessionConfigured、active map 可查、 SessionSource/permission/model snapshot 正确、需要时 rollout/store 已建立,以及 shutdown 后不残留 tracked thread。若注入 startup blocker,还应验证创建完成前 map 不可见,释放后才原子发布。
11. 创建边界测试
11.1 Session首事件
finalize_thread_spawn() 要求首事件同时满足两个条件:Event ID等于 INITIAL_SUBMIT_ID,消息是 SessionConfigured。Session测试从 Session::new() 返回的 receiver读取第一项,并校验 root resume 的 session/thread identity。
源码位置:codex-rs/core/src/session/tests.rs
// :: resumed_root_session_uses_thread_id_as_session_id(断言节选)
let thread_id = ThreadId::new();
let (session, rx_event) = make_session_with_history_source_and_agent_control_and_rx(
InitialHistory::Resumed(ResumedHistory {
conversation_id: thread_id,
history: Arc::new(Vec::new()),
rollout_path: None,
}),
SessionSource::Exec,
AgentControl::default(),
)
.await
.expect("resume should succeed");
assert_eq!(session.thread_id(), thread_id);
assert_eq!(session.session_id(), SessionId::from(thread_id));
// receiver的第一项必须直接匹配SessionConfigured,不能先出现warning或MCP事件。
let event = rx_event.recv().await.expect("session configured event");
let EventMsg::SessionConfigured(event) = event.msg else {
panic!("expected session configured event");
};
assert_eq!(event.session_id, SessionId::from(thread_id));
assert_eq!(event.thread_id, thread_id);它覆盖 Session::new() 的正常首事件和身份内容,没有构造恶意/错误事件来验证 manager 返回 SessionConfiguredNotFirstEvent。后者目前只有 finalize_thread_spawn() 的 match分支作为正向证据。
Session内部的LiveThread discard、config-lock和required MCP失败测试已经集中到 Session依赖装配,本篇不再复制相同源码块。
11.2 Session阻塞
mcp_invalidation_refreshes_threads_that_are_still_starting 用 thread lifecycle contributor 阻塞 Session::spawn,此时 manager 的 active map 必须为空。释放 blocker 后创建成功,启动期间收到的 MCP invalidation 再补做 refresh。
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: mcp_invalidation_refreshes_threads_that_are_still_starting(关键断言)
let starting = tokio::spawn({
let manager = Arc::clone(&manager);
async move { manager.start_thread(StartThreadOptions::new(config)).await }
});
tokio::time::timeout(Duration::from_secs(5), observer.entered.notified())
.await
.expect("thread should enter its startup lifecycle");
// Session内部callback已经运行,manager仍未执行finalize/register。
assert!(manager.list_thread_ids().await.is_empty());
manager.invalidate_mcp_runtimes().await;
observer.release.notify_one();
starting
.await
.expect("thread startup task should finish")
.expect("thread should start");
tokio::time::timeout(Duration::from_secs(5), observer.refreshed.notified())
.await
.expect("invalidation during startup should refresh the newly published thread");这组测试属于 ThreadManager创建,因为断言对象是 manager registry 与启动期 marker;callback 内部为什么阻塞、Session 如何装配 extension,则属于 Session依赖装配。
11.3 Resume复用Thread
常规 active resume不会创建第二个 Session再等最终 map冲突。测试先启动并materialize source Thread,再用 同一 rollout path resume;返回的 ThreadId相同,并且 Arc::ptr_eq 证明是原 CodexThread。
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: resume_active_thread_from_rollout_returns_running_thread(断言节选)
let source = manager
.start_thread(StartThreadOptions::new(config.clone()))
.await
.expect("start source thread");
source.thread.ensure_rollout_materialized().await;
source.thread.flush_rollout().await.expect("flush source rollout");
let rollout_path = source
.thread
.rollout_path()
.expect("source rollout path should exist");
let resumed = manager
.resume_thread_from_rollout(
config,
rollout_path,
auth_manager,
/*parent_trace*/ None,
ClientMcpExtensions::default(),
)
.await
.expect("resume active source thread");
// fast path复用原Arc,不建立第二条submission loop。
assert_eq!(resumed.thread_id, source.thread_id);
assert!(Arc::ptr_eq(&resumed.thread, &source.thread));这不等于 Entry::Occupied 分支已有测试。两个并发启动仍可能同时越过早期查重,最终由 vacant entry串行化; 源码会 shutdown输掉竞态的第二个 Session并返回 InvalidRequest,但当前没有 barrier test强制两个 spawn同时 到达该写锁。这个缺口应保留在后续测试任务中。
12. 创建事务验证
- 从
start_thread(StartThreadOptions)开始,按顺序说出 options归一化、Session调用边界、SessionConfigured首事件、vacant entry和补偿refresh五个节点。 - 解释为什么 Session已经是
Arc<Session>仍不代表 Thread可被get_thread()查到。 - 给定“创建返回错误,但下一次创建报duplicate writer”,判断应优先检查 LiveThread guard还是 manager map, 并指出各自的测试覆盖和当前缺口。
- 用只读搜索核对三道提交边界:
rg -n "LiveThreadInitGuard|SessionConfiguredNotFirstEvent|finalize_thread_spawn" \
codex-rs/core/src codex-rs/thread-store/src
rg -n "resumed_root_session_uses|discard_thread_drops|resume_active_thread" \
codex-rs/core/src/session/tests.rs codex-rs/core/src/thread_manager_tests.rs \
codex-rs/thread-store/src/local/mod.rs继续阅读 ThreadManager恢复 时,重点观察相同创建事务如何更换 InitialHistory 和 LiveThread open方式;阅读 Session依赖装配 时,则把本文的“Session内部黑盒”展开成具体依赖屏障。
