Skip to content

ThreadManager依赖

逐字段解析 ThreadManager 的活跃线程注册表、共享服务、派生管理器、宿主 provider、持久化与多 Agent 依赖。

基于rust-v0.150.0
CodexRustThreadManagerDependency Injection

ThreadManager依赖 ​

ThreadManager 是 Core 的进程级 thread factory 与 live registry。它不只是一个 HashMap<ThreadId, CodexThread>:创建 Session 还需要认证、模型目录、执行环境、Plugin/Skill/MCP、 Code Mode、Extension、ThreadStore、AgentGraph、attestation、time provider、安装身份和 telemetry。

这些依赖的生命周期比单个 Thread 更长,因此集中保存在 Arc<ThreadManagerState>,每次 spawn 时再 clone 到新的 Session。理解字段所有权,才能判断配置刷新应改 manager、已有 Session,还是只影响下 一次 Step。

本文面向已经了解 Thread、Session、Turn 基本区别,但第一次阅读 ThreadManager 字段的读者。建议先读 Core运行时架构总览 和 Thread与Turn概念模型。本文只解释 manager 的字段、owner、共享方式和 invalidation 边界;不展开创建事务、恢复和分叉的完整流程,它们分别由 ThreadManager创建~ThreadManager分叉 负责。

读完后,应能回答三个问题:为什么 live map 不是持久化事实源;为什么 Plugin、Skill、MCP 处于三层 owner;以及为什么全局 MCP invalidation 需要同时处理 loaded Thread 和尚未注册的启动中 Session。

1. 外层 handle 很薄 ​

源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManager

rust
pub struct ThreadManager {
    // 所有 clone 共享同一 State,manager handle 本身不复制 registry。
    state: Arc<ThreadManagerState>,
    // guard 仅让测试临时 CODEX_HOME 随 manager 生命周期自动清理。
    _test_codex_home_guard: Option<TempCodexHomeGuard>,
}

ThreadManager 没有实现自己的并发 map 或业务缓存;生产状态都在 ThreadManagerState。第二个字段只 服务测试:测试构造器创建临时 CODEX_HOME,guard drop 时删除目录,普通业务 manager 为 None。

State 使用 Arc 有两个原因:

  1. App Server 等宿主通常再用 Arc<ThreadManager> 共享 manager;
  2. AgentControl 需要 Weak<ThreadManagerState> 回到全局 thread registry,又不能形成 State → CodexThread → Session → AgentControl → State 强引用环。

这也解释了 with_code_mode_session_provider() 的限制:它通过 Arc::get_mut() 替换 provider,只能在 State 尚未共享前调用。manager 一旦被 thread 或宿主共享,进程级 provider 不能再原地换掉。

2. Manager状态 ​

源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManagerState

rust
pub(crate) struct ThreadManagerState {
    // map 是 loaded runtime 的权威集合;ThreadStore 另管未加载历史。
    threads: Arc<RwLock<HashMap<ThreadId, Arc<CodexThread>>>>,
    // broadcast 只做增量提示,lag 后必须回到 registry 查询。
    thread_created_tx: broadcast::Sender<ThreadId>,
    thread_id_generator: ThreadIdGenerator,
    // 以下 Arc 服务跨 Thread 复用,spawn 时再 clone 给 Session。
    auth_manager: Arc<AuthManager>,
    models_manager: SharedModelsManager,
    environment_manager: Arc<EnvironmentManager>,
    // Weak marker 覆盖“正在启动但尚未注册”的 MCP refresh 窗口。
    starting_mcp_runtimes: std::sync::Mutex<Vec<std::sync::Weak<AtomicBool>>>,
    skills_service: Arc<HostSkillsService>,
    plugins_manager: Arc<PluginsManager>,
    mcp_manager: Arc<McpManager>,
    code_mode_session_provider: Arc<dyn CodeModeSessionProvider>,
    extensions: Arc<ExtensionRegistry<Config>>,
    user_instructions_provider: Arc<dyn UserInstructionsProvider>,
    thread_store: Arc<dyn ThreadStore>,
    agent_graph_store: Option<Arc<dyn AgentGraphStore>>,
    attestation_provider: Option<Arc<dyn AttestationProvider>>,
    external_time_provider: Option<Arc<dyn TimeProvider>>,
    session_source: SessionSource,
    installation_id: String,
    analytics_events_client: Option<AnalyticsEventsClient>,
    ops_log: Option<SharedCapturedOps>,
}

