Skip to content

ThreadManager管理

分析 ThreadManager 的查找、列表、移除、批量关闭与元数据分流,厘清 loaded、running 和 persisted 三套状态。

基于rust-v0.150.0
CodexRustThreadManagerConcurrency

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“存在”的含义典型反例
PersistedThreadStoremetadata/rollout 或远端记录仍可读取进程重启后没有 live handle
LoadedThreadManagerState.threads当前进程已发布 Arc<CodexThread>Session loop 可能已经停止
RunningSessionIo.tx_subsubmission receiver 尚未关闭running 不等于当前有 ActiveTurn
Publicly visiblemanager 的 list/get 过滤loaded 且 SessionSource 非 internalinternal 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

rust
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

rust
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

rust
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

rust
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 或 internalrollout 已删除
get 返回 ThreadNotFoundmap 是否未加载;source 是否 internalstore 中不存在
list 中有 ID 但提交失败is_running()、submission receiver、loop terminationmap 自动保证 running
remove 后仍有事件/进程是否还有 Arc owner;是否执行 shutdownRust Arc 泄漏或 store 错误
shutdown batch 后仍有 IDreport 的 submit_failed/timed_outremove 逻辑失效
shutdown 期间出现新 ID是否有在途 start/resumebatch 一定覆盖全部启动任务
metadata 更新绕过 live ordering是否错误直写 ThreadStoreJSONL 与 SQLite 会自动排序
internal Thread 未列出但被关闭batch 读取底层 map,不走 public filtervisibility 规则不一致

验证线程管理至少需要三类断言: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

rust
// :: 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

rust
// :: 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

rust
// :: 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 三个尚无直接测试 ​

当前源码和测试中没有找到直接覆盖以下场景的测试:

  1. 构造 submit失败和永不终止的fake CodexThread,断言batch report分类并保留map项;
  2. remove_thread() 返回Arc后,继续通过外部Arc提交Op并证明manager list已不可见;
  3. 让旧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. 线程移除验证 ​

  1. 给定一个ThreadId,分别说明 Persisted、Loaded、Running、Publicly visible四种“存在”由谁回答。
  2. 解释为什么 remove_thread() 后Session可能继续运行,以及为什么这不是Rust内存泄漏。
  3. 给定 batch report 含 timed_out,说明该 Thread 为什么必须留在 map;再说明目前哪条结论只有源码路径支撑。
  4. 使用只读命令核对成功测试、active/stopped Resume和按ID remove:
bash
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为何仍可能产生事件。