ThreadManager恢复
恢复 Thread 不是“读取 JSONL 后创建一条新会话”。Codex 要保留原 ThreadId 与 agent-tree 身份,重新 打开持久化 writer,把 rollout 归约成模型 history、world state、压缩窗口和上一 Turn 设置,再建立新的 Session、submission loop 与产品事件流。
这条流程还要处理一个看似矛盾的情况:用户要求恢复的 Thread 可能仍在当前进程运行。此时正确结果不是 再启动一套 Session,而是返回已有 Arc<CodexThread>;只有旧 handle 已停止,manager 才会移除它并用 同一个 ThreadId 创建新的运行时。
本文面向已理解 Thread/Session/Turn 概念,并读过 ThreadManager创建 的读者。本文只讲 Resume:如何选择历史、 重开 writer、归约 rollout并恢复同一身份;不会展开 Fork的历史截断,也不重复 Session每项服务的构造细节。
读完后,应能从 rollout path或预加载 history走到 InitialHistory::Resumed、active-map fast path、 LiveThread::resume、RolloutReconstruction 和重新注册,并能区分“模型history恢复成功”与“WorldState baseline可信”是两个不同结论。
1. 恢复保留身份
先把恢复与创建、分叉区分开:
| 操作 | ThreadId | 初始历史 | 持久化动作 | manager 结果 |
|---|---|---|---|---|
| 新建 | 新 UUIDv7 | InitialHistory::New | 创建新 Thread | 新 Arc<CodexThread> |
| 恢复 | 保留 rollout 中的 ID | InitialHistory::Resumed | 重开原 Thread writer | 复用活跃 handle,或重建已停止 handle |
| 分叉 | 新 UUIDv7 | InitialHistory::Forked | 创建 child Thread | 新 handle,并记录 lineage |
恢复保留的是业务身份和可重放状态,不是旧进程中的 Tokio 对象。channel、lock、MCP connection、模型 client、取消令牌和 pending oneshot 都不能跨进程恢复;它们由新 Session 重新构造。当前 paginated history 不能通过 rollout-path eager API 强行恢复:当 include_history: true 与 paginated 模式冲突时, store 返回 Unsupported,调用方必须先取得正确历史范围再使用 resume_thread_with_history()。
图中的“重建状态”发生在 Session 初始化期间,“注册”发生在首个 SessionConfigured 事件通过检查之后。 因此 rollout 已成功解析,不代表恢复的 Thread 已经对 manager 调用方可见。
2. spawn入口
ThreadManager 提供两种恢复入口:一种接收 rollout path,由 store 负责读取;另一种接收已经加载的 InitialHistory。二者最终都构造 StartThreadOptions,进入同一个 spawn_thread()。
源码位置:codex-rs/core/src/thread_manager.rs :: resume_thread_from_rollout, resume_thread_with_history
pub async fn resume_thread_from_rollout(
&self,
config: Config,
rollout_path: PathBuf,
auth_manager: Arc<AuthManager>,
parent_trace: Option<W3cTraceContext>,
client_mcp_extensions: ClientMcpExtensions,
) -> CodexResult<NewThread> {
// 路径入口先通过 ThreadStore 取得 StoredThread,不直接在 manager 中逐行解析 JSONL。
let initial_history = self.initial_history_from_rollout_path(rollout_path).await?;
Box::pin(self.resume_thread_with_history(
config,
initial_history,
auth_manager,
parent_trace,
client_mcp_extensions,
))
.await
}
pub async fn resume_thread_with_history(
&self,
config: Config,
initial_history: InitialHistory,
auth_manager: Arc<AuthManager>,
parent_trace: Option<W3cTraceContext>,
client_mcp_extensions: ClientMcpExtensions,
) -> CodexResult<NewThread> {
let agent_control = self.agent_control_for_config(&config);
// source 以历史 SessionMeta 为准;缺失时才回退到当前 manager 的产品来源。
let (session_source, thread_source) = initial_history
.get_resumed_session_sources()
.unwrap_or_else(|| (self.state.session_source.clone(), None));
if let InitialHistory::Resumed(resumed) = &initial_history
&& initial_history.get_multi_agent_version() == Some(MultiAgentVersion::V2)
&& !session_source.is_non_root_agent()
{
// V2 root 恢复时先补回持久化 agent metadata,再启动 Session。
agent_control
.restore_v2_agent_metadata(&config, resumed.conversation_id)
.await;
}
let options = StartThreadOptions {
initial_history,
session_source: Some(session_source),
thread_source,
parent_trace,
client_mcp_extensions,
..StartThreadOptions::new(config)
};
Box::pin(self.state.spawn_thread(ThreadSpawnRequest::new(
options,
auth_manager,
agent_control,
)))
.await
}auth_manager 是恢复调用显式传入的依赖,不一定等于 manager 构造时保存的实例。这允许账户切换等宿主 流程在恢复时选择当前认证状态。历史负责恢复 source、lineage 和会话语义;当前 Config 仍决定这次 运行使用的模型、权限、cwd、feature 与其他运行时能力。
V2 agent metadata 的恢复只在非 subagent source 上执行。child agent 的 control/lineage 由带 source 的 内部恢复入口和父级 AgentControl 负责,不能把 root 的恢复逻辑无条件套到每个 descendant。
下面的时序图把路径读取、active-map 快路径和 Session 重建放在同一条时间线上。它强调 store 读取 发生在 map 判断之前,而 resume lifecycle 发生在注册之后。
时序中的“返回原 handle”不会再次执行 reconstruction;“返回新 handle”才会重新打开 writer、重建 SessionState 并触发 extension resume callback。
3. rollout path
路径入口始终请求 include_archived: true 和 include_history: true。这意味着用户可以显式恢复已归档 Thread,同时要求 store 返回可重建 Session 的完整历史。
源码位置:codex-rs/core/src/thread_manager.rs :: initial_history_from_rollout_path, stored_thread_to_initial_history
async fn initial_history_from_rollout_path(
&self,
rollout_path: PathBuf,
) -> CodexResult<InitialHistory> {
let requested_rollout_path = rollout_path.clone();
let stored_thread = self
.state
.thread_store
.read_thread_by_rollout_path(ReadThreadByRolloutPathParams {
rollout_path,
include_archived: true,
include_history: true,
})
.await
.map_err(thread_store_rollout_read_error)?;
stored_thread_to_initial_history(stored_thread, Some(requested_rollout_path))
}
fn stored_thread_to_initial_history(
stored_thread: StoredThread,
rollout_path: Option<PathBuf>,
) -> CodexResult<InitialHistory> {
let thread_id = stored_thread.thread_id;
// 调用方请求了完整历史却没有得到 history,不能退化成空会话继续运行。
let history = stored_thread.history.ok_or_else(|| {
CodexErr::Fatal(format!(
"thread {thread_id} did not include persisted history"
))
})?;
Ok(InitialHistory::Resumed(ResumedHistory {
conversation_id: thread_id,
history: Arc::new(history.items),
// 用户请求路径优先,store path 只在入口未提供路径时作为后备。
rollout_path: rollout_path.or(stored_thread.rollout_path),
}))
}Local store 对相对路径先拼接 codex_home,随后拒绝目录、非普通文件和不存在的路径,最后执行 canonicalize。store 内部用解析后的实际文件打开 writer;但 ResumedHistory 仍保留调用方请求路径, active-map fast path 也直接拿这个值与运行 handle 的 rollout path 比较。因此调用方应沿用 handle/store 返回的规范路径;同一文件的另一种相对路径或别名仍可能被判定为“不同 rollout path”。
源码位置:codex-rs/thread-store/src/local/read_thread.rs :: resolve_requested_rollout_path
async fn resolve_requested_rollout_path(
store: &LocalThreadStore,
rollout_path: std::path::PathBuf,
) -> ThreadStoreResult<std::path::PathBuf> {
let path = if rollout_path.is_relative() {
store.config.codex_home.join(rollout_path)
} else {
rollout_path
};
match tokio::fs::metadata(path.as_path()).await {
Ok(metadata) if metadata.is_dir() => {
// 目录不能被猜测成“其中最新的 rollout”,目标必须唯一。
return Err(ThreadStoreError::InvalidRequest {
message: format!(
"failed to resolve rollout path `{}`: path is a directory",
path.display()
),
});
}
Ok(metadata) if !metadata.is_file() => {
return Err(ThreadStoreError::InvalidRequest {
message: format!(
"failed to resolve rollout path `{}`: path is not a file",
path.display()
),
});
}
_ => {}
}
let Some(path) = codex_rollout::existing_rollout_path(path.as_path()).await else {
return Err(ThreadStoreError::InvalidRequest {
message: format!(
"failed to resolve rollout path `{}`: file does not exist",
path.display()
),
});
};
// canonical path 让 writer ownership 与同文件别名收敛到同一目标。
std::fs::canonicalize(path.as_path()).map_err(|err| ThreadStoreError::InvalidRequest {
message: format!("failed to resolve rollout path `{}`: {err}", path.display()),
})
}按路径读取还会把 rollout 中的 legacy metadata 与 SQLite metadata 合并:paginated Thread 的显示信息 主要来自 SQLite;legacy Thread 保留 rollout 派生内容,同时补充 section、recency 和 Git 字段。
这里存在一个明确限制:include_history: true 会拒绝 ThreadHistoryMode::Paginated,返回 Unsupported { operation: "paginated_threads" }。原因是 paginated rollout 可能只保存历史后缀,不能把 单个文件当成完整 eager history。此类 Thread 应由已经取得正确历史范围的上层调用 resume_thread_with_history(),而不是强行走 rollout-path eager API。
4. InitialHistory
ResumedHistory 只包含三个字段,却承担恢复路径的关键交接:原 ThreadId、不可变共享的 rollout items, 以及可选持久化路径。
源码位置:codex-rs/protocol/src/protocol.rs :: ResumedHistory, InitialHistory helpers
pub struct ResumedHistory {
pub conversation_id: ThreadId,
pub history: Arc<Vec<RolloutItem>>,
pub rollout_path: Option<PathBuf>,
}
pub enum InitialHistory {
New,
Cleared,
Resumed(ResumedHistory),
Forked(Vec<RolloutItem>),
}
impl InitialHistory {
pub fn get_history_mode(&self, default_history_mode: ThreadHistoryMode) -> ThreadHistoryMode {
match self {
InitialHistory::New | InitialHistory::Cleared => default_history_mode,
// 恢复与分叉必须保留源历史格式,不能按目标配置静默转换。
InitialHistory::Resumed(_) | InitialHistory::Forked(_) => self
.get_session_meta()
.map(|meta| meta.history_mode)
.unwrap_or(default_history_mode),
}
}
pub fn get_resumed_session_sources(&self) -> Option<(SessionSource, Option<ThreadSource>)> {
let meta = self.get_resumed_session_meta()?;
// 产品 SessionSource 与用户/子 agent ThreadSource 分开恢复。
Some((meta.source.clone(), meta.thread_source.clone()))
}
pub fn get_resumed_parent_thread_id(&self) -> Option<ThreadId> {
self.get_resumed_session_meta()
.and_then(|meta| meta.parent_thread_id)
}
}其他 helper 还会从首个可用 SessionMeta 读取 base instructions、dynamic tools、selected capability roots、 multi-agent version、originator、fork parent 和 session identity。恢复不是把所有旧配置覆盖到当前 Config; 它只恢复需要保持会话连续性的语义字段,再由 Session 构造逻辑决定与当前配置的优先级。
Rollout loader 把文件中遇到的第一个 SessionMeta 作为 canonical ThreadId。后续 SessionMeta 可能来自 fork 时复制的历史,所以不能用“最后一个 metadata”决定当前 Thread 身份。普通损坏 JSON 行会记录 warning、增加 parse_errors 并跳过;空文件或始终没有可解析 ThreadId 的文件不能形成恢复信封。
恢复过程中几类同名对象属于不同层次。下面的类图展示数据从 store 快照进入协议信封,再分别投影到 运行 Session 和持久化 handle;它们不是一个对象的多种别名。
StoredThread 是一次读取结果,ResumedHistory 是不可变恢复输入,SessionState 是重新计算出的可变 运行状态,LiveThread 则负责恢复后的追加写入。区分这四层可以避免把 store metadata 误当成完整 Session snapshot。
5. active map
spawn_thread() 在任何昂贵 Session 初始化之前检查恢复 ID。三种分支分别是:复用运行中的相同 Thread、拒绝同 ID 不同 rollout、移除已经停止的旧 handle。
源码位置:codex-rs/core/src/thread_manager.rs :: ThreadManagerState::spawn_thread resumed fast path
let is_resumed_thread = matches!(&initial_history, InitialHistory::Resumed(_));
if let InitialHistory::Resumed(resumed) = &initial_history {
let mut threads = self.threads.write().await;
if let Some(thread) = threads.get(&resumed.conversation_id).cloned() {
if thread.is_running() {
if let Some(requested_rollout_path) = resumed.rollout_path.as_deref()
&& thread.rollout_path().as_deref() != Some(requested_rollout_path)
{
// 同一 ThreadId 指向不同 rollout 会造成身份与历史分裂,必须拒绝。
return Err(CodexErr::InvalidRequest(format!(
"thread {} is already running with a different rollout path",
resumed.conversation_id
)));
}
// 相同恢复请求返回同一个 Arc,不重复创建 Session、writer 或 event listener。
return Ok(NewThread {
thread_id: resumed.conversation_id,
session_configured: thread.session_configured(),
thread,
});
}
// stopped handle 不能接收新 Op;先从 registry 移除,再按原 ID 重建。
threads.remove(&resumed.conversation_id);
}
}检查期间持有 write lock,保证两个并发恢复不会都通过“map 中不存在”的判断。返回已有 handle 时不会 再次触发 rollout reconstruction 或 resume lifecycle;调用方得到的就是当前活动会话。
仓库测试直接验证这两个结果:恢复 active Thread 时 Arc::ptr_eq 为 true;先 shutdown 再恢复时 ThreadId 相同但 Arc::ptr_eq 为 false。这个区别比只断言 ID 更重要,它证明恢复的幂等边界位于 运行 handle,而不是业务身份。
active map 的决策可以归约成下面的状态机。状态名描述 manager 观察到的运行状态,终态描述本次恢复 请求的结果,而不是 Thread 永久生命周期。
只有 RunningSame 可以无副作用返回;Stopped 必须先清理 registry,RunningDifferent 则不能用 “后请求覆盖先请求”的方式自动纠正。
6. Session 保留 ID
进入新 Session 后,InitialHistory::Resumed 不分配 UUID。Core 直接使用 conversation_id,再从历史 首个 SessionMeta 恢复 SessionId。legacy subagent rollout 可能把自身 ThreadId 错当 SessionId,源码 会过滤该值并改为继承当前 AgentControl 的 tree identity。
源码位置:codex-rs/core/src/session/session.rs :: resumed thread_id and session_id
let thread_id = match &initial_history {
InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => {
ThreadId::default()
}
// 恢复必须复用原 conversation ID,否则 store、agent graph 与 UI 会分裂成两条 Thread。
InitialHistory::Resumed(resumed_history) => resumed_history.conversation_id,
};
let resumed_session_id = match &initial_history {
InitialHistory::Resumed(resumed) => {
resumed.history.iter().find_map(|item| match item {
RolloutItem::SessionMeta(meta_line) => Some(meta_line.meta.session_id),
_ => None,
})
}
InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => None,
};
let resumed_session_id = resumed_session_id.filter(|session_id| {
// legacy child 的自指 session ID 不能覆盖共享 agent-tree identity。
!session_configuration.session_source.is_non_root_agent()
|| *session_id != SessionId::from(thread_id)
});
let session_id = resumed_session_id.unwrap_or_else(|| {
if session_configuration.session_source.is_non_root_agent() {
agent_control.session_id()
} else {
SessionId::from(thread_id)
}
});非 ephemeral Session 随后调用 LiveThread::resume()。传入的 history 是 manager 已经加载的同一个 Arc<Vec<RolloutItem>>,所以 persistence 层不重复读文件;它只校准 metadata sync,并让具体 store 重新取得 writer ownership。
源码位置:codex-rs/core/src/session/session.rs :: ResumeThreadParams, LiveThread::resume
let params = ResumeThreadParams {
thread_id: resumed_history.conversation_id,
rollout_path: resumed_history.rollout_path.clone(),
// 历史已经由 ThreadManager/上层提供,LiveThread 不再执行第二次 load_history。
history: Some(resumed_history.history.clone()),
include_archived: true,
metadata: ThreadPersistenceMetadata {
// writer 使用本次运行的 cwd/provider/memory 配置继续追加,而不是盲用旧宿主环境。
cwd: Some(config.cwd.to_path_buf()),
model_provider: config.model_provider_id.clone(),
memory_mode: if config.memories.generate_memories {
ThreadMemoryMode::Enabled
} else {
ThreadMemoryMode::Disabled
},
},
};
LiveThread::resume(
Arc::clone(&thread_store),
session_configuration.history_mode,
params,
)
.await?Local store 在打开 recorder 前执行三层互斥:进程内 live_writer_locks、已有 recorder 检查,以及跨进程 writer lock。另一个进程仍拥有同一 rollout writer 时,恢复返回 conflict;原 owner shutdown 释放锁后, 新进程才可以接管并继续 append。
writer 使用哪个 ThreadHistoryMode 也不是只看当前配置:已有 history 时从 rollout items 计算 canonical mode;只有没有预加载 history 时才从 rollout/store metadata 读取。这样 resume 不会把 legacy 与 paginated 格式悄悄混写到同一个 Thread。
7. Rollout
Session 不会把所有 RolloutItem 原样塞进模型 history。reconstruct_history_from_rollout() 返回一组 相互关联的恢复结果:模型可见 ResponseItem、上一 Turn 设置、reference context、world-state baseline 和 auto-compaction window identity。
源码位置:codex-rs/core/src/session/rollout_reconstruction.rs :: RolloutReconstruction
pub(super) struct RolloutReconstruction {
// history 与恢复元数据由同一次 replay 产生,避免分别扫描后得到不一致快照。
pub(super) history: Vec<ResponseItem>,
pub(super) previous_turn_settings: Option<PreviousTurnSettings>,
pub(super) reference_context_item: Option<TurnContextItem>,
pub(super) world_state_baseline: Option<WorldStateSnapshot>,
pub(super) window_number: u64,
pub(super) first_window_id: Option<Uuid>,
pub(super) previous_window_id: Option<Uuid>,
pub(super) window_id: Option<Uuid>,
}算法先从新到旧扫描 rollout,把记录归入 Turn segment。这样可以尽早找到最新仍有效的 compaction replacement history、上一用户 Turn 设置和 reference context,然后只顺序重放 checkpoint 之后的后缀。
源码位置:codex-rs/core/src/session/rollout_reconstruction.rs :: reconstruct_history_from_rollout reverse scan(节选)
for (index, item) in rollout_items.iter().enumerate().rev() {
match item {
RolloutItem::Compacted(compacted) => {
let active_segment =
active_segment.get_or_insert_with(ActiveReplaySegment::default);
active_segment.world_state_replay.push(item);
if active_segment.base_replacement_history.is_none()
&& let Some(replacement_history) = &compacted.replacement_history
{
// 最新 surviving checkpoint 成为 history base,更早消息无需再次物化。
active_segment.base_replacement_history = Some(replacement_history);
rollout_suffix = &rollout_items[index + 1..];
}
}
RolloutItem::EventMsg(EventMsg::ThreadRolledBack(rollback)) => {
// 反向扫描把“删除最新 N 个用户 Turn”转换成跳过接下来 N 个 segment。
pending_rollback_turns = pending_rollback_turns
.saturating_add(usize::try_from(rollback.num_turns).unwrap_or(usize::MAX));
}
RolloutItem::TurnContext(ctx) => {
let active_segment =
active_segment.get_or_insert_with(ActiveReplaySegment::default);
if active_segment.turn_id.is_none() {
active_segment.turn_id = ctx.turn_id.clone();
}
if turn_ids_are_compatible(
active_segment.turn_id.as_deref(),
ctx.turn_id.as_deref(),
) {
// 只让同一 Turn 的 context 恢复 model/comp hash/realtime 设置。
active_segment.previous_turn_settings = Some(PreviousTurnSettings {
model: ctx.model.clone(),
comp_hash: ctx.comp_hash.clone(),
realtime_active: ctx.realtime_active,
});
}
}
RolloutItem::EventMsg(EventMsg::TurnStarted(event)) => {
// TurnStarted 是反向 segment 的最老边界,到达后才能安全 finalize。
if active_segment.as_ref().is_some_and(|active_segment| {
turn_ids_are_compatible(
active_segment.turn_id.as_deref(),
Some(event.turn_id.as_str()),
)
}) && let Some(active_segment) = active_segment.take()
{
finalize_active_segment(
active_segment,
&mut base_replacement_history,
&mut previous_turn_settings,
&mut reference_context_item,
&mut world_state_replay,
&mut window,
&mut pending_rollback_turns,
);
}
}
// ... 其他 RolloutItem 分支按同一 segment 规则处理。
}
}随后算法按时间正序重放 surviving suffix:ResponseItem 进入 ContextManager,inter-agent communication 转换为模型输入,新的 compaction replacement 直接替换历史,legacy compaction 重建压缩历史, ThreadRolledBack 删除最后 N 个用户 Turn。
world state 单独重放。full snapshot 建立 baseline,后续 merge patch 在其上应用;没有 full snapshot 的 patch 会被忽略。snapshot 反序列化或 merge patch 失败时记录 warning 并清空 baseline,而不是让全部 Thread 恢复失败。这是一种有意降级:模型历史仍可使用,但外部世界状态不再伪装成可信快照。
8. 重建结果安装
apply_rollout_reconstruction() 先处理历史图片与音频,再在一个 SessionState 临界区里替换 history、 world-state baseline、auto-compaction window 和 previous-turn settings。
源码位置:codex-rs/core/src/session/mod.rs :: Session::apply_rollout_reconstruction
let rollout_reconstruction::RolloutReconstruction {
mut history,
previous_turn_settings,
reference_context_item,
world_state_baseline,
window_number,
first_window_id,
previous_window_id,
window_id,
} = self
.reconstruct_history_from_rollout(turn_context, rollout_items)
.await;
// replay 不回填新的图片 resize notice,避免修改历史模型前缀。
let _ = prepare_image_response_items(&mut history, ImageResizeNoticeMode::Disabled);
prepare_audio_response_items(&mut history);
{
// 相关恢复字段在同一锁内提交,外部不会看到“新 history + 旧 window”的中间状态。
let mut state = self.state.lock().await;
state.replace_history(history, reference_context_item);
if let Some(world_state) = world_state_baseline {
state.history.set_world_state_baseline(world_state);
}
let fallback_ids = state.auto_compact_window_ids();
let window_id = window_id.unwrap_or(fallback_ids.window_id);
state.restore_auto_compact_window(
window_number,
AutoCompactWindowIds {
first_window_id: first_window_id.unwrap_or(window_id),
previous_window_id,
window_id,
},
);
state.set_previous_turn_settings(previous_turn_settings.clone());
}若当前 token-limit scope 是 BodyAfterPrefix,安装后还会根据重建 history 与 base instructions 估算 prefill token,并写回 auto-compact window。恢复旧 window ID 的目的不是展示历史编号,而是让下一次 compaction 延续正确 lineage 和预算窗口。
record_initial_history() 还会从 rollout 最后一个 TokenCount 事件恢复 token usage,UI 因此可以在 新 Turn 开始前显示已有消耗。如果历史最后使用的 model 与当前 model 不同,Core 发送 warning,但允许 恢复继续;这是性能/上下文兼容风险,不是结构损坏。
9. 配置完成事件
恢复不会把历史事件重新发给客户端,也不会用 ThreadResumed 取代启动握手。新 Session 仍先排队 INITIAL_SUBMIT_ID 的 SessionConfigured,随后才调用 record_initial_history();manager 校验首事件并 注册 handle 后,再触发 thread-resume extension lifecycle。
源码位置:codex-rs/core/src/thread_manager.rs :: spawn_thread post-registration resume lifecycle
let new_thread = self
.finalize_thread_spawn(session, io, tracked_session_source)
.await?;
if source_changed_during_startup.load(Ordering::Acquire) {
new_thread.thread.session.request_mcp_runtime_refresh();
}
if is_resumed_thread {
// extension 的 resume callback 只在 handle 已通过首事件检查并进入 registry 后执行。
new_thread.thread.emit_thread_resume_lifecycle().await;
}
Ok(new_thread)普通 root 恢复在 history reconstruction 后会 flush rollout;subagent 跳过这次 flush,避免每个 child startup 独立引入不必要的持久化等待。next_turn_is_first 则根据历史中是否已有用户 Turn 恢复,不能把 “新 Session 对象”误判成“会话的第一轮用户交互”。
10. 失败清理分三层
恢复的错误不是都发生在同一个位置,也不能全部用“rollout 损坏”解释:
| 层次 | 典型失败 | 结果与清理 |
|---|---|---|
| 路径解析 | 不存在、目录、非普通文件、canonicalize 失败 | InvalidRequest,尚未创建 Session/writer |
| Store 读取 | Thread 不存在、metadata/ThreadId 不一致、paginated eager history 不支持 | 映射为 ThreadNotFound、InvalidRequest 或 fatal store error |
| 历史信封 | 请求完整历史却返回 history: None | CodexErr::Fatal,不创建空会话 |
| active map | 同 ID 已运行且路径一致 | 返回原 Arc,不属于失败 |
| active map | 同 ID 已运行但路径不同 | InvalidRequest,保护身份唯一性 |
| writer ownership | 进程内重复 recorder、跨进程 writer lock 冲突 | Session startup 失败,原 owner 保持运行 |
| 状态重建 | world-state JSON/patch 无效 | warning 并丢弃 world baseline,history 继续恢复 |
| Session 后续初始化 | shell、Hook、network、MCP 或 extension 初始化失败 | LiveThreadInitGuard discard 已打开的 persistence |
| 首事件/注册 | 首事件不是 SessionConfigured、并发注册冲突 | 关闭未注册 Session,避免孤儿 submission loop |
LiveThreadInitGuard 同时覆盖 create 与 resume。只有 Session 初始化完全成功才 commit();任一后续 步骤返回错误都会 discard().await,释放刚取得的 writer ownership。意外 drop 时 guard 也会尝试在 当前 Tokio runtime 异步 discard。
源码位置:codex-rs/core/src/session/session.rs :: LiveThreadInitGuard commit/discard
let session_result: anyhow::Result<Arc<Self>> = async {
// ... SessionConfigured、MCP、历史重建与启动状态安装
Ok(sess)
}
.await;
match session_result {
Ok(sess) => {
// commit 发生在所有 fallible 初始化之后,writer 才正式归 Session 所有。
live_thread_init.commit();
Ok(sess)
}
Err(err) => {
// resume 已经打开的 writer 必须释放,否则下一次恢复会误报重复 owner。
live_thread_init.discard().await;
Err(err)
}
}11. 恢复阶段定位
| 现象 | 优先检查 |
|---|---|
| 恢复后得到的还是原运行会话 | active map fast path、Arc::ptr_eq,这是预期幂等行为 |
| 同一 ID 报不同 rollout path | 调用方是否用别名或错误历史指向已运行 Thread |
| rollout 能列出但不能按路径恢复 | history mode 是否为 paginated、是否请求 eager complete history |
| 恢复后模型上下文缺少旧消息 | compaction checkpoint、rollback segment、ResponseItem 是否进入 surviving suffix |
| UI token usage 归零 | rollout 末尾是否有 TokenCount,以及 reconstruction 后是否 seed state |
| 恢复后 model warning | 历史 TurnContext model 与当前 Config model 不同 |
| world state 丢失但对话仍在 | full snapshot/patch 是否损坏,日志是否出现 replay warning |
| 第二个进程无法恢复 | 原进程 writer 是否 shutdown,跨进程 writer lock 是否仍被持有 |
| 初始化失败后下一次仍报 duplicate writer | LiveThreadInitGuard discard 或 store shutdown 是否完成 |
验证恢复语义时至少覆盖四组测试:active handle 复用、stopped handle 重建、ThreadSource/SessionId 保留, 以及 rollout reconstruction 对 compaction、rollback、world state 和 previous-turn settings 的归约。只验证 “返回相同 ThreadId”无法证明 writer、history 与运行状态真的恢复正确。
12. 恢复边界测试
12.1 Replacement历史
Compaction记录带 replacement_history 时,reconstruction应直接采用该 Vec作为历史基底,并恢复 window number、first/previous/current window ID。测试故意在 replacement中放入 summary和“stale developer instructions”,验证恢复不会在此阶段擅自清洗或重建它。
源码位置:codex-rs/core/src/session/tests.rs
// :: reconstruct_history_uses_replacement_history_verbatim(核心断言)
let replacement_history = vec![
summary_item.clone(),
ResponseItem::Message {
id: None,
role: "developer".to_string(),
content: vec![ContentItem::InputText {
text: "stale developer instructions".to_string(),
}],
phase: None,
internal_chat_message_metadata_passthrough: None,
},
];
let rollout_items = vec![RolloutItem::Compacted(CompactedItem {
message: String::new(),
replacement_history: Some(replacement_history.clone()),
window_number: Some(42),
first_window_id: Some(first_window_id.to_string()),
previous_window_id: Some(previous_window_id.to_string()),
window_id: Some(window_id.to_string()),
})];
let reconstructed = session
.reconstruct_history_from_rollout(&turn_context, &rollout_items)
.await;
// verbatim断言证明checkpoint是完整历史基底,不是摘要提示的再计算输入。
assert_eq!(reconstructed.history, replacement_history);
assert_eq!(42, reconstructed.window_number);
assert_eq!(Some(first_window_id), reconstructed.first_window_id);
assert_eq!(Some(previous_window_id), reconstructed.previous_window_id);
assert_eq!(Some(window_id), reconstructed.window_id);该测试不覆盖 checkpoint之后的 surviving suffix;suffix顺序由更完整的 compaction/rollback tests验证。
12.2 正确性边界
测试先从当前 TurnContext构造 WorldState,把其完整模型可见 fragments和 full snapshot一起放入已完成 Turn rollout;Resume后再次调用 context update。若 baseline恢复正确,history必须保持原 fragments,不追加第二 份相同环境、权限或AGENTS上下文。
源码位置:codex-rs/core/src/session/rollout_reconstruction_tests.rs
// :: record_initial_history_restores_world_state_baseline(断言节选)
let world_state = build_world_state_from_turn_context(&session, &turn_context).await;
let expected_history = world_state
.render_full()
.into_iter()
.map(ContextualUserFragment::into_boxed_response_item)
.collect::<Vec<_>>();
let mut world_state_items = expected_history
.iter()
.cloned()
.map(RolloutItem::ResponseItem)
.collect::<Vec<_>>();
world_state_items.push(RolloutItem::WorldState(WorldStateItem::full(
world_state.snapshot().into_value(),
)));
let rollout_items = completed_user_turn_rollout(
turn_context.to_turn_context_item(),
world_state_items,
);
session
.record_initial_history(InitialHistory::Resumed(ResumedHistory {
conversation_id: ThreadId::default(),
history: Arc::new(rollout_items),
rollout_path: Some(PathBuf::from("/tmp/resume.jsonl")),
}))
.await;
let step_context = StepContext::for_test(Arc::clone(&turn_context));
session
.record_context_updates_and_set_reference_context_item(&step_context)
.await
.expect("world state should build");
// 没有重复items证明full snapshot已成为下一次diff的可信baseline。
assert_eq!(
session.clone_history().await.raw_items(),
expected_history.as_slice(),
);它证明合法 full snapshot恢复;没有测试直接写入无法反序列化的 full state、没有 baseline的 patch或无法应用 的 merge patch。源码对三种情况会 warning并清空 baseline,使下一次 context update回退完整注入,但这条 降级目前缺少专门的反向测试。
12.3 Writer恢复
LocalThreadStore同时防止同进程重复 recorder和跨进程 writer lock冲突。最直接的测试先 create一个 live writer,再用相同 ThreadId和 rollout path调用 resume;第二次恢复必须失败,并给出“already has a live local writer”。
源码位置:codex-rs/thread-store/src/local/mod.rs
// :: resume_thread_rejects_duplicate_live_writer
let thread_id = ThreadId::default();
store
.create_thread(create_thread_params(thread_id))
.await
.expect("create live thread");
let rollout_path = store
.live_rollout_path(thread_id)
.await
.expect("live rollout path");
let err = store
.resume_thread(ResumeThreadParams {
thread_id,
rollout_path: Some(rollout_path),
history: None,
include_archived: true,
metadata: thread_metadata(),
})
.await
.expect_err("duplicate live resume should fail");
// 原writer保持owner,失败恢复不能抢占或替换它。
assert!(matches!(err, ThreadStoreError::InvalidRequest { .. }));
assert!(err.to_string().contains("already has a live local writer"));这证明 store owner约束,不等于 manager active-map幂等路径。正常 manager Resume会先发现 running handle并 返回同一 Arc,不会故意走到这个错误;跨进程或绕过manager的竞争才由 writer lock负责。
13. 恢复身份验证
- 从
resume_thread_from_rollout()开始,说明路径如何进入 ThreadStore、形成InitialHistory::Resumed,以及 active map在哪一步决定复用还是重建。 - 解释 replacement history、surviving suffix、TurnContextItem和WorldState baseline分别恢复什么,为什么 “对话消息存在”不能证明动态上下文恢复正确。
- 给定“恢复后环境上下文重复出现”,判断应检查 WorldState baseline还是 writer ownership,并说明合法 baseline测试能证明什么。
- 使用只读命令定位本文三组证据和损坏降级分支:
rg -n "replacement_history|world_state_baseline|failed to apply world-state" \
codex-rs/core/src/session
rg -n "reconstruct_history_uses_replacement|restores_world_state_baseline|duplicate_live_writer" \
codex-rs/core/src/session codex-rs/thread-store/src/local/mod.rs继续学习 ThreadManager分叉 时,对比 Resume保留 ThreadId与Fork分配 新身份的差异;需要理解重建结果如何写入 SessionState时,可回看 Session启动上下文 的 reference context和WorldState章节。
