Skip to content

ThreadManager恢复

追踪 ThreadManager 从 rollout 或预加载历史恢复 Thread,解释身份复用、持久化重开、状态重建与失败清理。

基于rust-v0.150.0
CodexRustThreadManagerSession

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 结果
新建新 UUIDv7InitialHistory::New创建新 Thread新 Arc<CodexThread>
恢复保留 rollout 中的 IDInitialHistory::Resumed重开原 Thread writer复用活跃 handle,或重建已停止 handle
分叉新 UUIDv7InitialHistory::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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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(节选)

rust
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

rust
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

rust
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: NoneCodexErr::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

rust
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 writerLiveThreadInitGuard 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

rust
// :: 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

rust
// :: 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

rust
// :: 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. 恢复身份验证 ​

  1. 从 resume_thread_from_rollout() 开始,说明路径如何进入 ThreadStore、形成 InitialHistory::Resumed,以及 active map在哪一步决定复用还是重建。
  2. 解释 replacement history、surviving suffix、TurnContextItem和WorldState baseline分别恢复什么,为什么 “对话消息存在”不能证明动态上下文恢复正确。
  3. 给定“恢复后环境上下文重复出现”,判断应检查 WorldState baseline还是 writer ownership,并说明合法 baseline测试能证明什么。
  4. 使用只读命令定位本文三组证据和损坏降级分支:
bash
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章节。