ThreadManager管理
ThreadManager 的 threads map 很容易被误读成“全部 Thread”。它实际只保存当前进程已经发布的 Arc<CodexThread>;ThreadStore 另行保存未加载或已归档记录,正在启动但尚未通过首事件握手的 Session 也还没有进入 map。
因此,线程管理不是一个 CRUD 表。list、get、remove、shutdown 和 delete 分别作用于不同 owner:有的只改内存索引,有的驱动 Session loop,有的修改持久化数据。把它们合并成一个“关闭线程” 操作,会制造孤儿 task、丢失 handle,或者错误删除仍被 lineage 引用的历史。
本文面向已读过 ThreadManager依赖 的读者,默认知道 threads map与ThreadStore是不同owner。本文只讲运行期管理操作和并发窗口,不重复创建、恢复和Fork的 内部流程;需要理解 stopped handle为何能用同一ThreadId重建,可先回看 ThreadManager恢复。
读完后,应能根据现象选择 live map、submission loop还是ThreadStore作为事实源,并能解释remove、shutdown 和delete为何没有可交换性,以及批量shutdown失败项为何必须继续留在map。
1. 四种存在状态
同一个 ThreadId 可以同时处在多套视图中:
| 视图 | 权威 owner | “存在”的含义 | 典型反例 |
|---|---|---|---|
| Persisted | ThreadStore | metadata/rollout 或远端记录仍可读取 | 进程重启后没有 live handle |
| Loaded | ThreadManagerState.threads | 当前进程已发布 Arc<CodexThread> | Session loop 可能已经停止 |
| Running | SessionIo.tx_sub | submission receiver 尚未关闭 | running 不等于当前有 ActiveTurn |
| Publicly visible | manager 的 list/get 过滤 | loaded 且 SessionSource 非 internal | internal maintenance Thread 仍在 map |
“NotLoaded”是产品状态,不是“不存在”;“Idle”是没有活动工作,不是“不在 map”;“remove”只是解除 manager 索引,不是停止运行或删除 rollout。
2. 一张操作地图
下面的 mindmap 先按管理问题选择入口。它替代固定的线性流程,因为 Thread 管理本来就有多条正交路径。
查“当前进程能否提交 Op”进入 runtime lifecycle;查“历史是否仍存在”进入 persistence;查“为什么列表 看不到”进入 visibility。先选权威视图,才能避免拿 map 结果推断磁盘状态。
3. Live注册表
ThreadManager 外层方法只委托给共享 State。真正的 map 是 Arc<RwLock<HashMap<ThreadId, Arc<CodexThread>>>>:读路径复制 ID 或 clone Arc,写路径注册和移除。
源码位置:codex-rs/core/src/thread_manager.rs :: list_thread_ids, get_thread, remove_thread
pub async fn list_thread_ids(&self) -> Vec<ThreadId> {
self.state.list_thread_ids().await
}
pub async fn get_thread(&self, thread_id: ThreadId) -> CodexResult<Arc<CodexThread>> {
self.state.get_thread(thread_id).await
}
pub async fn remove_thread(&self, thread_id: &ThreadId) -> Option<Arc<CodexThread>> {
// remove 返回原 Arc,让调用方仍有机会显式 shutdown 或检查状态。
self.state.threads.write().await.remove(thread_id)
}
impl ThreadManagerState {
pub(crate) async fn list_thread_ids(&self) -> Vec<ThreadId> {
self.threads
.read()
.await
.iter()
// public list 隐藏 internal source,但不会过滤 stopped handle。
.filter_map(|(thread_id, thread)| {
(!thread.session_source.is_internal()).then_some(*thread_id)
})
.collect()
}
pub(crate) async fn get_thread(
&self,
thread_id: ThreadId,
) -> CodexResult<Arc<CodexThread>> {
let threads = self.threads.read().await;
match threads.get(&thread_id) {
// clone Arc 后释放 read lock,后续异步操作不占用 registry 锁。
Some(thread) if !thread.session_source.is_internal() => Ok(thread.clone()),
Some(_) | None => Err(CodexErr::ThreadNotFound(thread_id)),
}
}
}get_thread() 对 internal 与 absent 返回同一个 ThreadNotFound,避免公共调用方探测 internal maintenance Thread。list_thread_ids() 也只回答“manager 当前公开索引了哪些 ID”,不调用 ThreadStore,更不检查 is_running()。
map 与 store 都能按 ThreadId 查找,但返回对象不同:前者是可提交 Op 的进程内 handle,后者是可序列化 记录。CodexThread clone 不复制 Session;多个调用方共享同一 Arc。
3.1 Running
源码没有通过 ActiveTurn、AgentStatus 或 task handle 定义 is_running()。它只检查 submission sender 是否已经关闭。
源码位置:codex-rs/core/src/codex_thread.rs :: CodexThread::is_running, SessionIo::shutdown_and_wait
pub(crate) fn is_running(&self) -> bool {
// channel open 表示 session loop 仍能接收控制消息;Idle Thread 仍属于 running。
!self.io.tx_sub.is_closed()
}
pub(crate) async fn shutdown_and_wait(&self) -> CodexResult<()> {
let session_loop_termination = self.session_loop_termination.clone();
match self.submit(Op::Shutdown).await {
Ok(_) => {}
// receiver 已经消失时仍等待 termination;重复 shutdown 不因此失败。
Err(err) if matches!(err.details(), CodexErrorDetails::InternalAgentDied) => {}
Err(err) => return Err(err),
}
session_loop_termination.await;
Ok(())
}shutdown_and_wait() 的完成条件是 session loop termination,不是 Shutdown submission 入队。它可以容忍 receiver 已关闭,因为这种情况下真正需要的仍是等待后台 loop 收尾。
4. 移除关闭删除
下面的状态图刻意把内存索引、运行 loop 和持久化记录拆开。某些转换由不同产品层 API 触发,不是 ThreadManager::remove_thread() 一次完成。
Detached 最容易产生误解:manager 已无法通过 ID 找到 Thread,但持有返回 Arc 的调用方仍可提交或关闭 它。反过来,LoadedStopped 仍可能出现在 list/get 中,因为这些方法不检查 sender 是否关闭。
| 操作 | 修改 map | 发送 Shutdown | 等待 loop | 删除持久化 |
|---|---|---|---|---|
remove_thread | 是 | 否 | 否 | 否 |
CodexThread::shutdown_and_wait | 否 | 是 | 是 | 否 |
shutdown_all_threads_bounded 成功项 | 是 | 是 | 是 | 否 |
| ThreadStore delete | 不直接修改 | 否 | 否 | 是 |
因此手工管理单个 Thread 时,常见安全顺序是先保留 Arc、执行 shutdown_and_wait(),再从 manager map 移除。若先 remove,调用方必须保证仍持有返回 handle,否则失去等待 loop 的直接入口。
5. 三个并发窗口
5.1 启动中的 Session
Session 只有在首个 SessionConfigured 通过 finalize_thread_spawn() 后才进入 map。启动 lifecycle、 MCP 或 persistence 阶段阻塞时,list_thread_ids() 和批量 shutdown 的 map snapshot 都看不到它。
这是宿主级协调要求:进入全局 shutdown 前应停止接受新 start/resume,并等待在途创建结束。manager 的 starting_mcp_runtimes 只补 MCP invalidation 窗口,不是通用 startup registry。
5.2 Remove
remove_thread() 在 write lock 内只执行一次 HashMap remove,不等待 Arc strong count 归零。listener、 App Server state 或 agent control 仍可能持有 handle。manager 也不会向这些 owner 广播“你必须停止”。
对于异步延迟清理,当前实现还提供 remove_thread_if_matches():它在同一个 write lock 内用 Arc::ptr_eq 比较预期 handle,只有仍然对应原运行时对象时才删除。这使旧 shutdown 任务不会误删后来 以相同 ThreadId 注册的新 handle;需要无条件 detach 时才使用 remove_thread()。
这不是 Rust 内存安全问题,而是业务生命周期问题:被 detach 的 Session 仍可产生事件和持有 MCP/ exec 资源。调用方若把 remove 当 shutdown,就会制造 manager 无法再枚举的后台工作。
5.3 ID替换关闭
批量 shutdown 先 clone (ThreadId, Arc<CodexThread>) 快照,等待所有结果后,再按完成 ID 从当前 map 删除。源码没有比较 map 中的 Arc 是否仍等于快照 Arc。
当前版本提供 remove_thread_if_matches() 修复这个竞态:延迟清理可以把预期的 Arc<CodexThread> 传入, 只有 map 中仍由同一个 Arc::ptr_eq 对象占据该 ThreadId 时才移除。无条件 remove_thread() 仍适合 明确的宿主 teardown,但 replacement-sensitive 的异步清理应使用 identity-checked 入口。
6. 批量回收完成项
shutdown_all_threads_bounded() 不串行等待每个 Thread。它先释放 map read lock,再用 FuturesUnordered 并发关闭,最后只移除 Complete 项。
源码位置:codex-rs/core/src/thread_manager.rs :: shutdown_all_threads_bounded
pub async fn shutdown_all_threads_bounded(&self, timeout: Duration) -> ThreadShutdownReport {
let threads = {
let threads = self.state.threads.read().await;
// 只在 clone Arc 快照时持 read lock,shutdown await 不阻塞其他 map 访问。
threads
.iter()
.map(|(thread_id, thread)| (*thread_id, Arc::clone(thread)))
.collect::<Vec<_>>()
};
let mut shutdowns = threads
.into_iter()
.map(|(thread_id, thread)| async move {
let outcome = match tokio::time::timeout(timeout, thread.shutdown_and_wait()).await {
Ok(Ok(())) => ShutdownOutcome::Complete,
Ok(Err(_)) => ShutdownOutcome::SubmitFailed,
Err(_) => ShutdownOutcome::TimedOut,
};
(thread_id, outcome)
})
.collect::<FuturesUnordered<_>>();
let mut report = ThreadShutdownReport::default();
while let Some((thread_id, outcome)) = shutdowns.next().await {
match outcome {
ShutdownOutcome::Complete => report.completed.push(thread_id),
ShutdownOutcome::SubmitFailed => report.submit_failed.push(thread_id),
ShutdownOutcome::TimedOut => report.timed_out.push(thread_id),
}
}
let mut tracked_threads = self.state.threads.write().await;
// 失败和超时项继续留在 map,调用方可以重试或检查。
for thread_id in &report.completed {
tracked_threads.remove(thread_id);
}
report.completed.sort_by_key(std::string::ToString::to_string);
report.submit_failed.sort_by_key(std::string::ToString::to_string);
report.timed_out.sort_by_key(std::string::ToString::to_string);
report
}每个 Thread 都拥有完整的 timeout,而不是整批共享一个总 deadline;两条慢 Thread 可以并行超时。最终 排序让 report 不依赖 future 完成顺序,便于日志、测试和上层重试使用稳定结果。
internal Thread 虽然不出现在 public list/get 中,仍存在于底层 map,所以批量 shutdown 会包含它。仓库 测试用一个 internal memory-consolidation Thread 验证 list 为空,但 report 仍包含其 ID。
7. 元数据操作分流
metadata update 不能无条件直写 ThreadStore。loaded Thread 必须经过 CodexThread/LiveThread,让 metadata 变更与 rollout writer 保持顺序;cold Thread 才直接调用 store。ephemeral Thread 没有 durable owner, 因此拒绝更新。
源码位置:codex-rs/core/src/thread_manager.rs :: update_thread_metadata
pub async fn update_thread_metadata(
&self,
thread_id: ThreadId,
patch: ThreadMetadataPatch,
include_archived: bool,
) -> CodexResult<StoredThread> {
if let Ok(thread) = self.get_thread(thread_id).await {
if thread.config_snapshot().await.ephemeral {
// ephemeral 没有可承诺的 durable metadata,不能假装更新成功。
return Err(CodexErr::InvalidRequest(format!(
"ephemeral thread does not support metadata updates: {thread_id}"
)));
}
return thread
.update_thread_metadata(patch, include_archived)
.await
.map_err(|err| thread_store_metadata_update_error(thread_id, err));
}
// get_thread 的 NotFound 还可能表示 cold Thread;继续让 store 做权威存在性判断。
self.state
.thread_store
.update_thread_metadata(UpdateThreadMetadataParams {
thread_id,
patch,
include_archived,
})
.await
.map_err(|err| match err {
ThreadStoreError::ThreadNotFound { thread_id } => CodexErr::ThreadNotFound(thread_id),
err => thread_store_metadata_update_error(thread_id, err),
})
}这里不能把第一次 get_thread() 失败直接返回给用户,因为 absent-from-map 不等于 absent-from-store。 不过 internal Thread 同样被 public get 隐藏,因此若内部调用要管理 internal metadata,应使用明确的内部 路径,而不是误走公共 fallback。
move_thread_to_section() 采用相似的 ephemeral 检查,但即使目标 Thread 已加载,也直接交给 ThreadStore 统一处理;它可能同时影响多个 cold/live Thread 的 section position,因此不是经由 CodexThread/LiveThread 排队的 rollout 更新。
8. 运维诊断卡片
| 现象 | 先查什么 | 不要立即推断 |
|---|---|---|
| list 中没有 Thread | 是否 cold、archived 或 internal | rollout 已删除 |
| get 返回 ThreadNotFound | map 是否未加载;source 是否 internal | store 中不存在 |
| list 中有 ID 但提交失败 | is_running()、submission receiver、loop termination | map 自动保证 running |
| remove 后仍有事件/进程 | 是否还有 Arc owner;是否执行 shutdown | Rust Arc 泄漏或 store 错误 |
| shutdown batch 后仍有 ID | report 的 submit_failed/timed_out | remove 逻辑失效 |
| shutdown 期间出现新 ID | 是否有在途 start/resume | batch 一定覆盖全部启动任务 |
| metadata 更新绕过 live ordering | 是否错误直写 ThreadStore | JSONL 与 SQLite 会自动排序 |
| internal Thread 未列出但被关闭 | batch 读取底层 map,不走 public filter | visibility 规则不一致 |
验证线程管理至少需要三类断言:public list/get 的过滤语义、shutdown report 与 map 清理一致性,以及 remove 后外部 Arc 的生命周期。只测试 HashMap 长度,无法证明后台 loop、持久化顺序和 internal Thread 都被正确管理。
9. 测试与并发缺口
9.1 成功批量关闭
测试启动两个Thread后调用 shutdown_all_threads_bounded(10s)。完成列表先按ThreadId字符串排序,再与预期 比较;submit_failed和timed_out必须为空,manager public list最终也必须为空。
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: shutdown_all_threads_bounded_submits_shutdown_to_every_thread(断言节选)
let thread_1 = manager
.start_thread(StartThreadOptions::new(config.clone()))
.await
.expect("start first thread")
.thread_id;
let thread_2 = manager
.start_thread(StartThreadOptions::new(config.clone()))
.await
.expect("start second thread")
.thread_id;
let report = manager
.shutdown_all_threads_bounded(Duration::from_secs(10))
.await;
let mut expected_completed = vec![thread_1, thread_2];
expected_completed.sort_by_key(std::string::ToString::to_string);
// completed顺序稳定,且只有完成项才应从map移除。
assert_eq!(report.completed, expected_completed);
assert!(report.submit_failed.is_empty());
assert!(report.timed_out.is_empty());
assert!(manager.list_thread_ids().await.is_empty());它证明成功路径的fan-out、排序和map清理,没有证明失败/超时项保留:当前测试没有注入关闭submission失败或 超过timeout的CodexThread。因此正文对 submit_failed、timed_out 的保留语义来自实现分支,仍缺直接 反向测试。
9.2 Active与Stopped
两个Resume测试使用同一rollout path。source仍running时,返回同一个 Arc<CodexThread>;source先 shutdown_and_wait() 后,Resume保留ThreadId但创建新的Arc。这直接锁定Loaded/Running两个状态不能合并。
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: resume_active_thread_from_rollout_returns_running_thread(核心断言)
let resumed = manager
.resume_thread_from_rollout(
config,
rollout_path,
auth_manager,
/*parent_trace*/ None,
ClientMcpExtensions::default(),
)
.await
.expect("resume active source thread");
assert_eq!(resumed.thread_id, source.thread_id);
// running fast path复用同一handle,不新建Session loop。
assert!(Arc::ptr_eq(&resumed.thread, &source.thread));停止后的测试重新建立独立fixture,先关闭source再执行同样的Resume入口:
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: resume_stopped_thread_from_rollout_spawns_new_thread(核心断言)
source
.thread
.shutdown_and_wait()
.await
.expect("shutdown source thread");
let resumed = manager
.resume_thread_from_rollout(
config,
rollout_path,
auth_manager,
/*parent_trace*/ None,
ClientMcpExtensions::default(),
)
.await
.expect("resume stopped source thread");
assert_eq!(resumed.thread_id, source.thread_id);
// stopped handle被清理后,同一业务身份绑定到新运行时Arc。
assert!(!Arc::ptr_eq(&resumed.thread, &source.thread));两段分别来自独立测试,不能连在同一作用域执行;它们共同证明业务ThreadId与运行时Arc是两条身份轴。
9.3 三个尚无直接测试
当前源码和测试中没有找到直接覆盖以下场景的测试:
- 构造 submit失败和永不终止的fake
CodexThread,断言batch report分类并保留map项; remove_thread()返回Arc后,继续通过外部Arc提交Op并证明manager list已不可见;- 让旧handle shutdown与同ThreadId replacement在最终remove前精确交错,验证或暴露按ID删除新handle的 竞态。
第三项尤其重要:当前实现只按 completed ThreadId执行 tracked_threads.remove(thread_id),不比较Arc身份。 因此正文把“宿主停止阶段不得并发replacement”标成使用前提,而不是已由manager内部保证的不变量。
internal Thread 的另一条测试提供补充观察:public list/get 隐藏它,但 batch shutdown report 仍包含其 ID,说明 批量关闭读取底层map而不是公共过滤视图。该测试已在 ThreadManager依赖 展开,本文不重复粘贴。
10. 线程移除验证
- 给定一个ThreadId,分别说明 Persisted、Loaded、Running、Publicly visible四种“存在”由谁回答。
- 解释为什么
remove_thread()后Session可能继续运行,以及为什么这不是Rust内存泄漏。 - 给定 batch report 含
timed_out,说明该 Thread 为什么必须留在 map;再说明目前哪条结论只有源码路径支撑。 - 使用只读命令核对成功测试、active/stopped Resume和按ID remove:
rg -n "shutdown_all_threads_bounded|tracked_threads.remove|remove_thread" \
codex-rs/core/src/thread_manager.rs
rg -n "shutdown_all_threads_bounded_submits|resume_active_thread|resume_stopped_thread" \
codex-rs/core/src/thread_manager_tests.rs继续学习时,阅读 CodexThread公共API,理解外部Arc实际还能调用哪些方法; 再阅读 CodexThread背压,观察detached handle为何仍可能产生事件。
