Thread归档与删除
一个对话从列表中消失,至少可能发生了三件不同的事:它的运行对象被卸载;持久日志移到了归档区;日志和关联状态被删除。只盯着界面上的一行,很难判断它还能不能恢复、正在执行的工具是否已经退出,以及引用它的另一个分支是否仍然可读。
本文面向掌握 Rust async、Arc、锁和 Result 的读者。建议先读ThreadResume处理流程,理解持久记录如何重新建立执行者;再读ThreadFork处理流程,理解分支为什么可能依赖源日志。这里沿三个公开 API 追踪运行资源、文件和数据库的变化,最后能解释“恢复归档后仍是 NotLoaded”“归档成功但一个子线程还在原目录”“删除因另一个分支而失败”这三种现象。
Thread ID 是逻辑对话身份,rollout 是持久记录文件,运行对象是 Core 中的 CodexThread 及其 Session;它们不是一一同生共死的对象。下文的文件移动和文件锁特指 LocalThreadStore。ThreadStore 是存储抽象,其他后端可以返回无本地 path 的摘要,不能把本地目录布局写成全部后端的 API 契约。
1. 三种资源操作
1.1 请求与结果
三个方法只接收 Thread ID,没有“强制忽略引用”或“递归恢复归档”参数。
源码文件:codex-rs/app-server-protocol/src/protocol/common.rs
相关函数/类型:ClientRequest 生命周期方法;行号:532-541,652-656。
ThreadArchive => "thread/archive" {
params: v2::ThreadArchiveParams,
serialization: thread_id(params.thread_id),
response: v2::ThreadArchiveResponse,
},
ThreadDelete => "thread/delete" {
params: v2::ThreadDeleteParams,
serialization: thread_id(params.thread_id),
response: v2::ThreadDeleteResponse,
},
// ...
ThreadUnarchive => "thread/unarchive" {
params: v2::ThreadUnarchiveParams,
serialization: thread_id(params.thread_id),
response: v2::ThreadUnarchiveResponse,
},serialization: thread_id(...) 为同一 ID 的请求指定排队范围。处理子树时还会进入处理器级 permit 和 Store 的锁;仅有根 ID 排队不能保护另一进程的 writer,也不能单独解决父子 ID 的交叉操作。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:归档、删除与恢复归档请求响应;行号:666-668,673-673,678-680,685-685,761-763,1106-1108。
pub struct ThreadArchiveParams {
pub thread_id: String,
}
// ...
pub struct ThreadArchiveResponse {}
// ...
pub struct ThreadDeleteParams {
pub thread_id: String,
}
// ...
pub struct ThreadDeleteResponse {}
// ...
pub struct ThreadUnarchiveParams {
pub thread_id: String,
}
// ...
pub struct ThreadUnarchiveResponse {
pub thread: Thread,
}thread/archive、thread/delete 成功响应都是空对象;thread/unarchive 返回一个 Thread 摘要。这些空对象不携带逐个后代的完整结果,后续通知及再次读取才是客户端观察集合变化的入口。Rust 字段 thread_id 在线上采用 threadId,方法名使用斜线。
公开请求由 MessageProcessor 分派给专属 Thread 处理器。
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数/类型:MessageProcessor::process_request;行号:1161-1170。
ClientRequest::ThreadArchive { params, .. } => {
self.thread_processor
.thread_archive(request_id.clone(), params)
.await
}
ClientRequest::ThreadDelete { params, .. } => {
self.thread_processor
.thread_delete(request_id.clone(), params)
.await
}收到请求不等于修改持久状态。下一节会看到,解析身份、发现子树、关闭运行资源和存储操作分属不同步骤,每一步都有不同失败出口。
1.2 状态坐标
不要将 Archived 塞进运行状态枚举后画成一个不可返回的终态。实际需要同时检查以下坐标。
| 坐标 | 主要所有者 | 归档的影响 | 恢复归档的影响 | 删除的影响 |
|---|---|---|---|---|
| 加载与事件订阅 | ThreadManager、ThreadStateManager | 清理运行对象和订阅 | 不主动创建运行对象 | 清理运行对象和订阅 |
| 文件位置 | LocalThreadStore | 移入 archived_sessions | 按文件名日期移回 sessions | 移除已发现的关联文件 |
archived_at 与路径 | StateRuntime | 更新为归档时间和新路径 | 清空归档时间、更新路径 | 最后删除主 Thread 行 |
| 历史依赖 | history_base | 引用仍须可解析 | 不改变分叉前缀 | 外部引用会阻止删除 |
下面的图只表达这些存储操作的入口和消费者,不表示一种额外的 Thread 枚举。
归档和删除共享运行资源清理,但持久结果不同;恢复归档主要改变存储坐标。图中没有 ThreadManager::resume_thread_with_history 调用,因此成功 unarchive 并不承诺模型已经可以直接接收下一轮输入。
2. 归档的集合
2.1 后代的来源
根线程之外,还要发现它通过 agent spawn 拥有的后代。处理器中的 state_db_spawn_subtree_thread_ids 最终委托给下面的 Core 方法,不能根据 helper 名字断言它只读取 SQLite。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:list_agent_subtree_thread_ids;行号:870-904。
pub async fn list_agent_subtree_thread_ids(
&self,
thread_id: ThreadId,
) -> CodexResult<Vec<ThreadId>> {
let mut subtree_thread_ids = Vec::new();
let mut seen_thread_ids = HashSet::new();
subtree_thread_ids.push(thread_id);
seen_thread_ids.insert(thread_id);
if let Some(agent_graph_store) = self.state.agent_graph_store() {
for descendant_id in agent_graph_store
.list_thread_spawn_descendants(thread_id, /*status_filter*/ None)
.await
.map_err(|err| {
CodexErr::Fatal(format!("failed to load thread-spawn descendants: {err}"))
})?
{
if seen_thread_ids.insert(descendant_id) {
subtree_thread_ids.push(descendant_id);
}
}
}
for descendant_id in self
.agent_control()
.list_live_agent_subtree_thread_ids(thread_id)
.await?
{
if seen_thread_ids.insert(descendant_id) {
subtree_thread_ids.push(descendant_id);
}
}
Ok(subtree_thread_ids)
}集合首先包含根;随后合并持久 agent graph 与当前 live agent 子树,并用 HashSet 去重。status_filter=None 使持久图查询不只选活动边。由用户 fork 得到的分支和 agent spawn 后代是两种关系:分叉来源不自动成为这个删除或归档集合的成员。
这里的 Vec 保留发现顺序,后文会反转后代部分。不要未经核对查询实现就把任意后代列表宣称为严格的拓扑排序;当前处理器直接依赖这份发现结果。
2.2 候选与锁定
归档先读取根,再尽力读取后代的存储摘要。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_archive_response 的候选集合;行号:1624-1675。
async fn thread_archive_response(
&self,
params: ThreadArchiveParams,
) -> Result<(ThreadArchiveResponse, Vec<String>), JSONRPCErrorError> {
let thread_id = ThreadId::from_string(¶ms.thread_id)
.map_err(|err| invalid_request(format!("invalid session id: {err}")))?;
let subtree_thread_ids = self.state_db_spawn_subtree_thread_ids(thread_id).await?;
let mut archive_thread_ids = Vec::new();
match self
.thread_store
.read_thread(StoreReadThreadParams {
thread_id,
include_archived: false,
include_history: false,
})
.await
{
Ok(thread) => {
if thread.archived_at.is_none() {
archive_thread_ids.push(thread_id);
}
}
Err(err) => return Err(thread_store_mutation_error("archive", err)),
}
for descendant_thread_id in subtree_thread_ids.iter().copied().skip(1) {
match self
.thread_store
.read_thread(StoreReadThreadParams {
thread_id: descendant_thread_id,
include_archived: true,
include_history: false,
})
.await
{
Ok(thread) => {
if thread.archived_at.is_none() {
archive_thread_ids.push(descendant_thread_id);
}
}
Err(err) => {
warn!(
"failed to read spawned descendant thread {descendant_thread_id} while archiving {thread_id}: {err}"
);
}
}
}
if archive_thread_ids.is_empty() {
return Ok((ThreadArchiveResponse {}, Vec::new()));
}根读取使用 include_archived=false,失败直接结束。后代读取使用 include_archived=true,只把尚未归档的成员加入移动候选;某个后代读取失败只记录 warning。于是有两个集合:需要移动文件的 archive_thread_ids,以及整个 spawn 子树 subtree_thread_ids。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_archive_response 的关闭与批处理;行号:1677-1696。
archive_thread_ids[1..].reverse();
// Collaboration may resume an archived descendant without unarchiving it.
self.prepare_thread_for_archive(thread_id).await;
for &descendant_thread_id in subtree_thread_ids.iter().skip(1).rev() {
self.prepare_thread_for_archive(descendant_thread_id).await;
}
let archived_thread_ids = self
.thread_store
.archive_threads(StoreArchiveThreadsParams {
thread_ids: archive_thread_ids,
writer_lock_thread_ids: subtree_thread_ids,
})
.await
.map_err(|err| thread_store_mutation_error("archive", err))?
.into_iter()
.map(|thread_id| thread_id.to_string())
.collect();
Ok((ThreadArchiveResponse {}, archived_thread_ids))
}移动列表保持根在第一个位置,将后代部分反转。关闭运行对象时则先根、再逆序后代;传给 Store 的 writer_lock_thread_ids 仍包含整个子树。理由写在原注释里:协作路径可能恢复一个已经归档的后代而不先 unarchive。即使该后代不需要再次移动,也必须检查它的运行与写入所有权。
thread_archive_rejects_owned_unmaterialized_paginated_descendant 用另一个进程持有尚未物化的后代 writer,要求根归档失败。只锁那些扫描到文件的成员,会漏掉这种“有执行者、还没有文件”的对象。
3. 关闭与互斥
3.1 处理器串行化
归档内层在调用上述集合处理之前取得共享的列表状态 permit。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_archive_inner;行号:1616-1622。
async fn thread_archive_inner(
&self,
params: ThreadArchiveParams,
) -> Result<(ThreadArchiveResponse, Vec<String>), JSONRPCErrorError> {
let _thread_list_state_permit = self.acquire_thread_list_state_permit().await?;
self.thread_archive_response(params).await
}源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:acquire_thread_list_state_permit;行号:948-957。
pub(super) async fn acquire_thread_list_state_permit(
&self,
) -> Result<SemaphorePermit<'_>, JSONRPCErrorError> {
self.thread_list_state_permit
.acquire()
.await
.map_err(|err| {
internal_error(format!("failed to acquire thread list state permit: {err}"))
})
}permit 的生命周期覆盖业务修改,外层发送响应前已经释放。删除和恢复归档也使用同一处理器字段;这是服务器内的业务协调,无法代替 Store 跨进程保护。若 semaphore 已关闭,错误在开始变更前返回。
3.2 关闭超时
运行对象清理最容易误读的一行是 remove_thread:它发生在等待关闭之前。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:prepare_thread_for_removal;行号:1023-1040。
pub(super) async fn prepare_thread_for_removal(&self, thread_id: ThreadId, operation: &str) {
let removed_conversation = self.thread_manager.remove_thread(&thread_id).await;
if let Some(conversation) = removed_conversation {
info!("thread {thread_id} was active; shutting down");
match wait_for_thread_shutdown(&conversation).await {
ThreadShutdownResult::Complete => {}
ThreadShutdownResult::SubmitFailed => {
error!(
"failed to submit Shutdown to thread {thread_id}; proceeding with {operation}"
);
}
ThreadShutdownResult::TimedOut => {
warn!("thread {thread_id} shutdown timed out; proceeding with {operation}");
}
}
}
self.finalize_thread_teardown(thread_id).await;
}manager 已不再按 ID 返回该对象,但本函数仍持有 Arc<CodexThread> 并调用关闭。提交失败或超时都会记日志,然后继续 teardown 和后续存储操作。它与恢复文章中“关闭失败则保留缓存执行者”的替换策略不同,不能抽象成所有生命周期操作都相同的 shutdown 流程。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:wait_for_thread_shutdown;行号:439-445。
pub(super) async fn wait_for_thread_shutdown(thread: &Arc<CodexThread>) -> ThreadShutdownResult {
match tokio::time::timeout(Duration::from_secs(10), thread.shutdown_and_wait()).await {
Ok(Ok(())) => ThreadShutdownResult::Complete,
Ok(Err(_)) => ThreadShutdownResult::SubmitFailed,
Err(_) => ThreadShutdownResult::TimedOut,
}
}10 秒限制只包住 shutdown_and_wait。它不是整个归档 RPC 的超时,也没有为随后等待 lifecycle reservation 或文件 I/O 设定同样上限。超时意味着未确认 Core 已完整关闭;Store 仍可能通过 writer 检查拒绝操作,不能从“proceeding with archive”日志推断文件一定会被移动。
3.3 订阅与回调
结束阶段清除的不止 Core 注册项。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:finalize_thread_teardown;行号:978-989。
async fn finalize_thread_teardown(&self, thread_id: ThreadId) {
self.pending_thread_unloads.lock().await.remove(&thread_id);
self.outgoing
.cancel_requests_for_thread(thread_id, /*error*/ None)
.await;
self.thread_state_manager
.remove_thread_state(thread_id)
.await;
self.thread_watch_manager
.remove_thread(&thread_id.to_string())
.await;
}pending_thread_unloads 避免保留旧卸载任务标记;出站请求、ThreadState 和 watch 则分别承担待审批回调、listener/连接关系与状态观察。旧对象被移除后,不能让这些索引继续宣称它仍然存在。
源码文件:codex-rs/app-server/src/thread_state.rs
相关函数/类型:remove_thread_state;行号:445-471。
pub(crate) async fn remove_thread_state(&self, thread_id: ThreadId) {
let thread_state = {
let mut state = self.state.lock().await;
let thread_state = state
.threads
.remove(&thread_id)
.map(|thread_entry| thread_entry.state);
state.thread_ids_by_connection.retain(|_, thread_ids| {
thread_ids.remove(&thread_id);
!thread_ids.is_empty()
});
thread_state
};
self.unregister_listener_command_tx(thread_id);
if let Some(thread_state) = thread_state {
let mut thread_state = thread_state.lock().await;
tracing::debug!(
thread_id = %thread_id,
listener_generation = thread_state.listener_generation,
had_listener = thread_state.cancel_tx.is_some(),
had_active_turn = thread_state.active_turn_snapshot().is_some(),
"clearing thread listener during thread-state teardown"
);
thread_state.clear_listener();
}
}在 manager 锁下移走 ThreadState 并从每个 connection 的订阅集合删除 ID,之后才锁具体 ThreadState 清除 listener。将两个锁的工作拆开,避免持有全局索引锁等待单线程 listener 状态。clear_listener 取消旧监听,不能只删 Core 对象而让客户端继续等待旧事件流。
源码文件:codex-rs/app-server/src/outgoing_message.rs
相关函数/类型:cancel_requests_for_thread;行号:480-513。
pub(crate) async fn cancel_requests_for_thread(
&self,
thread_id: ThreadId,
error: Option<JSONRPCErrorError>,
) {
let entries = {
let mut request_id_to_callback = self.request_id_to_callback.lock().await;
let request_ids = request_id_to_callback
.iter()
.filter_map(|(request_id, entry)| {
(entry.thread_id == Some(thread_id)).then_some(request_id.clone())
})
.collect::<Vec<_>>();
let mut entries = Vec::with_capacity(request_ids.len());
for request_id in request_ids {
if let Some(entry) = request_id_to_callback.remove(&request_id) {
entries.push(entry);
}
}
entries
};
for entry in entries {
self.analytics_events_client
.track_server_request_aborted(now_unix_timestamp_ms(), entry.request.id().clone());
if let Some(error) = error.as_ref()
&& entry.callback.send(Err(error.clone())).is_err()
{
let request_id = entry.request.id();
warn!("could not notify callback for {request_id:?}: receiver dropped");
}
}
}回调表先在锁内筛选并移出条目,后续通知统计在锁外进行。这里传的是 error=None,因此不会为每个回调发送一个伪造 JSON-RPC 错误,而是在条目释放时关闭对应发送端。客户端曾经看到的审批卡片和服务器仍有可完成的原请求,是两件不同的事。
3.4 Store 锁序
本地批量归档在移动第一个文件前,对整个锁集合排序、去重并取得所有权。
源码文件:codex-rs/thread-store/src/local/archive_thread.rs
相关函数/类型:archive_threads 的三层互斥;行号:24-46。
let mut lock_thread_ids = params.writer_lock_thread_ids;
lock_thread_ids.extend(thread_ids.iter().copied());
lock_thread_ids.sort_unstable_by_key(ToString::to_string);
lock_thread_ids.dedup();
let mut _lifecycle_guards = Vec::with_capacity(lock_thread_ids.len());
for thread_id in &lock_thread_ids {
_lifecycle_guards.push(store.live_writer_locks.lock_lifecycle(*thread_id).await);
}
let mut _live_writer_guards = Vec::with_capacity(lock_thread_ids.len());
for thread_id in &lock_thread_ids {
_live_writer_guards.push(store.live_writer_locks.lock(*thread_id).await);
if store.live_recorders.lock().await.contains_key(thread_id) {
return Err(ThreadStoreError::Conflict {
message: format!("thread {thread_id} already has an active writer"),
});
}
}
let _writer_guards = store.acquire_writer_locks(&lock_thread_ids).await?;
let reference_index = RolloutReferenceIndex::scan(store.config.codex_home.as_path())
.await
.map_err(|err| ThreadStoreError::Internal {
message: format!("failed to scan thread rollout files: {err}"),
})?;先取所有 lifecycle 独占 guard,再取所有进程内 writer mutex,检查 live recorder,最后取得跨进程 writer 锁。排序保证两个批次不会仅因集合顺序相反而互相等锁;lifecycle 必须先于 writer,否则归档可能占着 writer 等 fork 释放 reservation,而 fork 又需要 writer 才能完成准备。
跨进程锁使用文件锁而非“看到 .lock 文件就算占用”。
源码文件:codex-rs/thread-store/src/local/writer_lock.rs
相关函数/类型:acquire;行号:39-87。
pub(super) fn acquire(
self: &Arc<Self>,
thread_id: ThreadId,
) -> ThreadStoreResult<WriterLockGuard> {
let coordination_lock = self.lock_coordination()?;
if !self.cleanup_attempted.swap(true, Ordering::Relaxed)
&& let Err(err) = self.remove_stale_thread_locks()
{
warn!("failed to clean up stale thread writer locks: {err}");
}
let path = self.directory.join(format!("{thread_id}.lock"));
let file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&path)
.map_err(|err| ThreadStoreError::Internal {
message: format!(
"failed to open thread writer lock {}: {err}",
path.display()
),
})?;
match file.try_lock() {
Ok(()) => {}
Err(std::fs::TryLockError::WouldBlock) => {
return Err(ThreadStoreError::Conflict {
message: format!("thread {thread_id} already has an active writer"),
});
}
Err(std::fs::TryLockError::Error(err)) => {
return Err(ThreadStoreError::Internal {
message: format!(
"failed to acquire thread writer lock {}: {err}",
path.display()
),
});
}
}
drop(coordination_lock);
Ok(WriterLockGuard {
coordinator: Arc::clone(self),
path,
file: Some(file),
})
}try_lock 遇到 WouldBlock 返回 Conflict,其他 OS 错误映射为内部错误;成功返回的 guard 持有文件句柄。旧锁文件是否存在不能替代锁是否仍被进程持有的判断,也不应通过手动删锁文件来测试活跃 writer。
下面的时序对应一个 fork 尚持有源 reservation 的情况。
图中归档等待时没有持有 writer。它只刻画取得锁的关系,不承诺源追加任意时长都能完成。以下单元测试直接对这个不变量设限。
源码文件:codex-rs/thread-store/src/local/archive_thread.rs
相关函数/类型:archive_waits_for_fork_reservation_without_holding_writer_lock;行号:173-204。
async fn archive_waits_for_fork_reservation_without_holding_writer_lock() {
let home = TempDir::new().expect("temp dir");
let store = LocalThreadStore::new(test_config(home.path()), /*state_db*/ None);
let uuid = Uuid::from_u128(205);
let thread_id = ThreadId::from_string(&uuid.to_string()).expect("valid thread id");
let active_path = write_session_file_with_history_mode(
home.path(),
"2025-01-03T12-00-00",
uuid,
ThreadHistoryMode::Paginated,
)
.expect("session file");
let reservation = store.live_writer_locks.reserve_lifecycle(thread_id).await;
let mut archive = Box::pin(store.archive_thread(ArchiveThreadParams { thread_id }));
tokio::select! {
biased;
result = &mut archive => panic!("archive completed while the source was reserved: {result:?}"),
_ = tokio::task::yield_now() => {}
}
let writer_guard = tokio::time::timeout(
Duration::from_secs(1),
store.live_writer_locks.lock(thread_id),
)
.await
.expect("pending archive should not hold the writer lock");
drop(writer_guard);
drop(reservation);
archive.await.expect("archive reserved thread");
assert!(!active_path.exists());
}测试先持有 reservation,轮询归档确认它尚未完成,再要求一秒内取得 writer mutex。只有释放 reservation 后归档才完成,原路径消失。它证明等待顺序,不能外推成所有文件系统上的一秒性能保证。
4. 文件移动与补偿
4.1 先选路径
一个逻辑 Thread 可能拥有多个物理 rollout,尤其发生过历史替换后。归档以 reference index 找到其拥有的文件,并额外保证当前选定路径进入候选。
源码文件:codex-rs/thread-store/src/local/archive_thread.rs
相关函数/类型:archive_thread_with_paths 的移动计划;行号:68-112。
let state_db_ctx = store.state_db().await;
let selected_rollout_path = thread_rollout_resolver::resolve_current(store, thread_id)
.await?
.map(|resolved| resolved.path)
.ok_or_else(|| ThreadStoreError::InvalidRequest {
message: format!("no rollout found for thread id {thread_id}"),
})?;
let archive_folder = store
.config
.codex_home
.join(codex_rollout::ARCHIVED_SESSIONS_SUBDIR);
std::fs::create_dir_all(&archive_folder).map_err(|err| ThreadStoreError::Internal {
message: format!("failed to archive thread: {err}"),
})?;
if !rollout_paths.contains(&selected_rollout_path) {
rollout_paths.push(selected_rollout_path.clone());
}
let mut archived_path = None;
let mut rollout_moves = Vec::new();
for rollout_path in rollout_paths {
if rollout_path_is_archived(store.config.codex_home.as_path(), rollout_path.as_path()) {
continue;
}
let canonical_rollout_path = scoped_rollout_path(
store.config.codex_home.join(codex_rollout::SESSIONS_SUBDIR),
rollout_path.as_path(),
"sessions",
)?;
let file_name =
validated_rollout_file_name(canonical_rollout_path.as_path(), rollout_path.as_path())?;
let destination = archive_folder.join(&file_name);
if rollout_path == selected_rollout_path {
archived_path = Some(destination.clone());
}
if !rollout_moves
.iter()
.any(|(source, _)| source == &canonical_rollout_path)
{
rollout_moves.push((canonical_rollout_path, destination));
}
}
let archived_path = archived_path.ok_or_else(|| ThreadStoreError::Internal {
message: format!("failed to archive selected rollout for thread {thread_id}"),
})?;这里归档的是所有权属于该 Thread 的文件,不能沿 history_base 把祖先的文件也一起搬走。canonical_rollout_path 用于去重,selected_rollout_path 用来确定主 metadata 应指向哪个归档目标。先完整构造 rollout_moves,再执行 rename,可让路径验证错误在移动之前暴露。
路径校验先解析真实目录和真实文件。
源码文件:codex-rs/thread-store/src/local/helpers.rs
相关函数/类型:scoped_rollout_path;行号:32-61。
pub(super) fn scoped_rollout_path(
root: PathBuf,
rollout_path: &Path,
root_name: &str,
) -> ThreadStoreResult<PathBuf> {
let canonical_root =
std::fs::canonicalize(&root).map_err(|err| ThreadStoreError::Internal {
message: format!(
"failed to resolve {root_name} directory `{}`: {err}",
root.display()
),
})?;
let canonical_rollout_path =
std::fs::canonicalize(rollout_path).map_err(|_| ThreadStoreError::InvalidRequest {
message: format!(
"rollout path `{}` must be in {root_name} directory",
rollout_path.display()
),
})?;
if canonical_rollout_path.starts_with(&canonical_root) {
Ok(canonical_rollout_path)
} else {
Err(ThreadStoreError::InvalidRequest {
message: format!(
"rollout path `{}` must be in {root_name} directory",
rollout_path.display()
),
})
}
}比较的是 canonical 路径,避免把词法上位于 sessions 的软链接误认成受管理文件。它不是通用 OS 沙箱,只约束这次 Store 操作选中的来源路径。
源码文件:codex-rs/thread-store/src/local/helpers.rs
相关函数/类型:validated_rollout_file_name;行号:93-115。
pub(super) fn validated_rollout_file_name(
rollout_path: &Path,
display_path: &Path,
) -> ThreadStoreResult<std::ffi::OsString> {
let Some(file_name) = rollout_path.file_name().map(OsStr::to_owned) else {
return Err(ThreadStoreError::InvalidRequest {
message: format!(
"rollout path `{}` missing file name",
display_path.display()
),
});
};
if codex_rollout::rollout_id_from_path(rollout_path).is_some() {
Ok(file_name)
} else {
Err(ThreadStoreError::InvalidRequest {
message: format!(
"rollout path `{}` has an invalid filename",
display_path.display()
),
})
}
}文件名还必须能解析出 rollout ID;有合法目录但文件名不合法,同样不能进入移动或删除。后续恢复目录依赖文件名的时间部分,因此路径格式本身参与生命周期协议。
4.2 补偿的范围
归档的提交顺序是 rename 文件,再更新状态数据库。
源码文件:codex-rs/thread-store/src/local/archive_thread.rs
相关函数/类型:archive_thread_with_paths 的补偿;行号:114-145。
for (index, (source, destination)) in rollout_moves.iter().enumerate() {
if let Err(err) = std::fs::rename(source, destination) {
if let Err(restore_err) = restore_rollout_moves(&rollout_moves[..index]) {
return Err(ThreadStoreError::Internal {
message: format!(
"failed to archive thread: {err}; failed to restore moved rollouts: {restore_err}"
),
});
}
return Err(ThreadStoreError::Internal {
message: format!("failed to archive thread: {err}"),
});
}
}
if let Some(ctx) = state_db_ctx
&& let Err(err) = ctx
.mark_archived(thread_id, archived_path.as_path(), Utc::now())
.await
{
if let Err(restore_err) = restore_rollout_moves(&rollout_moves) {
return Err(ThreadStoreError::Internal {
message: format!(
"failed to update archived thread metadata: {err}; failed to restore moved rollouts: {restore_err}"
),
});
}
return Err(ThreadStoreError::Internal {
message: format!("failed to update archived thread metadata: {err}"),
});
}
Ok(())第 N 次 rename 失败时,仅反向移动前 N−1 个已经移动的文件;数据库更新失败则尝试恢复全部文件。补偿本身也可能失败,错误会同时包含原始原因与 restore 原因。不能把这段代码描述成文件系统事务,更不能声称崩溃或任意取消都自动走进这些 Err 分支。
源码文件:codex-rs/thread-store/src/local/helpers.rs
相关函数/类型:restore_rollout_moves;行号:122-127。
pub(super) fn restore_rollout_moves(moves: &[(PathBuf, PathBuf)]) -> std::io::Result<()> {
for (source, destination) in moves.iter().rev() {
std::fs::rename(destination, source)?;
}
Ok(())
}逆序恢复避免把后执行的移动留在前面,但遇到第一个恢复失败就通过 ? 返回,没有一个后台无限重试器。已经创建的空目录、已发生的运行对象关闭,也不由这个 helper 撤销。
数据库更新只改相应元数据,而非重建一个新对话。
源码文件:codex-rs/state/src/runtime/threads.rs
相关函数/类型:mark_archived;行号:1061-1082。
pub async fn mark_archived(
&self,
thread_id: ThreadId,
rollout_path: &Path,
archived_at: DateTime<Utc>,
) -> anyhow::Result<()> {
let Some(mut metadata) = self.get_thread(thread_id).await? else {
return Ok(());
};
metadata.archived_at = Some(archived_at);
metadata.rollout_path = rollout_path.to_path_buf();
if let Some(updated_at) = file_modified_time_utc(rollout_path).await {
metadata.updated_at = updated_at;
}
if metadata.id != thread_id {
warn!(
"thread id mismatch during archive: expected {thread_id}, got {}",
metadata.id
);
}
self.upsert_thread(&metadata).await
}archived_at 使用操作时刻,updated_at 尝试从文件 mtime 读取;原 recency_at 没有在此重置。数据库中不存在该行时返回 Ok(()),所以本地文件可作为没有现成 metadata 行时的来源;不能倒过来说归档成功就一定插入了一行新记录。
4.3 后代部分成功
批次对根和后代的错误处理不同。
源码文件:codex-rs/thread-store/src/local/archive_thread.rs
相关函数/类型:archive_threads 的部分成功;行号:48-60。
let parent_thread_id = thread_ids[0];
let mut archived_thread_ids = Vec::new();
for thread_id in thread_ids {
let rollout_paths = owned_rollout_paths_from_index(&reference_index, thread_id);
match archive_thread_with_paths(store, thread_id, rollout_paths).await {
Ok(()) => archived_thread_ids.push(thread_id),
Err(err) if archived_thread_ids.is_empty() => return Err(err),
Err(err) => warn!(
"failed to archive spawned descendant thread {thread_id} while archiving {parent_thread_id}: {err}"
),
}
}
Ok(archived_thread_ids)根排第一;在还没有任何成功项时遇到错误立即返回。根成功之后,某个后代失败只 warning,循环继续,最终返回实际成功 ID。每个 Thread 内部的文件补偿,与多个 Thread 之间的部分成功,是两个层次。
测试 thread_archive_succeeds_when_descendant_archive_fails 建立 parent、child、grandchild,并在 child 的目标位置创建同名目录,使该次 rename 失败。
源码文件:codex-rs/app-server/tests/suite/v2/thread_archive.rs
相关函数/类型:thread_archive_succeeds_when_descendant_archive_fails;行号:650-677。
let mut archived_ids = Vec::new();
for _ in 0..2 {
let archived_notification: ThreadArchivedNotification = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_notification("thread/archived"),
)
.await??;
archived_ids.push(archived_notification.thread_id);
}
assert_eq!(archived_ids, vec![parent_id.clone(), grandchild_id.clone()]);
assert!(
timeout(
std::time::Duration::from_millis(250),
mcp.read_stream_until_notification_message("thread/archived"),
)
.await
.is_err()
);
assert!(
child_rollout_path.exists(),
"child should stay active after descendant archive failure"
);
assert!(
archived_child_path.is_dir(),
"test conflict should remain in archived sessions"
);预期通知只有 parent 和 grandchild,child 原文件仍在,目标冲突目录也还在。RPC 成功因此只说明根操作成功并完成了这次后代处理,不代表每个后代都归档成功。重复 archive 已归档的根会被拒绝,不能指望原样重发根请求就补做失败 child。
外层先响应,再逐一发送成功成员的归档通知。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_archive;行号:594-615。
pub(crate) async fn thread_archive(
&self,
request_id: ConnectionRequestId,
params: ThreadArchiveParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
match self.thread_archive_inner(params).await {
Ok((response, archived_thread_ids)) => {
self.outgoing
.send_response(request_id.clone(), response)
.await;
for thread_id in archived_thread_ids {
self.outgoing
.send_server_notification(ServerNotification::ThreadArchived(
ThreadArchivedNotification { thread_id },
))
.await;
}
Ok(None)
}
Err(error) => Err(error),
}
}Ok(None) 表示当前分支已经发送响应,不让通用分派再发一次。客户端要通过通知中的 ID 或重新读取来更新列表;出站顺序也不等于所有客户端已渲染完成。
5. 恢复一个归档
5.1 单对象入口
恢复归档没有调用子树遍历。
源码文件:codex-rs/thread-store/src/local/unarchive_thread.rs
相关函数/类型:unarchive_thread 的来源和锁;行号:22-40。
let thread_id = params.thread_id;
let _lifecycle_guard = store.live_writer_locks.lock_lifecycle(thread_id).await;
// Archive, delete, and revert use the same cross-process lock while moving or selecting
// rollout files. Unarchive must participate before it moves those files back.
let _writer_lock = store.writer_lock_coordinator.acquire(thread_id)?;
let state_db_ctx = store.state_db().await;
let selected_archived_path =
thread_rollout_resolver::resolve_current_including_archived(store, thread_id)
.await?
.filter(|resolved| resolved.location == RolloutLocation::Archived)
.map(|resolved| resolved.path)
.ok_or_else(|| ThreadStoreError::InvalidRequest {
message: format!("no archived rollout found for thread id {thread_id}"),
})?;
let mut rollout_paths = owned_rollout_paths(store, thread_id).await?;
if !rollout_paths.contains(&selected_archived_path) {
rollout_paths.push(selected_archived_path.clone());
}它等待 lifecycle 独占 guard,并参与同一跨进程 writer 锁。来源必须被 resolver 判断为 Archived;普通活动记录不能当作归档记录再次恢复。与归档批次相比,这条路径直接取文件锁,没有机械复用全部 writer mutex 批量流程。
5.2 日期与时间
目标目录来自 rollout 文件名中的年、月、日,而不是恢复归档当天的日期。
源码文件:codex-rs/thread-store/src/local/unarchive_thread.rs
相关函数/类型:unarchive_thread 的目标路径;行号:41-88。
let mut restored_path = None;
let mut rollout_moves = Vec::new();
for rollout_path in rollout_paths {
if !rollout_path_is_archived(store.config.codex_home.as_path(), rollout_path.as_path()) {
continue;
}
let canonical_archived_path = scoped_rollout_path(
store
.config
.codex_home
.join(codex_rollout::ARCHIVED_SESSIONS_SUBDIR),
rollout_path.as_path(),
"archived",
)?;
let file_name =
validated_rollout_file_name(canonical_archived_path.as_path(), rollout_path.as_path())?;
let Some((year, month, day)) = rollout_date_parts(&file_name) else {
return Err(ThreadStoreError::InvalidRequest {
message: format!(
"rollout path `{}` missing filename timestamp",
rollout_path.display()
),
});
};
let dest_dir = store
.config
.codex_home
.join(codex_rollout::SESSIONS_SUBDIR)
.join(year)
.join(month)
.join(day);
std::fs::create_dir_all(&dest_dir).map_err(|err| ThreadStoreError::Internal {
message: format!("failed to unarchive thread: {err}"),
})?;
let destination = dest_dir.join(&file_name);
if rollout_path == selected_archived_path {
restored_path = Some(destination.clone());
}
if !rollout_moves
.iter()
.any(|(source, _)| source == &canonical_archived_path)
{
rollout_moves.push((canonical_archived_path, destination));
}
}
let restored_path = restored_path.ok_or_else(|| ThreadStoreError::Internal {
message: format!("failed to unarchive selected rollout for thread {thread_id}"),
})?;恢复时保持文件名和逻辑身份,按文件时间放回目录;同一 Thread 的多个归档文件逐个规划移动。注意目录日期和列表更新时间不同,前者由原文件名决定,后者会在移动后更新。
源码文件:codex-rs/thread-store/src/local/unarchive_thread.rs
相关函数/类型:unarchive_thread 的时间与状态提交;行号:90-142。
for (index, (source, destination)) in rollout_moves.iter().enumerate() {
if let Err(err) = std::fs::rename(source, destination) {
if let Err(restore_err) = restore_rollout_moves(&rollout_moves[..index]) {
return Err(ThreadStoreError::Internal {
message: format!(
"failed to unarchive thread: {err}; failed to restore moved rollouts: {restore_err}"
),
});
}
return Err(ThreadStoreError::Internal {
message: format!("failed to unarchive thread: {err}"),
});
}
}
if let Err(err) = touch_modified_time(restored_path.as_path()) {
if let Err(restore_err) = restore_rollout_moves(&rollout_moves) {
return Err(ThreadStoreError::Internal {
message: format!(
"failed to update unarchived thread timestamp: {err}; failed to restore moved rollouts: {restore_err}"
),
});
}
return Err(ThreadStoreError::Internal {
message: format!("failed to update unarchived thread timestamp: {err}"),
});
}
if let Some(ctx) = state_db_ctx
&& let Err(err) = ctx
.mark_unarchived(thread_id, restored_path.as_path())
.await
{
if let Err(restore_err) = restore_rollout_moves(&rollout_moves) {
return Err(ThreadStoreError::Internal {
message: format!(
"failed to update unarchived thread metadata: {err}; failed to restore moved rollouts: {restore_err}"
),
});
}
return Err(ThreadStoreError::Internal {
message: format!("failed to update unarchived thread metadata: {err}"),
});
}
super::read_thread::read_thread(
store,
ReadThreadParams {
thread_id,
include_archived: false,
include_history: false,
},
)
.awaitrename 失败、mtime 修改失败、数据库更新失败,都有对应文件位置补偿。但最后 read_thread(...).await 失败直接返回,并没有再把文件归档回去。因此“unarchive RPC 返回错误”也可能发生在文件已经恢复、metadata 已经更新之后。
touch_modified_time 使用 FileTimes::new().set_modified(SystemTime::now()) 更新已选文件;它不是新增一轮对话。数据库层随后清除归档标记。
源码文件:codex-rs/state/src/runtime/threads.rs
相关函数/类型:mark_unarchived;行号:1085-1105。
pub async fn mark_unarchived(
&self,
thread_id: ThreadId,
rollout_path: &Path,
) -> anyhow::Result<()> {
let Some(mut metadata) = self.get_thread(thread_id).await? else {
return Ok(());
};
metadata.archived_at = None;
metadata.rollout_path = rollout_path.to_path_buf();
if let Some(updated_at) = file_modified_time_utc(rollout_path).await {
metadata.updated_at = updated_at;
}
if metadata.id != thread_id {
warn!(
"thread id mismatch during unarchive: expected {thread_id}, got {}",
metadata.id
);
}
self.upsert_thread(&metadata).await
}原 metadata 的其他字段沿用,尤其 section 与 recency 不因恢复归档就全部重置。mtime 已修改而后续数据库失败时,位置补偿也没有将原 mtime 一并恢复;补偿承诺应限定为源码实际处理的资源。
5.3 NotLoaded 的含义
处理器把返回的 Store 摘要投影为公开 Thread,并叠加当前 watch 状态。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_unarchive_response;行号:1999-2024。
async fn thread_unarchive_response(
&self,
params: ThreadUnarchiveParams,
) -> Result<(ThreadUnarchiveResponse, String), JSONRPCErrorError> {
let thread_id = ThreadId::from_string(¶ms.thread_id)
.map_err(|err| invalid_request(format!("invalid session id: {err}")))?;
let fallback_provider = self.config.model_provider_id.clone();
let stored_thread = self
.thread_store
.unarchive_thread(StoreArchiveThreadParams { thread_id })
.await
.map_err(|err| thread_store_mutation_error("unarchive", err))?;
let (mut thread, _) =
thread_from_stored_thread(stored_thread, fallback_provider.as_str(), &self.config.cwd);
thread.status = resolve_thread_status(
self.thread_watch_manager
.loaded_status_for_thread(&thread.id)
.await,
/*has_in_progress_turn*/ false,
);
self.attach_thread_name(thread_id, &mut thread).await;
let thread_id = thread.id.clone();
Ok((ThreadUnarchiveResponse { thread }, thread_id))
}这里没有创建 Core、恢复模型上下文或建立新的 listener。正常从已关闭归档恢复时,watch 中没有运行对象,因此响应通常是 NotLoaded;如果某后端或协作路径已有对应运行对象,则仍以真实 watch 结果为准,不能把该状态硬编码成恒定值。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_unarchive;行号:726-743。
pub(crate) async fn thread_unarchive(
&self,
request_id: ConnectionRequestId,
params: ThreadUnarchiveParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
match self.thread_unarchive_inner(params).await {
Ok((response, notification)) => {
self.outgoing
.send_response(request_id.clone(), response)
.await;
self.outgoing
.send_server_notification(ServerNotification::ThreadUnarchived(notification))
.await;
Ok(None)
}
Err(error) => Err(error),
}
}响应之后广播 thread/unarchived,它介绍的是存储变化,不是 thread/started。希望继续发消息的客户端还要按恢复 API 的运行对象契约处理。
上游测试先归档一个位于 pinned section 的 Thread,将归档文件 mtime 设置到早期时间,再恢复归档。
源码文件:codex-rs/app-server/tests/suite/v2/thread_unarchive.rs
相关函数/类型:thread_unarchive_moves_rollout_back_into_sessions_directory;行号:184-199,221-229。
let unarchived_notification: ThreadUnarchivedNotification = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_notification("thread/unarchived"),
)
.await??;
assert_eq!(unarchived_notification.thread_id, thread.id);
assert_eq!(unarchived_thread.section, Some(pinned_section.clone()));
assert_eq!(
unarchived_thread.section_entered_at,
Some(pinned_entered_at)
);
assert!(
unarchived_thread.updated_at > old_timestamp,
"expected updated_at to be bumped on unarchive"
);
assert_eq!(unarchived_thread.status, ThreadStatus::NotLoaded);
// ...
let rollout_path_display = rollout_path.display();
assert!(
rollout_path.exists(),
"expected rollout path {rollout_path_display} to be restored"
);
assert!(
!archived_path.exists(),
"expected archived rollout path {archived_path_display} to be moved"
);它同时要求 section 与进入 section 的时间不变、updated_at 变新、status 为 NotLoaded、文件回到原活动路径。这个组合正好说明“重新出现在活动列表”和“已经有模型执行者”是两种能力。另一个 pathless Store 测试证明协议摘要可没有本地路径,本文的日期目录算法仅属于本地后端。
6. 删除的引用集合
6.1 删除资格
删除允许清理活动或归档记录,还必须处理“主文件丢了,但数据库或后代还在”的情况。
源码文件:codex-rs/app-server/src/request_processors/thread_delete.rs
相关函数/类型:validate_root_thread_delete;行号:86-137。
async fn validate_root_thread_delete(
&self,
thread_id: ThreadId,
has_descendants: bool,
) -> Result<(), JSONRPCErrorError> {
if let Ok(thread) = self.thread_manager.get_thread(thread_id).await {
if !thread.config_snapshot().await.ephemeral {
return Ok(());
}
return Err(invalid_request(format!(
"thread is not persisted and cannot be deleted: {thread_id}"
)));
}
match self
.thread_store
.read_thread(StoreReadThreadParams {
thread_id,
include_archived: true,
include_history: false,
})
.await
{
Ok(_) => Ok(()),
Err(ThreadStoreError::ThreadNotFound { .. }) => {
if has_descendants {
return Ok(());
}
let Some(state_db) = self.state_db.as_ref() else {
return Err(thread_store_delete_error(
ThreadStoreError::ThreadNotFound { thread_id },
));
};
if state_db
.get_thread(thread_id)
.await
.map_err(|err| {
internal_error(format!(
"failed to read app-server state for {thread_id}: {err}"
))
})?
.is_some()
{
Ok(())
} else {
Err(thread_store_delete_error(
ThreadStoreError::ThreadNotFound { thread_id },
))
}
}
Err(err) => Err(thread_store_delete_error(err)),
}
}已加载且非 ephemeral 的 Thread 可继续删除,即使 rollout 尚未物化;临时对象明确拒绝。未加载时先读 Store,找不到文件但有后代,或仍有主状态行,也可继续。这是为清理残留和重试保留入口,并不是把完全未知的 ID 都视为删除成功。
6.2 两种顺序
App Server 先发现子树和清理运行资源,再把删除列表交给 Store。
源码文件:codex-rs/app-server/src/request_processors/thread_delete.rs
相关函数/类型:thread_delete_response;行号:31-74。
async fn thread_delete_response(
&self,
params: ThreadDeleteParams,
deleted_thread_ids: &mut Vec<String>,
) -> Result<ThreadDeleteResponse, JSONRPCErrorError> {
let thread_id = ThreadId::from_string(¶ms.thread_id)
.map_err(|err| invalid_request(format!("invalid thread id: {err}")))?;
let thread_ids = self.state_db_spawn_subtree_thread_ids(thread_id).await?;
self.validate_root_thread_delete(thread_id, thread_ids.len() > 1)
.await?;
for thread_id_to_delete in thread_ids.iter().copied() {
self.prepare_thread_for_delete(thread_id_to_delete).await;
}
let mut delete_order: Vec<_> = thread_ids.iter().skip(1).rev().copied().collect();
delete_order.push(thread_id);
self.thread_store
.delete_threads(StoreDeleteThreadsParams {
thread_ids: delete_order.clone(),
})
.await
.map_err(thread_store_delete_error)?;
if let Some(state_db) = self.state_db.as_ref() {
state_db
.delete_threads_strict(thread_ids.as_slice())
.await
.map_err(|err| {
internal_error(format!(
"failed to delete app-server state for {thread_id}: {err}"
))
})?;
}
deleted_thread_ids.extend(
delete_order
.into_iter()
.map(|thread_id| thread_id.to_string()),
);
Ok(ThreadDeleteResponse {})
}删除顺序是逆序后代、最后根;主 StateRuntime 的严格清理放在 Store 成功之后。与归档相比,删除不能在遇到后代错误后仍将整个结果报告成功。deleted_thread_ids 只有全部步骤成功后才填入。
每个成员关闭后还会先刷新日志缓存。
源码文件:codex-rs/app-server/src/request_processors/thread_delete.rs
相关函数/类型:prepare_thread_for_delete;行号:139-144。
async fn prepare_thread_for_delete(&self, thread_id: ThreadId) {
self.prepare_thread_for_removal(thread_id, "delete").await;
if let Some(log_db) = self.log_db.as_ref() {
log_db.flush().await;
}
}刷新发生在持久清理前,减少已排队旧日志在删除后才写回的窗口;这里没有声明一个阻止全系统所有日志生产者的总屏障。随后锁顺序和删除顺序分开计算。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:delete_threads;行号:78-115。
pub(super) async fn delete_threads(
store: &LocalThreadStore,
params: DeleteThreadsParams,
) -> ThreadStoreResult<()> {
let thread_ids = params.thread_ids;
if thread_ids.is_empty() {
return Ok(());
}
let mut lock_thread_ids = thread_ids.clone();
lock_thread_ids.sort_unstable_by_key(ToString::to_string);
lock_thread_ids.dedup();
let mut _lifecycle_guards = Vec::with_capacity(lock_thread_ids.len());
for thread_id in &lock_thread_ids {
_lifecycle_guards.push(store.live_writer_locks.lock_lifecycle(*thread_id).await);
}
let mut _live_writer_guards = Vec::with_capacity(lock_thread_ids.len());
for &thread_id in &lock_thread_ids {
_live_writer_guards.push(store.live_writer_locks.lock(thread_id).await);
}
let reference_index = scan_reference_index(store).await?;
let thread_rollouts = thread_ids
.iter()
.map(|thread_id| ThreadRollouts::from_index(&reference_index, *thread_id))
.collect::<Vec<_>>();
ensure_no_external_references(&reference_index, thread_rollouts.as_slice())?;
let mut writer_guards = store.acquire_writer_locks(&lock_thread_ids).await?;
for thread_rollouts in thread_rollouts {
match delete_thread_after_reference_check(store, thread_rollouts, &mut writer_guards).await
{
Ok(()) | Err(ThreadStoreError::ThreadNotFound { .. }) => {}
Err(err) => return Err(err),
}
}
Ok(())
}lock_thread_ids 排序去重用于一致取锁,原 thread_ids 顺序用于逐个删除。最关键的调用在循环前:ensure_no_external_references 检查整个集合,然后才开始删除第一个成员。集合内部相互引用不应导致一个本可一起删除的子树永远无法清理。
删除与归档还在 live recorder 处理上有所区别。归档提前拒绝仍存在的 recorder;删除允许在自己的 Store 内接管其既有锁。
源码文件:codex-rs/thread-store/src/local/mod.rs
相关函数/类型:acquire_writer_locks;行号:304-316。
async fn acquire_writer_locks(
&self,
thread_ids: &[ThreadId],
) -> ThreadStoreResult<Vec<WriterLockGuard>> {
let mut writer_locks = Vec::with_capacity(thread_ids.len());
for &thread_id in thread_ids {
if self.live_recorders.lock().await.contains_key(&thread_id) {
continue;
}
writer_locks.push(self.writer_lock_coordinator.acquire(thread_id)?);
}
Ok(writer_locks)
}helper 跳过当前 Store 已拥有的 live recorder,避免重复获取同一文件锁。后续清理会从 recorder 转移出该 guard,一直保留到文件清理结束;另一个进程的 writer 仍会触发文件锁冲突。这说明“shutdown 超时后继续”没有跳过所有权保护,但两种操作的具体拒绝位置不同。
6.3 物理引用计数
必须把逻辑 Thread 集合转换为物理 rollout 集合,才能检查真实依赖。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:ThreadRollouts;行号:28-62。
struct ThreadRollouts {
thread_id: codex_protocol::ThreadId,
rollout_ids: HashSet<codex_protocol::ThreadId>,
paths: Vec<PathBuf>,
}
impl ThreadRollouts {
fn from_index(
reference_index: &RolloutReferenceIndex,
thread_id: codex_protocol::ThreadId,
) -> Self {
let mut rollout_ids = HashSet::new();
let paths = reference_index
.rollouts_for_thread(thread_id)
.map(|(rollout_id, path)| {
rollout_ids.insert(rollout_id);
path.to_path_buf()
})
.collect();
Self {
thread_id,
rollout_ids,
paths,
}
}
fn add_path(&mut self, path: PathBuf) {
if let Some(rollout_id) = codex_rollout::rollout_id_from_path(path.as_path()) {
self.rollout_ids.insert(rollout_id);
}
if !self.paths.contains(&path) {
self.paths.push(path);
}
}
}ThreadRollouts 记录逻辑 ID、所属物理 ID 集合与路径列表。ID 集合用于引用算法,路径列表用于文件删除;不是只把请求 ID 当作一个文件名。后续 active/archived 路径发现还会通过 add_path 补足候选。
索引自身保存正向元数据与反向计数。
源码文件:codex-rs/rollout/src/rollout_reference_index.rs
相关函数/类型:RolloutReferenceIndex / IndexedRollout;行号:24-34。
pub struct RolloutReferenceIndex {
rollouts_by_id: HashMap<RolloutId, IndexedRollout>,
reference_counts_by_rollout: HashMap<RolloutId, usize>,
}
#[derive(Debug)]
struct IndexedRollout {
thread_id: ThreadId,
path: PathBuf,
history_base: Option<HistoryPosition>,
}rollouts_by_id 以物理 ID 寻找所属逻辑 Thread、路径和 history base;reference_counts_by_rollout 则以目标物理 ID 保存被引用次数。IndexedRollout.thread_id 是逻辑所有者,history_base.thread_id 是指向的物理 rollout,这两个同名字段处于不同结构中。
这张图解释集合转换,而非持久数据库表:ThreadRollouts 是当前删除调用的局部计算结果。失去物理 ID 集合只保留文件路径,会让引用预检和后面的投影清理无法使用同一个身份。
下面把两种关系放在一张图里。spawn 边决定删除成员,history base 边决定能否删。
P、C 同时删除时,C 对 P 的内部引用可以扣除;E 不在 spawn 删除集合中,它对 P 的引用必须保留。把 E 也当成 spawn 后代自动删除,会改变用户分叉的所有权含义。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:ensure_no_external_references;行号:117-148。
fn ensure_no_external_references(
reference_index: &RolloutReferenceIndex,
thread_rollouts: &[ThreadRollouts],
) -> ThreadStoreResult<()> {
let deletion_rollout_ids = thread_rollouts
.iter()
.flat_map(|thread_rollouts| thread_rollouts.rollout_ids.iter().copied())
.collect::<HashSet<_>>();
let mut internal_reference_counts = HashMap::new();
for source_rollout_id in &deletion_rollout_ids {
if let Some(history_base) = reference_index.history_base(*source_rollout_id)
&& history_base.thread_id != *source_rollout_id
&& deletion_rollout_ids.contains(&history_base.thread_id)
{
*internal_reference_counts
.entry(history_base.thread_id)
.or_default() += 1;
}
}
for thread_rollouts in thread_rollouts {
if thread_rollouts.rollout_ids.iter().any(|rollout_id| {
let internal_reference_count = internal_reference_counts
.get(rollout_id)
.copied()
.unwrap_or_default();
reference_index.reference_count(*rollout_id) > internal_reference_count
}) {
return Err(referenced_thread_error(thread_rollouts.thread_id));
}
}
Ok(())
}令 D 是待删物理 ID 集合。第一轮对每个 D 内源记录,数它指向 D 内目标的非自引用边;第二轮检查某目标的总入边数是否大于内部入边数。只要差值为正,就存在集合外依赖并拒绝。使用哈希集合和计数表,集合处理按物理记录数线性展开;总成本还包含扫描磁盘元数据,不能把整个 RPC 说成常数时间。
这个算法也能保护历史替换前的旧 rollout:外部分支可能引用的不是 Thread 当前选中的物理文件。只检查当前路径会遗漏这种依赖。
6.4 索引的边界
引用安全依赖于扫描实际可读的元数据。
源码文件:codex-rs/rollout/src/rollout_reference_index.rs
相关函数/类型:RolloutReferenceIndex::scan_with_deadline;行号:126-141,145-156。
let Some(rollout_file) = RolloutFile::from_path(path) else {
continue;
};
let Some(rollout_id) = crate::rollout_id_from_path(rollout_file.path()) else {
continue;
};
let Ok(meta) = crate::read_session_meta_line(rollout_file.path()).await else {
continue;
};
if let Entry::Vacant(entry) = rollouts_by_id.entry(rollout_id) {
entry.insert(IndexedRollout {
thread_id: meta.meta.id,
path: rollout_file.into_path(),
history_base: meta.meta.history_base,
});
}
// ...
let mut reference_counts_by_rollout = HashMap::new();
for (rollout_id, rollout) in &rollouts_by_id {
let Some(history_base) = rollout.history_base else {
continue;
};
if history_base.thread_id == *rollout_id {
continue;
}
*reference_counts_by_rollout
.entry(history_base.thread_id)
.or_default() += 1;
}扫描 active 和 archived 目录,识别 rollout、读取首条 SessionMeta;无法解析的条目直接跳过。计数按物理 rollout ID 去重,自引用不计入入边。这不是对损坏磁盘的完整恢复工具:某个文件的元数据无法读取时,它声明的 history base 也不在索引中。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:delete_thread_ignores_unreadable_reference_metadata;行号:443-465。
async fn delete_thread_ignores_unreadable_reference_metadata() {
let home = TempDir::new().expect("temp dir");
let store = LocalThreadStore::new(test_config(home.path()), /*state_db*/ None);
let source_uuid = Uuid::from_u128(305);
let source_thread_id =
ThreadId::from_string(&source_uuid.to_string()).expect("valid source thread id");
let source_path = write_session_file(home.path(), "2025-01-03T12-00-00", source_uuid)
.expect("source session file");
let unreadable_path = source_path.with_file_name(format!(
"rollout-2025-01-03T12-00-01-{}.jsonl",
Uuid::from_u128(306)
));
std::fs::write(unreadable_path, "{not json}\n").expect("unreadable rollout metadata");
store
.delete_thread(DeleteThreadParams {
thread_id: source_thread_id,
})
.await
.expect("unreadable metadata should not block delete");
assert!(!source_path.exists());
}测试故意写入一个看似 rollout 的文件名但正文是 not-json,要求删除另一个正常 Thread 仍成功。这说明坏元数据不会自动导致全局拒绝,不能宣传为“无论存储是否损坏都能发现全部依赖”。遇到损坏历史应先分析原文件,而不是将索引中没有边当成从未存在依赖的证明。
6.5 集合反例
集成测试建立 P 的 spawn 子 C,同时让集合外 E 的 history base 指向 P,然后请求删除 P。
源码文件:codex-rs/app-server/tests/suite/v2/thread_delete.rs
相关函数/类型:thread_delete_preflights_external_fork_references_for_spawned_subtrees;行号:238-261。
assert_eq!(
delete_err.error.message,
format!("cannot delete thread {parent_thread_id}: forked history still references it")
);
for thread_id in [parent_thread_id, child_thread_id, external_thread_id] {
assert!(
find_thread_path_by_id_str(
codex_home.path(),
&thread_id.to_string(),
/*state_db_ctx*/ None,
)
.await?
.is_some(),
"expected rollout for {thread_id} to remain"
);
}
assert_eq!(
state_db
.list_thread_spawn_descendants(parent_thread_id)
.await?,
vec![child_thread_id]
);
Ok(())预期返回明确的 referenced-history 错误,P、C、E 的文件都在,spawn 子树也仍可重新发现。这检验的是磁盘删除前的预检;App Server 的运行资源清理已经在 Store 调用前发生,不能外推成失败后所有运行订阅也完全不变。
反方向,Store 单元测试把 child 与 source 同时放进删除集合。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:delete_threads_allows_internal_history_references;行号:468-511。
async fn delete_threads_allows_internal_history_references() {
let home = TempDir::new().expect("temp dir");
let store = LocalThreadStore::new(test_config(home.path()), /*state_db*/ None);
let source_uuid = Uuid::from_u128(307);
let source_thread_id =
ThreadId::from_string(&source_uuid.to_string()).expect("valid source thread id");
let source_path = write_session_file_with_history_mode(
home.path(),
"2025-01-03T12-00-00",
source_uuid,
ThreadHistoryMode::Paginated,
)
.expect("source session file");
let child_uuid = Uuid::from_u128(308);
let child_thread_id =
ThreadId::from_string(&child_uuid.to_string()).expect("valid child thread id");
let child_path = write_session_file_with_history_mode(
home.path(),
"2025-01-03T12-00-01",
child_uuid,
ThreadHistoryMode::Paginated,
)
.expect("child session file");
set_history_base(
child_path.as_path(),
HistoryPosition {
thread_id: source_thread_id,
end_ordinal_exclusive: 1,
end_byte_offset: std::fs::metadata(source_path.as_path())
.expect("source rollout metadata")
.len(),
},
);
store
.delete_threads(DeleteThreadsParams {
thread_ids: vec![child_thread_id, source_thread_id],
})
.await
.expect("internal references should not block batch delete");
assert!(!source_path.exists());
assert!(!child_path.exists());
}它要求两份文件一起删除成功。单独证明“有引用就禁止删除”是不够的,还必须证明“只剩内部引用可以删除”,否则实现会过度拒绝合法的集合清理。
7. 不可逆的清理
7.1 文件与投影
通过引用预检后,单个 Thread 的清理先补足 active/archived 路径,再处理投影、writer、文件和名称索引。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:delete_thread_after_reference_check 的清理顺序;行号:204-227。
thread_rollouts.rollout_ids.insert(thread_id);
for rollout_id in thread_rollouts.rollout_ids {
super::thread_history::delete_thread(store, rollout_id).await?;
}
// Drop the recorder before removing files, but retain its writer lock until cleanup finishes.
if let Some(entry) = store.live_recorders.lock().await.remove(&thread_id) {
writer_guards.push(entry.writer_lock);
}
let found_rollout_path = !thread_rollouts.paths.is_empty();
for rollout_path in thread_rollouts.paths {
delete_rollout_file(store, rollout_path.as_path())?;
}
remove_thread_name_entries(store.config.codex_home.as_path(), thread_id)
.await
.map_err(|err| ThreadStoreError::Internal {
message: format!("failed to delete thread name index entries for {thread_id}: {err}"),
})?;
if !found_rollout_path {
return Err(ThreadStoreError::ThreadNotFound { thread_id });
}
Ok(())历史投影按物理 rollout ID 清理。live recorder 从表中移除时,writer lock 被转移到局部 guard 集合而不是提前释放,防止删文件过程中另一写入者接手。这个顺序也意味着文件删除失败时投影可能已删除,不能承诺 Store 失败会回到调用前的状态。
投影清理本身有独立事务,也有独立的不可用条件。
源码文件:codex-rs/thread-store/src/local/thread_history.rs
相关函数/类型:delete_thread;行号:244-286。
pub(super) async fn delete_thread(
store: &LocalThreadStore,
thread_id: ThreadId,
) -> ThreadStoreResult<()> {
let db_path = store.config.sqlite.thread_history_db_path();
if !tokio::fs::try_exists(db_path.as_path())
.await
.map_err(thread_history_delete_error)?
{
return Ok(());
}
let pool = store.thread_history_db().await?;
let mut transaction = pool
.begin_with("BEGIN IMMEDIATE")
.await
.map_err(thread_history_delete_error)?;
let thread_id = thread_id.to_string();
sqlx::query("DELETE FROM thread_items WHERE thread_id = ?")
.bind(thread_id.as_str())
.execute(&mut *transaction)
.await
.map_err(thread_history_delete_error)?;
sqlx::query("DELETE FROM thread_realtime_items WHERE thread_id = ?")
.bind(thread_id.as_str())
.execute(&mut *transaction)
.await
.map_err(thread_history_delete_error)?;
sqlx::query("DELETE FROM thread_turns WHERE thread_id = ?")
.bind(thread_id.as_str())
.execute(&mut *transaction)
.await
.map_err(thread_history_delete_error)?;
sqlx::query("DELETE FROM thread_history_projection_state WHERE thread_id = ?")
.bind(thread_id.as_str())
.execute(&mut *transaction)
.await
.map_err(thread_history_delete_error)?;
transaction
.commit()
.await
.map_err(thread_history_delete_error)
}历史库文件不存在时直接成功;存在时用 BEGIN IMMEDIATE 删除 items、realtime items、turns 和 projection state。它在这个历史库内部提交,不把后续文件删除纳入事务。
源码文件:codex-rs/thread-store/src/local/mod.rs
相关函数/类型:thread_history_db;行号:255-269。
async fn thread_history_db(&self) -> ThreadStoreResult<&sqlx::SqlitePool> {
if self.state_db.is_none() {
return Err(ThreadStoreError::Unsupported {
operation: "paginated_history",
});
}
self.thread_history_db
.get_or_try_init(|| async {
codex_state::open_thread_history_db(&self.config.sqlite).await
})
.await
.map_err(|err| ThreadStoreError::Internal {
message: format!("failed to open thread history database: {err}"),
})
}如果历史库已经存在而 Store 没有 state DB 能力,取得 pool 会返回 Unsupported(paginated_history),操作在删文件前停止。delete_thread_without_state_db_preserves_materialized_thread_history 正是构造已有四类投影各一行,要求失败后四个计数仍为 1、rollout 仍存在。不能将“缺 state DB”概括成忽略所有投影继续删除。
无已发现文件时仍尝试清理名称索引,再返回 ThreadNotFound;批量入口将该错误视为已清理,继续后续成员。这与单独删除未知根之前的资格检查组合起来,既支持残留清理,又不会无条件接受所有错误 ID。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:delete_rollout_file;行号:230-236。
fn delete_rollout_file(store: &LocalThreadStore, rollout_path: &Path) -> ThreadStoreResult<bool> {
let plain_path = codex_rollout::plain_rollout_path(rollout_path);
let compressed_path = plain_path.with_extension("jsonl.zst");
let deleted_plain = delete_rollout_path(store, plain_path.as_path())?;
let deleted_compressed = delete_rollout_path(store, compressed_path.as_path())?;
Ok(deleted_plain || deleted_compressed)
}普通 JSONL 和 .jsonl.zst 压缩兄弟都尝试删除。第一份成功而第二份失败时没有重建第一份文件的分支;这就是硬删除与归档补偿的本质差别。
源码文件:codex-rs/thread-store/src/local/delete_thread.rs
相关函数/类型:delete_rollout_path;行号:238-266。
fn delete_rollout_path(store: &LocalThreadStore, rollout_path: &Path) -> ThreadStoreResult<bool> {
let canonical_rollout_path = scoped_rollout_path(
store.config.codex_home.join(SESSIONS_SUBDIR),
rollout_path,
"sessions",
)
.or_else(|_| {
scoped_rollout_path(
store.config.codex_home.join(ARCHIVED_SESSIONS_SUBDIR),
rollout_path,
"archived sessions",
)
})
.or_else(|err| match rollout_path.try_exists() {
Ok(false) => Ok(rollout_path.to_path_buf()),
Ok(true) | Err(_) => Err(err),
})?;
validated_rollout_file_name(&canonical_rollout_path, rollout_path)?;
match std::fs::remove_file(&canonical_rollout_path) {
Ok(()) => Ok(true),
Err(err) if err.kind() == ErrorKind::NotFound => Ok(false),
Err(err) => Err(ThreadStoreError::Internal {
message: format!(
"failed to delete rollout file `{}`: {err}",
canonical_rollout_path.display()
),
}),
}
}先按活动目录校验,失败后尝试归档目录。若发现文件后它已经消失,try_exists=false 允许继续,最终 remove_file 的 NotFound 当作已删除;真正的访问、权限等错误则传播。它没有执行磁盘安全擦除,也不会删除用户工作目录里的代码文件。
7.2 最后删除主行
所有相关 rollout 清理成功之后,App Server 才调用 StateRuntime。
源码文件:codex-rs/state/src/runtime/threads.rs
相关函数/类型:delete_threads_strict;行号:1116-1162。
pub async fn delete_threads_strict(&self, thread_ids: &[ThreadId]) -> anyhow::Result<u64> {
if thread_ids.is_empty() {
return Ok(0);
}
let thread_id_strings = thread_ids
.iter()
.map(ThreadId::to_string)
.collect::<Vec<_>>();
for (thread_id, thread_id_string) in thread_ids.iter().zip(&thread_id_strings) {
sqlx::query("DELETE FROM logs WHERE thread_id = ?")
.bind(thread_id_string)
.execute(self.logs_pool.as_ref())
.await?;
self.thread_queue.delete_thread_queue(*thread_id).await?;
self.memories.delete_thread_memory(*thread_id).await?;
self.thread_goals.delete_thread_goal(*thread_id).await?;
}
let mut tx = self.pool.begin().await?;
for thread_id_string in &thread_id_strings {
sqlx::query("DELETE FROM thread_dynamic_tools WHERE thread_id = ?")
.bind(thread_id_string)
.execute(&mut *tx)
.await?;
}
for thread_id_string in &thread_id_strings {
sqlx::query(
"DELETE FROM thread_spawn_edges WHERE parent_thread_id = ? OR child_thread_id = ?",
)
.bind(thread_id_string)
.bind(thread_id_string)
.execute(&mut *tx)
.await?;
}
let mut rows_affected = 0;
for thread_id_string in &thread_id_strings {
rows_affected += sqlx::query("DELETE FROM threads WHERE id = ?")
.bind(thread_id_string)
.execute(&mut *tx)
.await?
.rows_affected();
}
tx.commit().await?;
Ok(rows_affected)
}日志、待执行队列、记忆、goal 依次清理,然后才开启主数据库事务:先 dynamic tools,后 spawn edges,最后 threads。保留主行和图到最后,是为了前面失败时仍有足够信息重新发现需要清理的子树。
这里涉及不同存储,并没有一个包住日志库、记忆库、文件系统和主库的共同事务。前面的日志已经删除,而后面的记忆删除失败,是源码允许的中间状态。把“strict”解释成所有资源一起提交或一起回滚,会直接误导故障恢复。
以下时序把不可逆清理与最后的事务分开。
图中的事务只覆盖主库最后一段。任何此前成功的文件删除都不会因主库失败自动撤销;重试依赖剩余主行和 spawn 图,以及批量删除对缺失文件的容忍。
上游测试主动关闭日志数据库来制造前置清理失败。
源码文件:codex-rs/state/src/runtime/threads.rs
相关函数/类型:delete_thread_keeps_retry_graph_on_cleanup_failure;行号:1943-1973。
async fn delete_thread_keeps_retry_graph_on_cleanup_failure() -> Result<()> {
let codex_home = unique_temp_dir();
let runtime = StateRuntime::init(
crate::SqliteConfig::new_for_testing(codex_home.as_path().abs()),
"test-provider".to_string(),
)
.await?;
let thread_id = ThreadId::from_string("00000000-0000-0000-0000-000000000405")?;
let child_thread_id = ThreadId::from_string("00000000-0000-0000-0000-000000000406")?;
runtime
.upsert_thread(&test_thread_metadata(
&codex_home,
thread_id,
codex_home.clone(),
))
.await?;
seed_thread_cleanup_state(&runtime, thread_id, child_thread_id).await?;
runtime.logs_pool.close().await;
runtime
.delete_thread(thread_id)
.await
.expect_err("closed log db should fail deletion");
assert!(runtime.get_thread(thread_id).await?.is_some());
assert_eq!(
runtime.list_thread_spawn_descendants(thread_id).await?,
vec![child_thread_id]
);
Ok(())
}断言主 Thread 行和 child 图仍在,证明后续可以重新发现目标。它没有模拟每一种磁盘失败,也不证明所有先前清理步骤都会撤销。定位时要记录最后完成的阶段,而不是笼统说“数据库删除失败,所以文件肯定还在”。
7.3 成功通知边界
删除外层只在整个业务结果成功时发送响应与通知。
源码文件:codex-rs/app-server/src/request_processors/thread_delete.rs
相关函数/类型:thread_delete;行号:7-29。
pub(crate) async fn thread_delete(
&self,
request_id: ConnectionRequestId,
params: ThreadDeleteParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
let mut deleted_thread_ids = Vec::new();
let result = {
let _thread_list_state_permit = self.acquire_thread_list_state_permit().await?;
self.thread_delete_response(params, &mut deleted_thread_ids)
.await
};
match result {
Ok(response) => {
self.outgoing
.send_response(request_id.clone(), response)
.await;
self.send_thread_deleted_notifications(deleted_thread_ids)
.await;
Ok(None)
}
Err(error) => Err(error),
}
}如果 Store 已删掉部分文件,后面某个清理步骤失败,客户端收到错误且没有这条分支的成功 thread/deleted 通知。因此没有通知不能作为文件仍在的证明;客户端需要重新查询服务器状态。
源码文件:codex-rs/app-server/src/request_processors/thread_delete.rs
相关函数/类型:thread_store_delete_error;行号:147-160。
fn thread_store_delete_error(err: ThreadStoreError) -> JSONRPCErrorError {
match err {
ThreadStoreError::ThreadNotFound { thread_id } => {
invalid_request(format!("thread not found: {thread_id}"))
}
ThreadStoreError::InvalidRequest { message } | ThreadStoreError::Conflict { message } => {
invalid_request(message)
}
ThreadStoreError::Unsupported { operation } => {
unsupported_thread_store_operation(operation)
}
err => internal_error(format!("failed to delete thread: {err}")),
}
}不存在、非法引用和 writer 冲突都落到请求错误一侧;后端不支持有独立映射;其余清理错误成为内部错误。排障时应保留错误消息中的阶段,不只按数字错误码统一重试。
8. 从现象反查
8.1 资源残留
下表依据前面的调用顺序选择检查点,不需要修改或删除用户文件。
| 现象 | 先检查 | 对应源码含义 |
|---|---|---|
| unarchive 成功但不能直接开始 Turn | loaded/watch、是否已经 resume | unarchive 只返回存储摘要,未创建 Core |
| archive 成功但 child 仍在活动目录 | child 移动失败日志、实际 archived 通知 ID | 根成功后后代允许部分失败 |
| archive 长时间没有响应 | prepared fork reservation、writer | 10 秒只限制 Core shutdown,未限制 lifecycle 等待 |
| delete 报 fork history referenced | 集合外 history base、被引用物理 ID | 用户 fork 不等于 spawn 后代 |
| delete 报内部错误但文件已少了一部分 | 投影、文件、名称索引、各 DB 的最后成功点 | 硬删除无跨资源回滚 |
| shutdown 超时后列表仍发生变化 | manager 与 Store writer 各自状态 | 注册移除、关闭确认与持久变更是不同阶段 |
本地存储算法没有按 Windows、macOS、Linux 分出不同的集合语义,但路径规范化、文件锁和 rename/remove 的结果由实际 OS 与文件系统决定。不能从一台机器的测试外推出所有平台上的占用文件行为。后端替换、reference index 扫描策略或 strict cleanup 顺序变化时,需重新审视这里的失败边界。
8.2 源码复现
在 Codex 仓库根目录执行以下测试。它们使用临时目录与模拟服务器,不以真实对话目录作为删除目标。
just test --locked -p codex-app-server --test all \
-E 'test(suite::v2::thread_archive::) or test(suite::v2::thread_unarchive::) or test(suite::v2::thread_delete::)'
just test --locked -p codex-thread-store --lib \
-E 'test(local::archive_thread::tests) or test(local::unarchive_thread::tests) or test(local::delete_thread::tests)'
just test --locked -p codex-state --lib \
-E 'test(delete_thread_cleans_associated_state) or test(delete_thread_keeps_retry_graph_on_cleanup_failure)'
rg -n 'prepare_thread_for_removal|thread_archive_response|thread_unarchive_response' \
codex-rs/app-server/src/request_processors/thread_processor.rs
rg -n 'ensure_no_external_references|delete_threads_strict' \
codex-rs/thread-store/src/local/delete_thread.rs \
codex-rs/state/src/runtime/threads.rs尝试追踪一个具体案例:P spawn 出 C,E 从 P 分叉,C 的归档目标又存在冲突目录。先说明归档 P 时哪些 ID 会成功、哪些运行对象会被移除;再说明为什么 unarchive P 不会把整棵树重新运行起来;最后判断删除 P 会在哪个检查点遇到 E 的引用。把答案分别落实到发现集合、关闭顺序、Store 预检与消费者通知,便能区分“文件没有移动”“运行对象没有恢复”和“持久依赖不允许删除”。