字段可以分为四组:

分组字段所有权来源
Live registrythreads、thread_created_txmanager 自己创建
派生管理器skills、plugins、MCP、Code Mode provider根据 Config、auth、extension 构造
宿主注入服务auth、models、environment、extension、instructions、storeThreadManager::new() 参数
可选宿主能力graph、attestation、time、analytics宿主按产品能力传入
身份与测试session source、installation ID、starting MCP markers、ops log构造输入或内部辅助

字段列表按类型展开后,可以看到 manager 同时承担 registry owner、共享服务容器和派生 manager 工厂, 但每个 Session 只 clone 自己需要的 handle。

ThreadManagerState → CodexThread 是 live ownership,ThreadManagerState → ThreadStore 则是跨加载状态 的持久化依赖。二者生命周期不同,不能用一个 map 同时替代。

3. 所有权图 ​

图中的“Manager → Session”表示 spawn 时 clone Arc 或值快照,不表示所有服务只有 manager 才能访问。 Session 创建后会在 SessionServices 中独立持有这些 Arc;即使 manager 从 active map 移除 thread, Session 仍可能因其他 handle 存活而继续持有服务。

4. Live registry ​

4.1 threads ​

类型是 Arc<RwLock<HashMap<ThreadId, Arc<CodexThread>>>>:

  • read lock:get、list、clone 全部 loaded handles;
  • write lock:注册、remove、resume 时清理 stopped handle;
  • value 用 Arc:listener、App Server state 和 agent control 可在释放 map lock 后使用 thread;
  • map 只保存 loaded thread,不等于 ThreadStore 的全部 persisted thread。

list_thread_ids() 还会过滤 session_source.is_internal() 的 thread;manager map 内存在不等于公共 list/get API 必须暴露。get_thread() 对 internal source 同样返回 ThreadNotFound。

源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManagerState::list_thread_ids

rust
pub(crate) async fn list_thread_ids(&self) -> Vec<ThreadId> {
    // 读锁只覆盖 map 遍历;collect 后不把 guard 泄漏给调用方。
    self.threads
        .read()
        .await
        .iter()
        // internal thread 可以存在于 map,但不会进入公共列表。
        .filter_map(|(thread_id, thread)| {
            (!thread.session_source.is_internal()).then_some(*thread_id)
        })
        .collect()
}

读锁只覆盖 map 遍历,返回的是复制的 ThreadId,不会把锁生命周期泄漏给调用方。过滤发生在 registry 读取处,也避免调用方先取得 internal handle 再自行判断。

4.2 创建通知与注册表 ​

thread_created_tx 是容量 1024 的 broadcast。订阅者用 subscribe_thread_created() 获得独立 receiver; notify_thread_created() 只发送 ID,不携带 handle 或完整配置。

源码位置:codex-rs/core/src/thread_manager.rs :: thread-created subscription and notification

rust
pub fn subscribe_thread_created(&self) -> broadcast::Receiver<ThreadId> {
    // 每次 subscribe 都取得独立游标,不会消费其他观察者的通知。
    self.state.thread_created_tx.subscribe()
}

pub(crate) fn notify_thread_created(&self, thread_id: ThreadId) {
    // 没有 receiver 时发送失败可忽略;active map 才是权威状态。
    let _ = self.thread_created_tx.send(thread_id);
}

pub fn auth_manager(&self) -> Arc<AuthManager> {
    // accessor clone Arc,使宿主复用身份服务而不转移 manager 所有权。
    self.state.auth_manager.clone()
}

pub fn skills_service(&self) -> Arc<HostSkillsService> {
    self.state.skills_service.clone()
}

pub fn plugins_manager(&self) -> Arc<PluginsManager> {
    self.state.plugins_manager.clone()
}

pub fn list_collaboration_modes(&self) -> Vec<CollaborationModeMask> {
    self.state.models_manager.list_collaboration_modes()
}

pub fn environment_manager(&self) -> Arc<EnvironmentManager> {
    self.state.environment_manager.clone()
}

pub fn get_models_manager(&self) -> SharedModelsManager {
    self.state.models_manager.clone()
}

