ThreadManager 不是一个只会 HashMap::insert 的目录。它同时承担 live Thread 注册、冷存储路由、共享依赖注入和有界批量关闭。理解它的关键是区分“表中可查”“对象仍存活”和“磁盘上存在”三个状态。
注册表只拥有进程内可服务的 Thread
内部状态通过异步 RwLock 保存 ThreadId -> Arc<CodexThread>。读操作可以并发,注册、移除和批量清理需要写锁;真正的 Session、Channel 和任务由 Arc<CodexThread> 间接保持。
ThreadManager 自身还持有模型目录、认证、环境、MCP、Skills、Plugins、网络代理和 ThreadStore 等共享服务。创建 Session 时把这些依赖组装进去,使不同 Thread 共享管理器而不共享自己的 SessionState。
class ThreadManagerState:
live: RWLock[dict[ThreadId, Arc[CodexThread]]]
thread_store: ThreadStore
environments: EnvironmentManager
models: ModelsManager
mcp: McpManager
skills: SkillsManager
async def get_thread(id):
thread = await live.read().get(id)
if thread is None or thread.source.is_internal():
raise ThreadNotFound(id)
return thread.clone_arc()
对外查询会隐藏 internal session。它们可以为 Review、Guardian 或其他内部流程服务,但不应出现在普通客户端的线程选择器中。隐藏是在访问边界执行,不是让这些对象脱离注册表。
为什么首事件到达后才登记
Session spawn 会先建立通道和输入循环。ThreadManager 随后等待第一条事件,并要求它是使用初始 submission id 发出的 SessionConfigured。只有初始化成功且事件顺序正确,才构造 CodexThread 并尝试写入表。
async def finalize_spawn(session, io, source):
first = await io.next_event()
require(first.id == INITIAL_SUBMIT_ID)
require(first.msg.kind == "SessionConfigured")
candidate = CodexThread(session, io, source, first.msg)
async with live.write():
if candidate.thread_id in live:
await candidate.shutdown_and_wait()
raise DuplicateThread(candidate.thread_id)
live[candidate.thread_id] = candidate
return candidate
这是一道发布屏障:调用方不会拿到半初始化 Thread。若同一个 ThreadId 已占用,不能简单覆盖旧 Arc,否则旧客户端仍持有原对象,新查询却进入另一个 Session,两套 Runtime 会同时写同一 Rollout。实现关闭失败候选并返回错误。
remove 不等于销毁,冷 Thread 也不等于不存在
remove_thread 只从 map 移除并返回 Arc。其他调用者若仍持有克隆,Session 继续运行。需要停止时必须显式调用 shutdown 协议。
元数据更新也根据冷热状态分流:live Thread 经 CodexThread/LiveThread 写入,以保持与当前事件流的顺序;未加载 Thread 才直接交给 ThreadStore。这防止客户端刚改名、活动 Session 随后又用旧快照覆盖。
有界批量关闭为什么保留失败项
进程退出或服务重载时,ThreadManager 先在读锁下复制当前 Arc,再并发发 Shutdown,每个等待受统一 timeout 约束。完成的 Thread 才从表移除;提交失败或超时项保留并进入报告。
async def shutdown_all(timeout):
snapshot = list((await live.read()).items())
outcomes = await bounded_join(
[thread.shutdown_and_wait() for _, thread in snapshot],
timeout=timeout,
)
async with live.write():
for thread_id, outcome in outcomes:
if outcome.completed and live.get(thread_id) is snapshot[thread_id]:
del live[thread_id]
return summarize(outcomes)
保留失败项不是泄漏,而是避免“表看起来已清空、后台任务实际仍活着”的假成功。使用快照还能避免关闭期间长时间持有写锁;最终删除时应确认表中仍是原 Arc,防止误删并发恢复出的新对象。
源码与测试锚点
codex-rs/core/src/thread_manager.rs:注册表、spawn finalization、冷热路由和批量关闭。codex-rs/core/src/thread_manager_tests.rs:恢复、重复注册和线程管理行为。codex-rs/core/src/live_thread.rs:live 元数据与顺序化写入接口。
评论
登录后即可评论