ThreadManager 与 CodexThread · 线程的生与死

线程在 app-server 里是一段可订阅的对话,在 codex-core 里是 ThreadManager 注册表中的一个 CodexThread:一端是提交队列,一端是事件队列。新建、恢复、分叉都汇入同一个 spawn_thread,差别只在初始历史;卸载、归档、删除则先关掉运行时,再交给 thread-store 处理持久化数据。

作者 David更新于 第 9 篇(共 57 篇)

ThreadManager 与 CodexThread · 线程的生与死

v2 协议里的 thread 只是一份数据;这一篇看它在内存里怎么活着。涉及三层:app-server 的 ThreadRequestProcessor 负责请求与订阅,codex-core 的 ThreadManager 负责创建和登记线程,codex-thread-store 负责线程的持久化边界。

用户看到的样子

codex resume、codex fork、codex archive、codex delete 与会话内的 /resume、/fork、/archive 等命令,背后都是 app-server 的 thread/resume、thread/fork、thread/archive、thread/delete;新会话则是 thread/start。命令参数与会话文件的位置见手册会话管理。

三层分工

图表加载中…

ThreadManager(codex-rs/core/src/thread_manager.rs:238)只是一层 Arc<ThreadManagerState>,真正的状态里有一张 HashMap<ThreadId, Arc<CodexThread>>,外加模型目录、MCP、技能、插件、扩展注册表、线程存储等所有线程共享的服务。app-server 在 MessageProcessor::new 里只建一个 ThreadManager,整个进程里的线程都挂在它上面。

CodexThread:一对队列

CodexThread(codex-rs/core/src/codex_thread.rs:186)持有运行时 Session 和一组队列端点 SessionIo:

/// Queue and lifecycle endpoints for a running [`Session`].
///
/// Runtime state lives on `Session`; keeping these endpoints separate lets all
/// submission senders be dropped to terminate the session loop. The shared
/// completion future observes that shutdown.
#[derive(Clone)]
pub(crate) struct SessionIo {
    pub(crate) tx_sub: Sender<Submission>,
    pub(crate) rx_event: Receiver<Event>,
    // Last known status of the agent.
    pub(crate) agent_status: watch::Receiver<AgentStatus>,
    // Shared future for the background submission loop completion so multiple
    // callers can wait for shutdown.
    pub(crate) session_loop_termination: SessionLoopTermination,
}

(codex-rs/core/src/session/mod.rs:397)

这就是 codex-rs/docs/protocol_v1.md 说的提交队列与事件队列:提交队列有界,容量 512(SUBMISSION_CHANNEL_CAPACITY),事件队列无界。submit(Op) 返回提交 id,它是 UUIDv7,注释说 app-server 把开启轮次的那个提交 id 直接当作公开的 turn id;next_event() 取下一个事件,第 6 篇的监听任务就在循环调用它;shutdown_and_wait() 提交 Op::Shutdown,再等会话循环结束。

生:start、resume、fork 殊途同归

三种“生”法最后都调用 ThreadManagerState::spawn_thread,区别只在带进去的初始历史:

#[derive(Debug, Clone, Deserialize, Serialize)]
pub enum InitialHistory {
    New,
    Cleared,
    Resumed(ResumedHistory),
    Forked(Vec<RolloutItem>),
}

(codex-rs/history/src/lib.rs:365)

thread/start 用 New,/clear 这类“清空重开”用 Cleared;冷恢复从存储读出历史,包成 Resumed;分叉把源线程的历史变成 Forked。spawn_thread 先处理一个捷径:要恢复的线程已经在注册表里且正在运行,就直接返回现有的 CodexThread,不再新建(正在工作的 Guardian 审查线程除外,它属于父线程)。否则它解析指令、多 agent 版本、originator 等,调用 Session::spawn 起会话,再由 finalize_thread_spawn 登记:

        let thread_id = session.thread_id();
        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,
            _ => {
                return Err(CodexErr::SessionConfiguredNotFirstEvent);
            }
        };
// ...
                let thread = Arc::new(CodexThread::new(
                    session,
                    io,
                    ThreadStartupMetadata::from(&session_configured),
                    session_configured.rollout_path.clone(),
                    session_source,
                ));
                e.insert(thread.clone());

(codex-rs/core/src/thread_manager.rs:2364)

会话的第一个事件必须是提交 id 为空串(INITIAL_SUBMIT_ID)的 SessionConfigured,否则视为启动失败;同一个 id 已被别的运行时占用时,新会话会被关掉并报 already running。持久化也在 Session::spawn 里按初始历史分流:临时(ephemeral)线程根本不建 LiveThread,New、Cleared、Forked 调 LiveThread::create,Resumed 调 LiveThread::resume 重新打开原来的记录继续追加。