pub async fn list_models(
    &self,
    refresh_strategy: RefreshStrategy,
    http_client_factory: codex_http_client::HttpClientFactory,
) -> Vec<ModelPreset> {
    self.state
        .models_manager
        .list_models(refresh_strategy, http_client_factory)
        .await
}

pub fn mcp_manager(&self) -> Arc<McpManager> {
    self.state.mcp_manager.clone()
}

这些方法返回 receiver、复制 ID 或 clone Arc,都没有把 State 内部锁与引用借用暴露到宿主层。

broadcast 可以 lag,也不重放订阅前的历史,所以权威查询仍是 active map/ThreadStore。通知的用途是 提示 App Server 等宿主增量 attach listener,而不是维护第二份 thread registry。

公共 get 与 agent-tree edge 查询都复用同一份 map,但过滤语义不同:

源码位置:codex-rs/core/src/thread_manager.rs :: list_live_thread_spawn_edges, get_thread

rust
pub(crate) async fn list_live_thread_spawn_edges(&self) -> Vec<(ThreadId, ThreadId)> {
    self.threads
        .read()
        .await
        .iter()
        .filter_map(|(thread_id, thread)| {
            // internal thread 不属于公开 agent tree,必须在解析 source 前排除。
            if thread.session_source.is_internal() {
                return None;
            }
            match &thread.session_source {
                SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
                    parent_thread_id,
                    ..
                }) => Some((*parent_thread_id, *thread_id)),
                _ => None,
            }
        })
        .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 后调用方可在 map 锁释放后安全使用 live handle。
        Some(thread) if !thread.session_source.is_internal() => Ok(thread.clone()),
        // 对外统一返回 NotFound,避免暴露 internal thread 是否存在。
        Some(_) | None => Err(CodexErr::ThreadNotFound(thread_id)),
    }
}

list_live_thread_spawn_edges() 只投影 ThreadSpawn parent-child 边;get_thread() 则把 internal 与不存在 统一成同一错误。两者都没有把 RwLockReadGuard 带出方法。

从 manager 视角看,一个 Thread 会在未加载、启动中、loaded 和 stopped 之间切换。下面的状态图强调 “仍在 map 中但 loop 已停止”是独立状态,下一次 resume 会先移除该 handle。

从 map remove 不会删除 ThreadStore 记录,也不保证其他 Arc<CodexThread> 已释放;它只结束 manager 对 loaded handle 的权威索引。

5. 认证与模型依赖 ​

5.1 auth_manager ​

Arc<AuthManager> 由宿主构造并注入。ThreadManager 用它:

  • 创建 PluginsManager 时读取当前 API auth mode;
  • 创建 model provider/models manager;
  • 每次 spawn 传入 Session,供模型、MCP、extension 和 refresh 使用;
  • 通过 public accessor 给 App Server/account processor 等宿主能力复用。

它是可刷新身份的共享 manager,不是 Thread 创建时复制的一份 token。多个 Session 因此可以观察到同一 账户的 auth refresh,但每次请求仍由 provider 解析当前 auth。

5.2 models_manager ​

SharedModelsManager 由宿主注入,常用 build_models_manager(config, auth_manager) 创建。helper 先构造 runtime model provider,再把 codex_home 和可选 model catalog 交给 provider 的 models manager。

它负责 provider model list、cache、default model 和 collaboration modes。ThreadManager 的 list_models()/list_collaboration_modes() 只是委托;Session spawn 再用同一 manager 解析本 thread model snapshot。把 models manager 放在 manager 级可以让多个 Thread 共享 catalog/cache,避免每个 Session 重复请求 /models。

6. Environment管理 ​

Arc<EnvironmentManager> 由宿主注入,因为 local、remote 或测试环境的注册方式属于产品层。 ThreadManager 用它:

  • 从 cwd/workspace roots 生成 default environment selections;
  • 校验 selection ID 是否重复、environment 是否存在;
  • 传入每个 Session,构造 thread-scoped ThreadEnvironments;
  • 供 App Server environment RPC 与 extension provider 使用同一注册表。

Manager 只验证“选择的环境存在”,不会在这里捕获每个 Turn 的 cwd、capability roots 或 readiness。 这些属于 Session/StepContext。环境注册表共享,环境快照按 Thread/Step 捕获。

7. 扩展管理器 ​

构造器不是让宿主分别传入三者,而是按依赖顺序创建:

text
Config.codex_home + SessionSource.restriction_product + auth mode
  → PluginsManager
  → HostSkillsService
  → McpManager(PluginsManager + ExtensionRegistry + Apps tool cache)

7.1 plugins_manager ​

PluginsManager 持有 marketplace/config/cache、auth policy 和可选 analytics。它按 PluginsConfigInput 解析 effective plugins,并为 Skill roots、MCP registrations、Hook sources、apps 和推荐能力提供统一 来源。

它是 process/thread-manager scoped cache。Session reload 可以调用 clear_cache(),但不为每个 Thread 重建一个 manager。

7.2 skills_service ​

HostSkillsService 由 codex_home、bundled-skills flag 和 restriction product 创建。它负责按 Config + plugin skill roots + executor filesystem 生成 skill snapshot,并缓存加载结果。

Skill content 最终按 Session/Turn 注入;服务本身可跨 Thread 复用发现与缓存。restriction product 在 manager 创建时固定,避免不同 SessionSource 意外看到不允许的系统/产品 Skill。

7.3 mcp_manager ​

McpManager 依赖 PluginsManager、ExtensionRegistry<Config> 与 Codex Apps tool cache。它不直接等于 某个 Thread 的 live MCP connections;它负责把配置、plugin registration、extension overlay、auth 和 capability discovery 投影成每个 Session/Step 所需的 McpConfig/McpRuntimeProjection。

live connections 由 Session 的 McpRuntime 拥有。manager 级 McpManager 是配置/目录工厂,Session 级 McpRuntime 才是连接 owner。

8. 创建竞态 ​

刷新所有 loaded thread 的 MCP runtime 很容易遗漏“正在 spawn、尚未插入 active map”的 Session。 starting_mcp_runtimes 专门覆盖这个窗口:

  1. spawn 创建 Arc<AtomicBool> 的 source_changed_during_startup;
  2. manager 把它的 Weak 放进同步 Mutex vector;
  3. invalidate_mcp_runtimes() 先 upgrade 所有 weak marker 并写 true,再遍历 active map 请求 refresh;
  4. Session 完成注册后检查 marker,若已变为 true,立即 request MCP refresh;
  5. 下次 invalidation 顺便清理无法 upgrade 的 marker。

marker 使用 Weak,不会因为 manager 的辅助列表让失败/取消的 Session 永久存活;使用 std::sync::Mutex 是因为临界区只做 retain/store,不跨 async await。

源码位置:codex-rs/core/src/thread_manager.rs :: invalidate_mcp_runtimes

rust
pub async fn invalidate_mcp_runtimes(&self) {
    // 先覆盖尚未进入 active map 的启动窗口。
    self.invalidate_starting_mcp_runtimes();
    let threads = self
        .state
        .threads
        .read()
        .await
        .values()
        .cloned()
        .collect::<Vec<_>>();
    // 释放 map 读锁后再逐个请求 refresh,避免锁跨调用链。
    for thread in threads {
        thread.session.request_mcp_runtime_refresh();
    }
}

fn invalidate_starting_mcp_runtimes(&self) {
    let mut starting = self
        .state
        .starting_mcp_runtimes
        .lock()
        .unwrap_or_else(std::sync::PoisonError::into_inner);
    starting.retain(|runtime| {
        let Some(runtime) = runtime.upgrade() else {
            return false;
        };
        runtime.store(true, Ordering::Release);
        true
    });
}

这段实现同时给出两个并发约束:active map 只在 clone handle 时持锁,启动中 marker 则用 Release 发布变化。Weak upgrade 失败的项会在同一次 invalidation 中清除。

当前 manager 还提供 refresh_hook_runtimes():它先复制 loaded thread handles,随后逐个读取 Session 配置并调用 Session::refresh_hooks。这与 invalidate_mcp_runtimes() 不同:前者只重建已加载 Thread 的 Hook runtime,后者同时覆盖启动中的 MCP marker 和已加载 Session。新增全局刷新入口时,必须明确它 覆盖的是 loaded、starting 还是两者。

9. Code Mode ​

构造器根据 CodeModeHost feature 或 disable_in_process_fallback 选择:

  • ProcessOwnedCodeModeSessionProvider;
  • DisabledCodeModeSessionProvider。

provider 保存于 manager,spawn 时 clone 给每个 Session 的 CodeModeService。这样 standalone host 进程和 session 创建策略在同一个 product manager 下保持一致。

with_code_mode_session_provider() 允许 App Server 注入自定义 provider,但要求 State 尚未共享。 测试还有 with_code_mode_host_program_for_tests(),用指定 binary 验证 host 路径。它们是构造期 strategy 替换,不是运行中 feature toggle。

10. Extension ​

10.1 extensions ​

Arc<ExtensionRegistry<Config>> 由宿主装配。ThreadManager 不知道具体产品要安装 Goal、Guardian、 Memories、MCP、Web Search、Image Generation 或 Skills 哪些 contributor;它只在 spawn 时把 registry 传入 Session。

App Server 使用 Arc::new_cyclic() 创建 manager,因为 extension 的 agent spawner/event sink 需要 Weak<ThreadManager>,同时 manager 又持有最终 registry。Weak 边消除了构造环的强所有权问题。

10.2 用户指令Provider ​

该 trait provider 负责进程/宿主级用户指令。fresh root、cold resume 和 root fork 通常加载新 snapshot; subagent 优先从 live parent Session 继承,避免每个 child 独立读取并漂移。

ThreadManager 只决定何时 load/inherit,并把 LoadedUserInstructions 交给 Session;AGENTS.md、Skill、 environment context 等后续拼装不属于这个 provider。

11. ThreadStore ​

11.1 thread_store ​

Arc<dyn ThreadStore> 是 process-scoped persistence backend。App Server 构造时明确固定它:配置 reload 可以改变 per-thread 行为,但不能把新建/恢复/fork thread 悄悄迁到另一个 root/backend。

thread_store_from_config() 当前选择:

  • LocalThreadStore:使用 codex_home、SQLite state,可选启动 rollout compression worker;
  • InMemoryThreadStore:按配置 ID 取得隔离的内存 store。

ThreadManager 用 store 读取 cold history、创建/fork/resume Session、更新 unloaded metadata;loaded metadata 更新则通过 CodexThread/LiveThread,保持与 live rollout writer 的顺序。

源码位置:codex-rs/core/src/thread_manager.rs :: thread_store_from_config

rust
pub fn thread_store_from_config(
    config: &Config,
    state_db: Option<StateDbHandle>,
) -> Arc<dyn ThreadStore> {
    // backend 在 manager 构造期选定,已有 Session 不随配置刷新迁移。
    match &config.experimental_thread_store {
        ThreadStoreConfig::Local => {
            // 压缩 worker 只属于 local rollout backend,并受独立 feature gate 控制。
            if config.features.enabled(Feature::LocalThreadStoreCompression) {
                codex_rollout::spawn_rollout_compression_worker(
                    config.codex_home.to_path_buf(),
                );
            }
            Arc::new(LocalThreadStore::new(
                LocalThreadStoreConfig::from_config(config),
                state_db,
            ))
        }
        ThreadStoreConfig::InMemory { id } => InMemoryThreadStore::for_id(id),
    }
}

backend 在 manager 构造时被收敛为 Arc<dyn ThreadStore>;后续 Session 只 clone trait object,不再读取 选择配置。压缩 worker 也只随 local backend 和 feature 同时满足时启动。

11.2 AgentGraphStore ​

Option<Arc<dyn AgentGraphStore>> 保存跨加载状态的 parent/descendant agent graph。helper 只在 state DB 存在时创建 LocalAgentGraphStore。

列出 agent subtree 时,manager 合并 persisted graph descendants 和 AgentControl 的 live subtree,再用 set 去重。ThreadStore 回答“thread 内容是什么”,AgentGraphStore 回答“agent thread 之间是什么 关系”,两者不能合并成一个 repository。

manager 对持久化请求会先判断 live/cold/ephemeral,再选择不同 owner。下面的流程图展示 metadata、 history 与 agent graph 为什么不能全部绕过 live handle 直接写 store。

loaded metadata 通过 LiveThread 保持与 rollout writer 的顺序,cold metadata 才直接进入 store;agent subtree 则需要把持久化边与仍在内存中的 child 一起合并。

12. 宿主Provider ​

字段用途缺失时语义
attestation_providerApp Server/desktop 在上游请求前生成 attestation不提供该宿主能力
external_time_provider允许宿主提供当前时间/时区来源使用 Core 默认 provider
analytics_events_clientthread/session/agent/plugin 产品 analytics不发送对应 analytics
session_source默认 SessionSource、restriction product、originator 语义构造器必填
installation_id跨 Thread 的安装身份、request metadata构造器必填