app-server 这一侧还有几件事:

  • thread/start:先预留 thread id,非临时线程还把 projectId 等宿主元数据暂存到线程存储;在真正启动之前就把发起请求的连接登记为订阅者,注释说启动时的预热握手可能需要向它请求 attestation;启动后挂上监听任务,先回响应,再广播 thread/started。
  • thread/resume:线程已在内存里,就是“热恢复”:请求交给该线程的监听任务处理,把这个连接加进订阅,用运行中的状态组装响应(进行到一半的当前轮也合并进去),并补发尚未答复的服务端请求;否则“冷恢复”:从线程存储读出历史,调用 ThreadManager::resume_thread_with_history。线程正在卸载时,恢复请求会被拒绝,让客户端稍后重试。
  • thread/fork:旧格式的源线程读出完整历史,走 fork_thread_from_history;分页格式的源线程先由线程存储 prepare_fork(分叉点可以是最新、截至某一轮,或某一轮之前),再走 fork_prepared_thread。后者的持久化方式是 ForkPersistence::Referenced,新线程引用源线程的历史,而不是复制一份。两条路都用 ForkSnapshot::Interrupted:源线程若正停在一轮中间,就在分叉历史末尾补上与真实中断相同的 <turn_aborted> 标记和 TurnAborted 事件,让新线程从一个完整的边界开始。

死:卸载、归档、删除

  • 卸载:线程没有订阅者、又不处于活动态,持续 thread_unload_delay_secs(默认 60 秒)后,监听任务取消该线程所有待答复的服务端请求,调 shutdown_and_wait(最多等 10 秒),从注册表移除并广播 thread/closed。移除用的是 remove_thread_if_matches,只在注册表里仍是同一个运行时时才删,防止误删 thread/revert 之后换上的新运行时。卸载不碰磁盘上的数据,之后可以随时冷恢复。
  • 归档:先经 list_agent_subtree_thread_ids 找出这个线程派生的整棵子代理树,逐个关掉仍在运行的运行时,再调 ThreadStore::archive_threads;按 trait 的约定,第一个(根线程)必须成功,其余尽力而为。
  • 删除:同样处理整棵子树,先关运行时并刷新日志库,再调 delete_threads,顺序是子线程在前、根线程在后,最后为每个被删的线程广播 thread/deleted。还在内存里的临时线程没有可删的持久化数据,会被直接拒绝。

归档和删除关运行时都走 ThreadManager::remove_thread_for_client,它拒绝移除正在运行的内部线程,错误信息是 live internal threads can only be removed by their owner;app-server 的 README 以 Guardian 审查线程为例说明,这类工作线程要由父会话自己收尾。

thread-store:线程的持久化边界

codex-thread-store 的文档开头定下原则:应用代码只把 ThreadId 当作线程的持久句柄,由实现决定它对应本地 rollout 文件、远程请求还是别的存储。ThreadStore trait(codex-rs/thread-store/src/store.rs:93)自称“storage-neutral thread persistence boundary”,方法覆盖创建、恢复、追加、落盘、读取、列表、归档、删除、分叉准备与回退。落盘请求带着 PersistContext 说明原因,其中 SubagentSpawn、TurnStart、SteeredUserInput 允许存储先排队、在后续的屏障处再确保写入,ThreadPreparation 与 Standard 必须同步完成。

活动线程并不直接调 trait,而是拿一个 LiveThread 句柄:注释说它把生命周期决策留给调用方、把存储细节交给 ThreadStore,会话代码只需要这一个句柄。配套的 LiveThreadInitGuard 负责初始化失败时丢弃已经打开的写入器,而不把内存里的半成品强行落盘。

0.158.0 里 ThreadStore 有两个实现:默认的 LocalThreadStore(rollout JSONL 文件加 SQLite 元数据,见下一篇),以及只供测试与调试配置使用的 InMemoryThreadStore;旧配置项 experimental_thread_store_endpoint 如今会直接报错“no longer supported”。选哪个在进程启动时就定了,MessageProcessor::new 的注释解释说,配置重载可以改变单个线程的行为,但不能把新建、恢复或分叉的线程挪到另一个持久化后端。

和《从 LLM 到 Coding Agent》对照

那本书的 session-history 把恢复会话写成“把 JSONL 一行行读回来、重建数组、接着循环”。Codex 的冷恢复本质相同,只是多了一层注册表:同一个线程在进程里只能有一个运行时,重复恢复直接复用,关闭则要等会话循环真正结束。分叉的做法与 Grok Build 的 Sessions 不同:Grok 把会话数据复制成新分支,Codex 的分页线程可以只引用源线程的历史,旧格式线程才整份复制。


上一篇:守护进程与传输 · 一个服务端,多种连法 · 下一篇:会话持久化 · rollout、thread-store 与 SQLite

本页目录