这些字段之所以放在 manager,而不是每次 StartThreadOptions 重复传入,是因为它们属于一个宿主进程或 manager 实例的身份。StartThreadOptions 可以覆盖 session/thread source 等少量 per-thread 语义,但 attestation implementation、installation ID 和默认 time provider 不应逐 Turn 改变。

13. 测试辅助状态 ​

ops_log 仅在全局 test mode 打开时存在,用 std::sync::Mutex<Vec<(ThreadId, Op)>> 捕获 ThreadManagerState::send_op(),用于断言 multi-agent/parent-child 控制消息。生产默认是 None。

测试构造器还会:

  • 创建 dummy AuthManager 和指定 ModelProviderInfo;
  • 创建临时 codex_home 与 guard;
  • 默认 local ThreadStore、可选 state DB/graph;
  • 使用 empty extension registry 与 disabled Code Mode provider;
  • 创建 test EnvironmentManager;
  • 不配置 attestation、external time 或 analytics。

因此测试 manager 不是生产构造的完全镜像。测试某个 extension、remote environment、非 local store 或 attestation 时,应调用生产 ThreadManager::new() 显式注入依赖,而不是继续扩张便捷测试构造器。

14. Manager依赖 ​

层次保存什么何时确定
ThreadManagerprocess-scoped manager、store、registry、host provider产品启动
StartThreadOptionsconfig、history、source、dynamic tools、environment selections每次 start/resume/fork
ThreadSpawnRequestauth/control/lineage/inherited env与policy等内部参数spawn 入口归一化
SessionServicesmanager 依赖的 Arc clone + thread-scoped runtimeSession::new
TurnContextmodel、permission、mode、environment snapshot每个 Turn
StepContextMCP binding、tool plan、ready capability每次 sampling

把依赖放错层会造成不同问题:把 tool catalog 固定到 manager 会阻止 Step refresh;把 ThreadStore 放进 StartThreadOptions 会允许同一进程配置 reload 后漂移 backend;把 auth token 值复制进 Session 会错过 refresh;把整个 Config 每次从全局读取则破坏 Turn 快照一致性。

15. 所有权测试 ​

15.1 Code Mode ​

字段类型是 Arc<dyn CodeModeSessionProvider> 只能说明“可以共享”,不能说明 spawn 时没有为每个 Thread 重建 provider。测试显式注入一个 provider,启动两个 Thread,再从两个 Session 的 CodeModeService 取回 provider并做指针身份比较。

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

rust
// :: code_mode_session_provider_is_shared_across_threads(断言节选)
let provider: Arc<dyn CodeModeSessionProvider> =
    Arc::new(DisabledCodeModeSessionProvider);
let manager = ThreadManager::with_models_provider_and_home_for_tests(
    CodexAuth::from_api_key("dummy"),
    config.model_provider.clone(),
    config.codex_home.to_path_buf(),
    Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()),
)
.with_code_mode_session_provider(Arc::clone(&provider));
let first = manager
    .start_thread(StartThreadOptions::new(config.clone()))
    .await
    .expect("start first thread");
let second = manager
    .start_thread(StartThreadOptions::new(config))
    .await
    .expect("start second thread");

let first_provider = first
    .thread
    .session
    .services
    .code_mode_service
    .session_provider();
let second_provider = second
    .thread
    .session
    .services
    .code_mode_service
    .session_provider();
// 三次ptr_eq同时证明两个Session和manager持有同一个provider对象。
assert!(Arc::ptr_eq(&first_provider, &second_provider));
assert!(Arc::ptr_eq(&first_provider, &provider));
assert!(Arc::ptr_eq(
    &first_provider,
    &manager.state.code_mode_session_provider,
));

该测试证明 provider identity 和 manager→Session clone 边,不证明 provider 内部创建的 Code Mode session 也跨 Thread 共享;后者由 provider 实现决定。

15.2 Weak Marker竞态 ​

测试注册一个同时贡献 Thread lifecycle 和 MCP projection 的 extension。on_thread_start() 用 Notify 卡住 Session 构造,使 Thread 尚未发布到 registry;此时调用 manager invalidation,再释放启动。关键断言有两 层:invalidation 时公开 live list仍为空,Session 发布后 MCP contributor必须观察到第二次 projection。

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

rust
// :: 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");
// 此刻active map没有Thread,普通map遍历无法覆盖这次invalidation。
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
    // 第二次projection证明Weak marker在发布后转化成真实refresh。
    .expect("invalidation during startup should refresh the newly published thread");

这不是依赖固定 sleep 的时序猜测:两个 Notify 精确控制“进入 startup callback”和“允许继续发布”。测试 说明 marker 不会遗漏这次变化,但没有直接检查 Weak 条目何时被清理;清理语义仍由 retain + upgrade 源码保证。

15.3 Registry ​

internal Thread必须进入 manager map,才能参与 shutdown和内部控制;但普通 list/get不能暴露它。测试启动 SessionSource::Internal(MemoryConsolidation),断言公共 list为空、get返回错误,随后批量shutdown报告中 又必须包含该 ThreadId。

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

rust
// :: start_thread_keeps_internal_threads_hidden_from_normal_lookups(断言节选)
let thread = manager
    .start_thread(StartThreadOptions {
        session_source: Some(SessionSource::Internal(
            InternalSessionSource::MemoryConsolidation,
        )),
        environments: Some(Vec::new()),
        ..StartThreadOptions::new(config)
    })
    .await
    .expect("internal thread should start");

// 对外查询隐藏internal Thread,不能由此推断它未被manager拥有。
assert_eq!(manager.list_thread_ids().await, Vec::new());
assert!(manager.get_thread(thread.thread_id).await.is_err());

let report = manager
    .shutdown_all_threads_bounded(Duration::from_secs(10))
    .await;
// shutdown遍历内部registry,因此仍能找到并关闭隐藏Thread。
assert_eq!(report.completed, vec![thread.thread_id]);
assert!(report.submit_failed.is_empty());
assert!(report.timed_out.is_empty());

这组断言直接区分了“map ownership”和“public visibility”。它没有验证冷历史;ThreadStore 中未加载 Thread 的读取和列表属于持久化系列。

16. 字段设计问题 ​

新增 ThreadManagerState 字段前,应先回答:

  1. 它是否真的跨多个 Thread 共享?如果只属于一个 Session,应放 SessionServices;
  2. 它是否需要 runtime refresh?若需要,现有 Session 如何收到 invalidation;
  3. 是否可能形成 State → Session → State 强引用环;
  4. 构造器是否必须由所有宿主提供,还是可以从现有输入派生;
  5. test constructor、thread-manager-sample 和 App Server 的生产装配是否都要更新;
  6. dependency 是否为 trait object,能否让低层 crate 保持独立;
  7. shutdown 由 manager、Session 还是 provider 自己负责。

验证入口包括 thread_manager_tests.rs 的 registration/shutdown/invalidation,Core integration tests 的 start/resume/subagent,App Server message processor 的生产构造,以及 thread-manager sample 的 facade 装配。字段“能编译”只说明构造参数齐全;还必须确认其生命周期、refresh 和 shutdown owner 明确。

17. 依赖发布验证 ​

  1. 从 ThreadManager { state: Arc<_> } 开始,解释为什么 threads map、ThreadStore 和 thread_created_tx 分别是 live authority、durable authority 和增量提示,三者不能相互替代。
  2. 给定“Plugin更新发生在 Session还没有注册进map时”,复述 Weak AtomicBool marker如何把变化带到发布后的 Session,并指出 Release store与后续读取之间的作用。
  3. 给定“list看不到Thread,但shutdown report包含它”,定位 session_source.is_internal() 的过滤位置,说明 这是可见性规则而不是registry损坏。
  4. 使用只读命令核对本文三组owner和测试:
bash
rg -n "struct ThreadManagerState|starting_mcp_runtimes|invalidate_mcp_runtimes" \
  codex-rs/core/src/thread_manager.rs
rg -n "code_mode_session_provider_is_shared|mcp_invalidation_refreshes|internal_threads_hidden" \
  codex-rs/core/src/thread_manager_tests.rs

继续学习时,阅读 ThreadManager创建,观察本文的共享依赖如何 进入一次创建事务;再阅读 ThreadManager管理,理解 live registry 在 remove、shutdown和cold metadata操作中的权威边界。