Skip to content

ThreadRead历史加载

沿 thread/read 的来源选择,解释未落盘摘要、实时 Store 历史、Legacy 事件归约、Paginated 全量补取及残留 Turn 状态的修正。

基于rust-v0.150.0
CodexRustAppServerThreadStoreHistory

ThreadRead历史加载 ​

客户端刚创建一个 Thread,就能读取到它的 path,但对应文件还不存在;把同一次读取改成 includeTurns=true,又可能收到“首条用户消息前不可用”的错误。另一种 Thread 的 path 始终为 null,却能正常返回历史。这三个现象并不矛盾:运行对象、持久化元数据、历史内容具有不同的可用时机,路径也不是所有存储后端共有的身份。

本文面向掌握 Rust Option、Result 与异步调用、开始阅读 App Server 的读者。可先读ThreadStart处理流程了解运行对象的创建,再读ThreadList与分页理解列表中的摘要。这里把范围收在单 Thread 读取:从请求分派追踪到 ThreadReadResponse,重点解释返回字段来自哪里、历史如何构造、失败时保留什么状态。恢复运行和订阅事件属于另一条流程,本文只在需要消歧时对照。

读代码时始终区分三个对象:Core 的 CodexThread 是进程内运行句柄;Store 的 StoredThread 是持久化视图;协议的 Thread 是本次响应的值。协议 Turn 表示一次任务轮次,ThreadItem 是轮次内用户消息、工具调用等显示条目。它们都不等同于终端输入检索使用的 history.jsonl。文中代码摘录保留原实现,跨区段省略以 // ... 标明;解释写在相邻正文中。

1. 读取契约 ​

thread/read 的请求只有线程 ID 和是否携带轮次两个业务字段。先读协议的默认值,再读服务端,才能判断空数组的含义。

源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs

相关函数/类型:ThreadReadParams / ThreadReadResponse。

rust
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct ThreadReadParams {
    pub thread_id: String,
    /// When true, include turns and their items from rollout history.
    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
    pub include_turns: bool,
}

#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct ThreadReadResponse {
    pub thread: Thread,
}

include_turns 缺省为 false;skip_serializing_if 只是发送方省略 false 的编码规则,不会让服务端默认读取历史。响应永远包在 thread 字段中,未请求历史时 thread.turns=[],不能据此推断这个 Thread 从未执行过任务。请求也没有分页 cursor;当它返回完整历史时,分页工作发生在服务端内部。

下面的对象关系只表达本次读取持有或构造的数据,不表示读取会创建新的 Core 会话。

Thread 的字段合成是本文主线:存储摘要负责历史身份和展示字段,运行对象只覆盖部分现场信息,watch manager 负责最后的运行状态。读取返回独立的响应值,不把 StoredThread 或 Core 的内部对象直接暴露给客户端。

源码文件:codex-rs/app-server/src/message_processor.rs

相关函数/类型:MessageProcessor::process_request。

rust
ClientRequest::ThreadRead { params, .. } => {
    self.thread_processor.thread_read(params).await
}

分派进入专门的 Thread 处理器;它通过通用结果通道交回响应,不自行发送 thread/started。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:thread_read / thread_read_response_inner。

rust
pub(crate) async fn thread_read(
    &self,
    params: ThreadReadParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
    self.thread_read_response_inner(params)
        .await
        .map(|response| Some(response.into()))
}
// ...
async fn thread_read_response_inner(
    &self,
    params: ThreadReadParams,
) -> Result<ThreadReadResponse, JSONRPCErrorError> {
    let ThreadReadParams {
        thread_id,
        include_turns,
    } = params;

    let thread_uuid = ThreadId::from_string(&thread_id)
        .map_err(|err| invalid_request(format!("invalid thread id: {err}")))?;

    let thread = self
        .read_thread_view(thread_uuid, include_turns)
        .await
        .map_err(thread_read_view_error)?;
    Ok(ThreadReadResponse { thread })
}

这里有两次类型转换:字符串 ID 先变成 ThreadId,读取错误再映射成 JSON-RPC 错误。外层 .map(|response| Some(response.into())) 表示“已经有可发送的结果”,不是“找不到则返回 null”。线程不存在属于错误路径。

源码文件:codex-rs/app-server/src/message_processor.rs

相关函数/类型:MessageProcessor::process_request 的响应出口。

rust
match result {
    Ok(Some(response)) => {
        self.outgoing
            .send_response_as(request_id.clone(), response)
            .await;
    }
    Ok(None) => {}
    Err(error) => {
        self.outgoing.send_error(request_id.clone(), error).await;
    }
}
Ok(())

这一出口保留原请求 ID。Ok(Some(...)) 入响应队列,Err(...) 入错误队列;它证明处理器的交付边界到队列为止,不能证明网络对端已收到消息。读取主链没有提交模型 Op,也不安排新的 Turn。

2. 来源裁决 ​

首先问 Core 中是否已经存在运行句柄,然后问 Store 能否提供元数据。get_thread 的失败被 .ok() 变为未加载;Store 的任意错误则不能同样吞掉。两种探测具有不同的错误含义。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:read_thread_view 的来源选择。

rust
async fn read_thread_view(
    &self,
    thread_id: ThreadId,
    include_turns: bool,
) -> Result<Thread, ThreadReadViewError> {
    let loaded_thread = self.thread_manager.get_thread(thread_id).await.ok();
    let mut thread = if include_turns {
        if let Some(loaded_thread) = loaded_thread.as_ref() {
            // Loaded thread with turns: use persisted metadata when it exists,
            // but reconstruct turns from the live ThreadStore history.
            let persisted_thread = self
                .load_persisted_thread_for_read(thread_id, /*include_turns*/ false)
                .await?;
            self.load_live_thread_view(
                thread_id,
                include_turns,
                loaded_thread,
                persisted_thread,
            )
            .await?
        } else if let Some(thread) = self
            .load_persisted_thread_for_read(thread_id, include_turns)
            .await?
        {
            // Unloaded thread with turns: load metadata and history together
            // from the ThreadStore.
            thread
        } else {
            return Err(ThreadReadViewError::InvalidRequest(format!(
                "thread not loaded: {thread_id}"
            )));
        }
    } else if let Some(thread) = self
        .load_persisted_thread_for_read(thread_id, include_turns)
        .await?
    {
        if let Some(loaded_thread) = loaded_thread.as_ref() {
            self.load_live_thread_view(thread_id, include_turns, loaded_thread, Some(thread))
                .await?
        } else {
            thread
        }
    } else if let Some(loaded_thread) = loaded_thread.as_ref() {
        // Loaded metadata-only read before persistence is materialized: build
        // the response from the live thread snapshot.
        self.load_live_thread_view(
            thread_id,
            include_turns,
            loaded_thread,
            /*persisted_thread*/ None,
        )
        .await?
    } else {
        return Err(ThreadReadViewError::InvalidRequest(format!(
            "thread not loaded: {thread_id}"
        )));
    };

这段代码值得按四格而不是逐行记忆。includeTurns=true 且已加载时,先取持久化摘要,再从 live Thread 的持久化接口取历史;true 且未加载时,一起从 Store 取得需要的视图。false 时,优先持久化摘要,若没有摘要但运行句柄还在,就用运行快照构造回退值。两处 thread not loaded 都表示“持久化来源也未命中、运行来源也不能满足本次读取”,不能把它解释成必须先调用 resume。

图中每一条成功分支最终都还要经过状态叠加,元数据命中并不是最终返回点。

这张分支图还有一个隐含的排错价值:当摘要已经读取失败,流程不会假装摘要不存在、改用 Core 快照。只有下层明确报告目标缺失,才允许回退。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:read_stored_thread_for_read。

rust
async fn read_stored_thread_for_read(
    &self,
    thread_id: ThreadId,
    include_history: bool,
) -> Result<Option<StoredThread>, ThreadReadViewError> {
    match self
        .thread_store
        .read_thread(StoreReadThreadParams {
            thread_id,
            include_archived: true,
            include_history,
        })
        .await
    {
        Ok(stored_thread) => Ok(Some(stored_thread)),
        Err(ThreadStoreError::InvalidRequest { message })
            if message == format!("no rollout found for thread id {thread_id}") =>
        {
            Ok(None)
        }
        Err(ThreadStoreError::ThreadNotFound {
            thread_id: missing_thread_id,
        }) if missing_thread_id == thread_id => Ok(None),
        Err(ThreadStoreError::InvalidRequest { message }) => {
            Err(ThreadReadViewError::InvalidRequest(message))
        }
        Err(ThreadStoreError::Unsupported { operation }) => {
            Err(ThreadReadViewError::Unsupported(operation))
        }
        Err(err) => Err(ThreadReadViewError::Internal(format!(
            "failed to read thread: {err}"
        ))),
    }
}

include_archived: true 使归档 Thread 可以按 ID 检查;读取并不解除归档。旧式 “no rollout found” 文本与新式 ThreadNotFound 都可转为 None,后者还核对缺失 ID 是否就是当前目标。其他非法请求、后端不支持和内部错误分别保留。这说明 None 是受限的“来源没有数据”,不是所有存储异常的统一兜底。

3. 摘要合成 ​

3.1 持久化字段 ​

Store 可以保存比客户端需要更多的信息。转协议对象时,历史被拆到旁边,状态先设成 NotLoaded,turns 先置空;这只是中间值。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:thread_from_stored_thread。

rust
let cwd = AbsolutePathBuf::relative_to_current_dir(path_utils::normalize_for_native_workdir(
    thread.cwd,
))
.unwrap_or_else(|err| {
    warn!("failed to normalize thread cwd while reading stored thread: {err}");
    fallback_cwd.clone()
});
let source = with_thread_spawn_agent_metadata(
    thread.source,
    thread.agent_nickname.clone(),
    thread.agent_role.clone(),
);
let history = thread.history;
let thread_id = thread.thread_id.to_string();
let thread = Thread {
    id: thread_id.clone(),
    extra: None,
    session_id: thread_id,
    forked_from_id: thread.forked_from_id.map(|id| id.to_string()),
    parent_thread_id: thread.parent_thread_id.map(|id| id.to_string()),
    preview: thread.preview,
    ephemeral: false,
    section: thread.section.map(|section| ThreadSection {
        id: section.id,
        name: section.name,
        appearance: section
            .appearance
            .map(|appearance| ThreadSectionAppearance {
                icon: appearance.icon,
                color: appearance.color,
            }),
    }),
    section_entered_at: thread
        .section_entered_at
        .map(|entered_at| entered_at.timestamp()),
    project_id: thread.project_id,
    history_mode: thread.history_mode.into(),
    model_provider: if thread.model_provider.is_empty() {
        fallback_provider.to_string()
    } else {
        thread.model_provider
    },
    created_at: thread.created_at.timestamp(),
    updated_at: thread.updated_at.timestamp(),
    recency_at: Some(thread.recency_at.timestamp()),
    status: ThreadStatus::NotLoaded,
    path,
    cwd,
    cli_version: thread.cli_version,
    agent_nickname: source.get_nickname(),
    agent_role: source.get_agent_role(),
    source: source.into(),
    can_accept_direct_input: None,
    thread_source: thread.thread_source.map(Into::into),
    git_info,
    name: thread.name,
    turns: Vec::new(),
};
(thread, history)

session_id 暂用持久化线程 ID,稍后有 live 对象时会替换为本次 Session 的 ID。model_provider 为空才使用服务器配置中的 fallback provider;cwd 归一化失败才回退服务器 cwd,因此 fallback 是修复不完整数据的规则,不代表用当前配置覆盖所有历史。创建、更新和 recency 时间从 Store 的时间值转换为秒,不能用此次读取的时间代替。

Local Store 的摘要也不是单纯的 SQLite 行。下面这段展示 Legacy 摘要合成中最容易忽略的优先级。

源码文件:codex-rs/thread-store/src/local/read_thread.rs

相关函数/类型:read_thread。

rust
let thread_id = params.thread_id;
if let Some(metadata) = read_sqlite_metadata(store, thread_id).await
    && (params.include_archived
        || (metadata.archived_at.is_none()
            && !rollout_path_is_archived(
                store.config.codex_home.as_path(),
                metadata.rollout_path.as_path(),
            )))
    && (!params.include_history
        || sqlite_rollout_path_can_load_history_for_thread(&metadata.rollout_path, thread_id)
            .await)
{
    let metadata_sandbox_policy = metadata.sandbox_policy.clone();
    let mut thread = stored_thread_from_sqlite_metadata(store, metadata).await?;
    // Paginated history may contain only a suffix, so its display metadata lives in SQLite.
    // Legacy display metadata remains rollout-derived.
    if thread.history_mode == ThreadHistoryMode::Legacy
        && !params.include_history
        && let Some(rollout_path) = thread.rollout_path.clone()
        && let Ok(mut rollout_thread) = read_thread_from_rollout_path(store, rollout_path).await
        && rollout_thread.thread_id == thread_id
        && (params.include_archived || rollout_thread.archived_at.is_none())
        && !rollout_thread.preview.is_empty()
    {
        rollout_thread.recency_at = thread.recency_at;
        rollout_thread.section = thread.section;
        rollout_thread.section_position = thread.section_position;
        rollout_thread.section_entered_at = thread.section_entered_at;
        if !thread.cwd.as_os_str().is_empty() {
            rollout_thread.cwd = thread.cwd;
        }
        if thread.name.is_some() {
            rollout_thread.name = thread.name;
        }
        rollout_thread.project_id = thread.project_id;
        rollout_thread.git_info = thread.git_info;
        rollout_thread.permission_profile = permission_profile_from_metadata_value(
            &metadata_sandbox_policy,
            rollout_thread.cwd.as_path(),
        );
        thread = rollout_thread;
    }
    reject_paginated_history(&thread, params.include_history)?;
    attach_history_if_requested(&mut thread, params.include_history).await?;

SQLite 分支先检查归档条件;需要历史时还要求路径能加载正确 Thread 的历史。对 Legacy 的纯摘要读取,如果 rollout 中有非空 preview,就采用 rollout 派生的展示数据,同时保留 SQLite 中的 recency、section、cwd、name、project 与 git 等特定字段。Paginated 的 rollout 可能只剩后缀,因而展示元数据保持 SQLite 为主。不能从“同一个 Store 接口”推出“两种 history mode 的字段来源完全相同”。

源码文件:codex-rs/thread-store/src/local/read_thread.rs

相关函数/类型:sqlite_rollout_path_can_load_history_for_thread。

rust
async fn sqlite_rollout_path_can_load_history_for_thread(
    path: &std::path::Path,
    thread_id: codex_protocol::ThreadId,
) -> bool {
    if codex_rollout::existing_rollout_path(path).await.is_none() {
        return false;
    }
    // SQLite metadata can outlive a moved/recreated rollout path. When history is
    // requested, verify the path still resolves to the requested thread before
    // trusting it as the source replay.
    read_session_meta_line(path)
        .await
        .is_ok_and(|metadata| metadata.meta.id == thread_id)
}

检查不仅判断文件存在,还读取 SessionMeta 比较 ID。这是防止 SQLite 旧路径指向被移动或重新创建的其他 rollout。若没有这一比较,拿到的文件能解析也可能属于另一个 Thread。

3.2 运行快照 ​

当 Thread 刚启动但尚未物化,服务器仍有配置快照和预计算路径。回退构造器使用这些信息组成最小摘要。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:build_thread_from_snapshot / build_thread_from_loaded_snapshot。

rust
fn build_thread_from_snapshot(
    thread_id: ThreadId,
    session_id: String,
    multi_agent_version: Option<codex_protocol::protocol::MultiAgentVersion>,
    config_snapshot: &ThreadConfigSnapshot,
    path: Option<PathBuf>,
) -> Thread {
    let now = time::OffsetDateTime::now_utc().unix_timestamp();
    Thread {
        id: thread_id.to_string(),
        extra: None,
        session_id,
        forked_from_id: None,
        parent_thread_id: config_snapshot.parent_thread_id.map(|id| id.to_string()),
        preview: String::new(),
        ephemeral: config_snapshot.ephemeral,
        section: None,
        section_entered_at: None,
        project_id: None,
        history_mode: config_snapshot.history_mode.into(),
        model_provider: config_snapshot.model_provider_id.clone(),
        created_at: now,
        updated_at: now,
        recency_at: Some(now),
        status: ThreadStatus::NotLoaded,
        path,
        cwd: config_snapshot.cwd().clone(),
        cli_version: env!("CARGO_PKG_VERSION").to_string(),
        agent_nickname: config_snapshot.session_source.get_nickname(),
        agent_role: config_snapshot.session_source.get_agent_role(),
        source: config_snapshot.session_source.clone().into(),
        can_accept_direct_input: Some(can_accept_direct_input(
            multi_agent_version,
            &config_snapshot.session_source,
        )),
        thread_source: config_snapshot.thread_source.clone().map(Into::into),
        git_info: None,
        name: None,
        turns: Vec::new(),
    }
}
// ...
fn build_thread_from_loaded_snapshot(
    thread_id: ThreadId,
    config_snapshot: &ThreadConfigSnapshot,
    loaded_thread: &CodexThread,
) -> Thread {
    build_thread_from_snapshot(
        thread_id,
        loaded_thread.session_configured().session_id.to_string(),
        loaded_thread.multi_agent_version(),
        config_snapshot,
        loaded_thread.rollout_path(),
    )
}

此时 preview、name、git_info 可能为空,时间使用构造时的当前时间,不能据此断言磁盘记录已经建立。path 来自 loaded_thread.rollout_path(),只是可选的路径提示;Store 中持久身份仍是 ThreadId。一旦持久化摘要可用,读取会优先使用摘要中的历史时间和展示信息。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:load_live_thread_view。

rust
async fn load_live_thread_view(
    &self,
    thread_id: ThreadId,
    include_turns: bool,
    loaded_thread: &CodexThread,
    persisted_thread: Option<Thread>,
) -> Result<Thread, ThreadReadViewError> {
    let config_snapshot = loaded_thread.config_snapshot().await;
    if include_turns && config_snapshot.ephemeral {
        return Err(ThreadReadViewError::InvalidRequest(
            "ephemeral threads do not support includeTurns".to_string(),
        ));
    }
    let fallback_thread =
        build_thread_from_loaded_snapshot(thread_id, &config_snapshot, loaded_thread);
    let mut thread = if let Some(mut thread) = persisted_thread {
        if thread.path.is_none() {
            thread.path = fallback_thread.path.clone();
        }
        thread.session_id.clone_from(&fallback_thread.session_id);
        thread.ephemeral = fallback_thread.ephemeral;
        thread.can_accept_direct_input = fallback_thread.can_accept_direct_input;
        thread
    } else {
        fallback_thread
    };
    self.apply_thread_read_store_fields(thread_id, &mut thread, include_turns, loaded_thread)
        .await?;
    Ok(thread)
}

已存在持久化摘要时,live 快照只补缺失路径,并覆盖 session_id、ephemeral、can_accept_direct_input。它不会在这个函数里把持久化 cwd、preview 全部替换为快照字段。ephemeral 的完整历史请求在这里直接拒绝;服务器可以展示临时 Thread 的状态,并不承诺为它提供持久化历史接口。

4. 实时历史 ​

load_live_thread_view 最后进入 apply_thread_read_store_fields。这里“live”指从运行对象关联的 Store 读取,不等于把模型当前上下文直接塞到响应里。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:apply_thread_read_store_fields。

rust
async fn apply_thread_read_store_fields(
    &self,
    thread_id: ThreadId,
    thread: &mut Thread,
    include_turns: bool,
    loaded_thread: &CodexThread,
) -> Result<(), ThreadReadViewError> {
    self.attach_thread_name(thread_id, thread).await;

    if include_turns {
        if matches!(
            thread.history_mode,
            codex_app_server_protocol::ThreadHistoryMode::Paginated
        ) {
            self.thread_store
                .persist_thread(thread_id, PersistContext::Standard)
                .await
                .map_err(|err| thread_read_history_load_error(thread_id, err))?;
            thread.turns = self
                .paginated_thread_full_turns(thread_id)
                .await
                .map_err(ThreadReadViewError::JsonRpc)?;
            return Ok(());
        }
        let history = loaded_thread
            .load_history(/*include_archived*/ true)
            .await
            .map_err(|err| thread_read_history_load_error(thread_id, err))?;
        thread.turns = build_legacy_api_turns_from_rollout_items(&history.items);
    }

    Ok(())
}

名字单独从 Store 补齐。Legacy 经过 loaded_thread.load_history(true);Paginated 则先 persist_thread(Standard),再从投影分页读取。这个 await 是读取过程中的持久化屏障,失败要返回错误,不能继续拿旧投影伪装新历史。因此此 API 没有启动模型的副作用,但不能描述为底层绝对不写存储。

把 Legacy 路径展开,可以看到跨层边界上的数据始终围绕同一个 Thread ID。

只有 Local Store 才需要图中的 rollout 路径。其他后端可以直接返回存储记录,App Server 不应自行拼接 $CODEX_HOME/sessions/... 来代替接口。

源码文件:codex-rs/core/src/codex_thread.rs

相关函数/类型:CodexThread::load_history。

rust
pub async fn load_history(
    &self,
    include_archived: bool,
) -> ThreadStoreResult<StoredThreadHistory> {
    let live_thread = self
        .session
        .live_thread_for_persistence("load history")
        .map_err(|err| ThreadStoreError::Internal {
            message: err.to_string(),
        })?;
    live_thread.load_history(include_archived).await
}

Core 从 Session 中取得 LiveThread 持久化句柄,失败包装为 Store 内部错误。这个方法没有返回 Session 的提示词历史,也没有把模型上下文里的压缩替换历史直接当作 UI 轮次。

源码文件:codex-rs/thread-store/src/live_thread.rs

相关函数/类型:LiveThread::load_history。

rust
pub async fn load_history(
    &self,
    include_archived: bool,
) -> ThreadStoreResult<StoredThreadHistory> {
    self.thread_store
        .load_history(LoadThreadHistoryParams {
            thread_id: self.thread_id,
            include_archived,
        })
        .await
}

LiveThread 只传递 ID 和是否允许归档。资源所有者是具体 Store;这里没有打开文件句柄,也没有新增历史读取任务来长期占有 Thread。

源码文件:codex-rs/thread-store/src/local/mod.rs

相关函数/类型:LocalThreadStore::load_history。

rust
async fn load_history(
    &self,
    params: LoadThreadHistoryParams,
) -> ThreadStoreResult<StoredThreadHistory> {
    if let Ok(rollout_path) = live_writer::rollout_path(self, params.thread_id).await {
        if !params.include_archived
            && helpers::rollout_path_is_archived(
                self.config.codex_home.as_path(),
                rollout_path.as_path(),
            )
        {
            return Err(ThreadStoreError::InvalidRequest {
                message: format!("thread {} is archived", params.thread_id),
            });
        }
        return read_thread::read_thread_by_rollout_path(
            self,
            rollout_path,
            /*include_archived*/ true,
            /*include_history*/ true,
        )
        .await?
        .history
        .ok_or_else(|| ThreadStoreError::Internal {
            message: format!("failed to load history for thread {}", params.thread_id),
        });
    }

    read_thread::read_thread(
        self,
        ReadThreadParams {
            thread_id: params.thread_id,
            include_archived: params.include_archived,
            include_history: true,
        },
    )
    .await?
    .history
    .ok_or_else(|| ThreadStoreError::Internal {
        message: format!("failed to load history for thread {}", params.thread_id),
    })
}

Local Store 优先查询 live writer 的路径;找不到 writer 时才按 ID 解析当前记录。路径可能随归档移动,优先级有实际意义。该函数中没有无条件 flush,因此不能把 Legacy 的 load_history 和 Paginated 显式的 persist_thread 当成相同一致性保证。只要读取跨越多个 await,响应就不应被解释为整台服务器的原子快照。

5. 历史归约 ​

5.1 记录与条目 ​

Legacy 历史是 RolloutItem 序列,而 UI 需要按 Turn 组织的条目。两者之间有过滤和状态归约,绝不是 JSON 反序列化后原样返回。

源码文件:codex-rs/app-server/src/request_processors.rs

相关函数/类型:build_legacy_api_turns_from_rollout_items。

rust
pub(crate) fn build_legacy_api_turns_from_rollout_items(items: &[RolloutItem]) -> Vec<Turn> {
    let mut builder = ThreadHistoryBuilder::new();
    for item in items {
        if is_persisted_rollout_item(item, codex_protocol::protocol::ThreadHistoryMode::Legacy) {
            builder.handle_rollout_item(item);
        }
    }
    builder.finish()
}

先按 Legacy 持久化规则筛选,再喂给 ThreadHistoryBuilder。这能避免把不属于此历史模式的记录混进显示历史;它也提醒我们,直接调用通用 builder 的结果不一定等同于 App Server 的 Legacy API 结果。

源码文件:codex-rs/app-server-protocol/src/protocol/thread_history.rs

相关函数/类型:ThreadHistoryBuilder / handle_rollout_item / finish。

rust
pub struct ThreadHistoryBuilder {
    turns: Vec<Turn>,
    current_turn: Option<PendingTurn>,
    next_item_index: i64,
    current_rollout_index: usize,
    next_rollout_index: usize,
    active_change_set: Option<ThreadHistoryChangeSet>,
}
// ...
pub fn finish(mut self) -> Vec<Turn> {
    self.finish_current_turn();
    self.turns
}
// ...
pub fn handle_rollout_item(&mut self, item: &RolloutItem) {
    self.current_rollout_index = self.next_rollout_index;
    self.next_rollout_index += 1;
    match item {
        RolloutItem::EventMsg(event) => self.handle_event(event),
        RolloutItem::Compacted(payload) => self.handle_compacted(payload),
        RolloutItem::ResponseItem(item) => self.handle_response_item(&item.item),
        RolloutItem::InterAgentCommunication(_)
        | RolloutItem::InterAgentCommunicationMetadata { .. }
        | RolloutItem::TurnContext(_)
        | RolloutItem::WorldState(_)
        | RolloutItem::RealtimeItem(_)
        | RolloutItem::SecurityRiskScore(_)
        | RolloutItem::SessionMeta(_) => {}
    }
}

builder 拥有已完成的 turns 和尚未收口的 current_turn,索引用于合成或定位条目,finish(self) 消费 builder 并交出结果。SessionMeta、TurnContext 等记录不会直接形成屏幕上的条目;它们存在于 rollout 中,不意味着客户端的 ThreadItem 必须一一对应。

5.2 轮次终态 ​

TurnStarted 给出明确的边界和 ID。读这一小段时,注意先结束前一轮,再创建当前轮;否则晚到的条目可能落到错误轮次。

源码文件:codex-rs/app-server-protocol/src/protocol/thread_history.rs

相关函数/类型:handle_turn_started / handle_turn_complete。

rust
fn handle_turn_started(&mut self, payload: &TurnStartedEvent) {
    self.finish_current_turn();
    let turn = self
        .new_turn(Some(payload.turn_id.clone()))
        .with_status(TurnStatus::InProgress)
        .with_started_at(payload.started_at)
        .opened_explicitly();
    self.record_changed_pending_turn(&turn);
    self.current_turn = Some(turn);
}

fn handle_turn_complete(&mut self, payload: &TurnCompleteEvent) {
    let terminal_error = payload.error.as_ref().map(|error| V2TurnError {
        message: error.message.clone(),
        codex_error_info: error.codex_error_info.clone().map(Into::into),
        additional_details: None,
    });
    let apply_completion = |turn: &mut PendingTurn| {
        if let Some(error) = terminal_error.as_ref() {
            turn.status = TurnStatus::Failed;
            turn.error = Some(error.clone());
        } else if matches!(turn.status, TurnStatus::Completed | TurnStatus::InProgress) {
            turn.status = TurnStatus::Completed;
        }
        turn.completed_at = payload.completed_at;
        turn.duration_ms = payload.duration_ms;
        ThreadHistoryTurnChange::from_pending_turn(turn)
    };

    // Prefer an exact ID match from the active turn and then close it.
    if let Some(current_turn) = self
        .current_turn
        .as_mut()
        .filter(|turn| turn.id == payload.turn_id)
    {
        let changed_turn = apply_completion(current_turn);
        self.record_changed_turn(changed_turn);
        self.finish_current_turn();
        return;
    }

    if let Some(turn) = self
        .turns
        .iter_mut()
        .find(|turn| turn.id == payload.turn_id)
    {
        if let Some(error) = terminal_error.as_ref() {
            turn.status = TurnStatus::Failed;
            turn.error = Some(error.clone());
        } else if matches!(turn.status, TurnStatus::Completed | TurnStatus::InProgress) {
            turn.status = TurnStatus::Completed;
        }
        turn.completed_at = payload.completed_at;
        turn.duration_ms = payload.duration_ms;
        let changed_turn = ThreadHistoryTurnChange::from_turn(turn);
        self.record_changed_turn(changed_turn);
        return;
    }

    // If the completion event cannot be matched, apply it to the active turn.
    if let Some(current_turn) = self.current_turn.as_mut() {
        let changed_turn = apply_completion(current_turn);
        self.record_changed_turn(changed_turn);
        self.finish_current_turn();
    }
}

完成事件携带 terminal error 时归约为 Failed;没有 error 时,只有原状态为 Completed 或 InProgress 才转成 Completed。这个条件保留之前已经确定的失败或中断,不会因为后面出现一个无错误完成事件就无条件洗掉终态。精确匹配当前 Turn ID 后会立刻收口;接着查找已经进入历史的同 ID 轮次;若完成事件的 ID 仍未匹配,才回退到当前轮次并收口。因此“按 ID 归约”有兼容旧事件的回退边界,不是未知 ID 一律忽略。

源码文件:codex-rs/app-server-protocol/src/protocol/thread_history.rs

相关函数/类型:handle_error / handle_turn_aborted。

rust
fn handle_error(&mut self, payload: &ErrorEvent) {
    if !payload.affects_turn_status() {
        return;
    }
    let tracking_changes = self.is_tracking_changes();
    let changed_turn = if let Some(turn) = self.current_turn.as_mut() {
        turn.status = TurnStatus::Failed;
        turn.error = Some(V2TurnError {
            message: payload.message.clone(),
            codex_error_info: payload.codex_error_info.clone().map(Into::into),
            additional_details: None,
        });
        tracking_changes.then(|| ThreadHistoryTurnChange::from_pending_turn(turn))
    } else {
        None
    };
    if let Some(changed_turn) = changed_turn {
        self.record_changed_turn(changed_turn);
    }
}

fn handle_turn_aborted(&mut self, payload: &TurnAbortedEvent) {
    let apply_abort = |turn: &mut PendingTurn| {
        turn.status = TurnStatus::Interrupted;
        turn.completed_at = payload.completed_at;
        turn.duration_ms = payload.duration_ms;
        ThreadHistoryTurnChange::from_pending_turn(turn)
    };
    if let Some(turn_id) = payload.turn_id.as_deref() {
        // Prefer an exact ID match so we interrupt the turn explicitly targeted by the event.
        if let Some(turn) = self.current_turn.as_mut().filter(|turn| turn.id == turn_id) {
            let changed_turn = apply_abort(turn);
            self.record_changed_turn(changed_turn);
            return;
        }

        if let Some(turn) = self.turns.iter_mut().find(|turn| turn.id == turn_id) {
            turn.status = TurnStatus::Interrupted;
            turn.completed_at = payload.completed_at;
            turn.duration_ms = payload.duration_ms;
            let changed_turn = ThreadHistoryTurnChange::from_turn(turn);
            self.record_changed_turn(changed_turn);
            return;
        }
    }

    // If the event has no ID (or refers to an unknown turn), fall back to the active turn.
    if let Some(turn) = self.current_turn.as_mut() {
        let changed_turn = apply_abort(turn);
        self.record_changed_turn(changed_turn);
    }
}

不是所有 ErrorEvent 都影响 Turn 状态,必须先过 affects_turn_status()。中断事件则优先命中当前轮次的 ID,再查已经存下的轮次;只有 ID 缺失或未知才回退到 active turn。由此可以定位一种问题:若旧 Turn 的中断错误地改变了新 Turn,先查 ID 匹配与事件来源,而不是只看最后一个 Interrupted 赋值。

5.3 回滚与收口 ​

源码文件:codex-rs/app-server-protocol/src/protocol/thread_history.rs

相关函数/类型:handle_thread_rollback / finish_current_turn。

rust
fn handle_thread_rollback(&mut self, payload: &ThreadRolledBackEvent) {
    self.finish_current_turn();

    let n = usize::try_from(payload.num_turns).unwrap_or(usize::MAX);
    let removed_turn_ids = if n >= self.turns.len() {
        self.turns.iter().map(|turn| turn.id.clone()).collect()
    } else if n == 0 {
        Vec::new()
    } else {
        self.turns[self.turns.len() - n..]
            .iter()
            .map(|turn| turn.id.clone())
            .collect()
    };
    self.record_removed_turn_ids(removed_turn_ids);

    if n >= self.turns.len() {
        self.turns.clear();
    } else {
        self.turns.truncate(self.turns.len().saturating_sub(n));
    }

    let item_count: usize = self.turns.iter().map(|t| t.items.len()).sum();
    self.next_item_index = i64::try_from(item_count.saturating_add(1)).unwrap_or(i64::MAX);
}

fn finish_current_turn(&mut self) {
    if let Some(turn) = self.current_turn.take() {
        if turn.items.is_empty() && !turn.opened_explicitly && !turn.saw_compaction {
            return;
        }
        self.turns.push(Turn::from(turn));
    }
}

回滚先将当前轮次收口,再按“轮次数”截断;它不是删最后 N 条消息。超大 N 清空全部已归约轮次,saturating_sub 防止下溢,随后重设合成 item 索引。收口时,没有条目、没有显式打开、也没有压缩标记的隐式空 Turn 会丢弃;显式打开的空 Turn 可以保留,用于表达已启动但没有正常输出的执行。

这里修改的是新建 builder 的视图,不会删除 rollout 文件中的旧行。归约结果可以不包含被回滚的消息,磁盘中仍保留记录;排查“文件里有、API 里没有”时,必须把回滚事件也纳入阅读范围。

6. 分页兼容 ​

Paginated 模式的历史投影已经按 Turn 和 item 保存。继续假设“完整历史来自一个 rollout 文件”会读丢前缀,所以处理器先查模式,再选择不同消费者。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:load_persisted_thread_for_read。

rust
async fn load_persisted_thread_for_read(
    &self,
    thread_id: ThreadId,
    include_turns: bool,
) -> Result<Option<Thread>, ThreadReadViewError> {
    let fallback_provider = self.config.model_provider_id.as_str();
    if include_turns {
        let Some(stored_thread) = self
            .read_stored_thread_for_read(thread_id, /*include_history*/ false)
            .await?
        else {
            return Ok(None);
        };
        if matches!(stored_thread.history_mode, ThreadHistoryMode::Paginated) {
            let (mut thread, _) =
                thread_from_stored_thread(stored_thread, fallback_provider, &self.config.cwd);
            thread.turns = self
                .paginated_thread_full_turns(thread_id)
                .await
                .map_err(ThreadReadViewError::JsonRpc)?;
            return Ok(Some(thread));
        }
    }
    let Some(stored_thread) = self
        .read_stored_thread_for_read(thread_id, /*include_history*/ include_turns)
        .await?
    else {
        return Ok(None);
    };
    let (mut thread, history) =
        thread_from_stored_thread(stored_thread, fallback_provider, &self.config.cwd);
    if include_turns && let Some(history) = history {
        thread.turns = build_legacy_api_turns_from_rollout_items(&history.items);
    }
    Ok(Some(thread))
}

includeTurns=true 时的第一次读取只取元数据,用它识别 history_mode。Paginated 转入全量分页;Legacy 才进行携带 history 的第二次读取。这个额外元数据探测是模式分流的一部分,不能删掉后直接请求 include_history=true。

完整兼容读取的成本具有两层循环:外层取轮次,内层给每个轮次补 items。图中的串行顺序来自逐个 await,不暗示并发预取。

前端只发送一个 thread/read,不代表后端只执行一次查询。随着轮次和工具输出增多,这条兼容路径会积累全部结果;需要渐进加载的消费者应使用 thread/turns/list 与 thread/items/list。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:paginated_thread_full_turns。

rust
async fn paginated_thread_full_turns(
    &self,
    thread_id: ThreadId,
) -> Result<Vec<Turn>, JSONRPCErrorError> {
    let mut cursor = None;
    let mut turns = Vec::new();
    loop {
        let page = self
            .paginated_thread_turns_list_response(
                thread_id,
                cursor.clone(),
                Some(THREAD_TURNS_MAX_LIMIT as u32),
                Some(SortDirection::Asc),
                Some(TurnItemsView::Full),
            )
            .await?;
        turns.extend(page.data);
        let Some(next_cursor) = page.next_cursor else {
            return Ok(turns);
        };
        if cursor.as_ref() == Some(&next_cursor) {
            return Err(internal_error(format!(
                "failed to load full thread turns for {thread_id}: thread store returned a repeated cursor"
            )));
        }
        cursor = Some(next_cursor);
    }
}

外层明确选择 Asc 和 Full,最终顺序按旧到新。最大 page size 限制每次请求,不限制返回历史的总量。检查的是本次 next_cursor 与当前 cursor 是否相同,用来阻止连续不推进;它不是维护已访问集合的通用循环检测,不能声称能识别任意 A→B→A 循环。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:paginated_thread_turns_list_response。

rust
async fn paginated_thread_turns_list_response(
    &self,
    thread_id: ThreadId,
    cursor: Option<String>,
    limit: Option<u32>,
    sort_direction: Option<SortDirection>,
    items_view: Option<TurnItemsView>,
) -> Result<ThreadTurnsListResponse, JSONRPCErrorError> {
    let items_view = items_view.unwrap_or(TurnItemsView::Summary);
    let page_size = thread_turns_page_size(limit);
    let sort_direction = match sort_direction.unwrap_or(SortDirection::Desc) {
        SortDirection::Asc => StoreSortDirection::Asc,
        SortDirection::Desc => StoreSortDirection::Desc,
    };
    // `Full` is only a temporary compatibility path. Keep it out of ThreadStore's API:
    // load turn shells here, then hydrate their items below.
    let stored_items_view = match items_view {
        TurnItemsView::NotLoaded => StoredTurnItemsView::NotLoaded,
        TurnItemsView::Summary => StoredTurnItemsView::Summary,
        TurnItemsView::Full => StoredTurnItemsView::NotLoaded,
    };
    let page = self
        .thread_store
        .list_turns(StoreListTurnsParams {
            thread_id,
            include_archived: true,
            cursor,
            page_size,
            sort_direction,
            items_view: stored_items_view,
        })
        .await
        .map_err(|err| match err {
            ThreadStoreError::InvalidRequest { message } => invalid_request(message),
            ThreadStoreError::Unsupported { operation } => {
                unsupported_thread_store_operation(operation)
            }
            ThreadStoreError::ThreadNotFound { thread_id } => {
                invalid_request(format!("no rollout found for thread id {thread_id}"))
            }
            err => internal_error(format!("failed to list thread history: {err}")),
        })?;
    let mut turns = Vec::with_capacity(page.turns.len());
    for turn in page.turns {
        let mut turn = stored_turn_to_api_turn(turn, items_view)?;
        if matches!(items_view, TurnItemsView::Full) {
            turn.items = self
                .paginated_turn_full_items(thread_id, turn.id.as_str())
                .await?;
        }
        turns.push(turn);
    }
    let loaded_thread = self.thread_manager.get_thread(thread_id).await.ok();
    let has_live_running_thread = match loaded_thread.as_ref() {
        Some(thread) => matches!(thread.agent_status().await, AgentStatus::Running),
        None => false,
    };
    normalize_thread_turns_status(
        &mut turns,
        self.thread_watch_manager
            .loaded_status_for_thread(&thread_id.to_string())
            .await,
        has_live_running_thread,
    );
    Ok(ThreadTurnsListResponse {
        data: turns,
        next_cursor: page.next_cursor,
        backwards_cursor: page.backwards_cursor,
    })
}

Full 在 Store 层被映射成 NotLoaded:先得到轮次外壳,再由 App Server 显式补齐。这是兼容职责留在协议层的体现;Store 的 StoredTurnItemsView 无需因此增加一条“读完所有 items”的重操作契约。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:paginated_turn_full_items。

rust
// app-server-only hydration path until those clients use `thread/items/list`.
async fn paginated_turn_full_items(
    &self,
    thread_id: ThreadId,
    turn_id: &str,
) -> Result<Vec<ThreadItem>, JSONRPCErrorError> {
    let mut cursor = None;
    let mut items = Vec::new();
    loop {
        let page = self
            .thread_store
            .list_items(StoreListItemsParams {
                thread_id,
                turn_id: Some(turn_id.to_string()),
                include_archived: true,
                cursor: cursor.clone(),
                page_size: THREAD_ITEMS_MAX_LIMIT,
                sort_direction: StoreSortDirection::Asc,
                sort_key: StoreItemSortKey::CreatedAtOrdinal,
                after_updated_at_ordinal: None,
            })
            .await
            .map_err(paginated_history_list_error)?;
        for item in page.items {
            items.push(deserialize_stored_thread_item(item)?);
        }
        let Some(next_cursor) = page.next_cursor else {
            return Ok(items);
        };
        if cursor.as_ref() == Some(&next_cursor) {
            return Err(internal_error(format!(
                "failed to load full turn items for {turn_id}: thread store returned a repeated cursor"
            )));
        }
        cursor = Some(next_cursor);
    }
}

内层按 CreatedAtOrdinal 升序取某个 Turn 的条目,每条调用反序列化函数。某页查询或某个条目转换失败,? 会让整个 RPC 失败,不会返回静默截断的部分成功数组。重复 cursor 也产生内部错误。这是“慢但完整”的兼容行为:不能为了容错在错误处直接 break,否则客户端无法区分完整历史与被截断历史。

7. 状态叠加 ​

历史中最后一条记录是 TurnStarted,只能说明持久化前缀中没有终态,不能说明当前进程里还在运行。响应必须把这个历史事实和实时状态结合。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:read_thread_view 的状态收尾。

rust
    let has_live_in_progress_turn = if let Some(loaded_thread) = loaded_thread.as_ref() {
        matches!(loaded_thread.agent_status().await, AgentStatus::Running)
    } else {
        false
    };

    let thread_status = self
        .thread_watch_manager
        .loaded_status_for_thread(&thread.id)
        .await;

    set_thread_status_and_interrupt_stale_turns(
        &mut thread,
        thread_status,
        has_live_in_progress_turn,
    );
    Ok(thread)
}

这里分别读取 Core 的 AgentStatus::Running 和 watch manager 的 ThreadStatus。它们可能由不同异步路径更新,因此下一步需要显式解决短暂不同步,而不是任选一方覆盖另一方。

源码文件:codex-rs/app-server/src/thread_status.rs

相关函数/类型:resolve_thread_status。

rust
pub(crate) fn resolve_thread_status(
    status: ThreadStatus,
    has_in_progress_turn: bool,
) -> ThreadStatus {
    // Running-turn events can arrive before the watch runtime state is observed by
    // the listener loop. In that window we prefer to reflect a real active turn as
    // `Active` instead of `Idle`/`NotLoaded`.
    if has_in_progress_turn && matches!(status, ThreadStatus::Idle | ThreadStatus::NotLoaded) {
        return ThreadStatus::Active {
            active_flags: Vec::new(),
        };
    }

    status
}

只有 watch 状态为 Idle 或 NotLoaded、且 Core 表示确有进行中的 Turn 时,才提升为 Active。SystemError 不会被这条规则抹掉,已有 Active 中的审批或输入等待标记也会保留。这是一个有条件的优先级规则。

下面画的是本次响应的修正规则,不是 Core 的生命周期状态机。

如果 Thread 已经不活跃,残留的 InProgress Turn 要显示为中断;已完成和失败的 Turn 不受影响。

源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs

相关函数/类型:set_thread_status_and_interrupt_stale_turns。

rust
pub(super) fn set_thread_status_and_interrupt_stale_turns(
    thread: &mut Thread,
    loaded_status: ThreadStatus,
    has_live_in_progress_turn: bool,
) {
    let status = resolve_thread_status(loaded_status, has_live_in_progress_turn);
    if !matches!(status, ThreadStatus::Active { .. }) {
        for turn in &mut thread.turns {
            if matches!(turn.status, TurnStatus::InProgress) {
                turn.status = TurnStatus::Interrupted;
            }
        }
    }
    thread.status = status;
}

函数只修改传入的协议 Thread 值,没有提交 Op::Interrupt,也没有追加中断记录。UI 看到 Interrupted 可以来自读取时归一化,不能反推运行时曾收到一个明确的用户取消命令。Paginated 路径会在每页构造时做类似归一化,最后外层再完成 Thread 状态赋值。

8. 错误分层 ​

为了判断能否重试,要保留“未物化”和真正损坏之间的边界。已加载不意味着已有可读历史:本地新 Thread 的预计算路径可先返回,第一次用户消息之前文件仍可能不存在。这个错误描述的是当前后端返回的缺失形态,不能把它推广为所有 Store 必须等用户消息后才可读。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:thread_read_history_load_error。

rust
fn thread_read_history_load_error(
    thread_id: ThreadId,
    err: ThreadStoreError,
) -> ThreadReadViewError {
    match err {
        ThreadStoreError::InvalidRequest { message }
            if message.starts_with("failed to resolve rollout path `") =>
        {
            ThreadReadViewError::InvalidRequest(format!(
                "thread {thread_id} is not materialized yet; includeTurns is unavailable before first user message"
            ))
        }
        ThreadStoreError::ThreadNotFound {
            thread_id: missing_thread_id,
        } if missing_thread_id == thread_id => ThreadReadViewError::InvalidRequest(format!(
            "thread {thread_id} is not materialized yet; includeTurns is unavailable before first user message"
        )),
        ThreadStoreError::InvalidRequest { message } => {
            ThreadReadViewError::InvalidRequest(message)
        }
        ThreadStoreError::Unsupported { operation } => ThreadReadViewError::Unsupported(operation),
        err => ThreadReadViewError::Internal(format!(
            "failed to load thread history for thread {thread_id}: {err}"
        )),
    }
}

路径解析失败的特定前缀,以及目标 Thread 自身的 ThreadNotFound,转成“尚未物化”的请求错误。其他非法请求保留原因,Unsupported 单独处理,其余错误视为内部加载失败。这个映射依赖错误形态,不能把所有 I/O 错误都建议成“先发一条消息”。

源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs

相关函数/类型:ThreadReadViewError / thread_read_view_error。

rust
enum ThreadReadViewError {
    InvalidRequest(String),
    Unsupported(&'static str),
    Internal(String),
    JsonRpc(JSONRPCErrorError),
}

fn thread_read_view_error(err: ThreadReadViewError) -> JSONRPCErrorError {
    match err {
        ThreadReadViewError::InvalidRequest(message) => invalid_request(message),
        ThreadReadViewError::Unsupported(operation) => {
            unsupported_thread_store_operation(operation)
        }
        ThreadReadViewError::Internal(message) => internal_error(message),
        ThreadReadViewError::JsonRpc(error) => error,
    }
}

最后的 RPC 映射区分请求错误、存储不支持、内部故障和已形成的 JSON-RPC 错误。错误码的共用定义见AppServer错误码体系;对读取来说,更实用的是根据触发阶段判断资源状态。

现象首先定位可以得出的结论
invalid thread id参数解析尚未进入 Store 读取
thread not loaded来源选择与缺失映射本次请求没有可用的持久化或 live 来源
ephemeral threads do not support includeTurnslive 配置快照可以读摘要,不能走该完整历史契约
includeTurns is unavailable before first user messagelive history 加载运行句柄存在,历史尚未物化
repeated cursor全量 Turn/item 循环后端分页没有连续推进,不能返回完整结果
Unsupported具体 Store 操作切换一个 RPC 参数不一定能补齐后端能力

这些失败不销毁原来的运行 Thread;返回过程中只持有临时快照和组装中的数组。若 Paginated 的持久化已经成功、后续分页失败,先前的持久化不会因为响应失败而回滚。读取过程也没有在这些函数里设置独立的业务超时,测试中的 timeout 是客户端等待上限,不能拿它当服务端的取消协议。

平台差异主要位于路径规范化和 Store 实现,本文主调用链没有按 Windows/Linux 切成两套。真正需要先识别的是 Local/InMemory 后端、Legacy/Paginated 模式及 ephemeral 状态;三组条件不能混为一个“本地线程”标签。

9. 读取取证 ​

9.1 路径与物化 ​

上游测试从真实的 thread/start 开始,断言路径尚不存在,再发送摘要读取。下面保留触发输入和关键断言。

源码文件:codex-rs/app-server/tests/suite/v2/thread_read.rs

相关函数/类型:thread_read_loaded_thread_returns_precomputed_path_before_materialization。

rust
let thread_path = thread.path.clone().expect("thread path");
assert!(
    !thread_path.exists(),
    "fresh thread rollout should not be materialized yet"
);

let read_id = mcp
    .send_thread_read_request(ThreadReadParams {
        thread_id: thread.id.clone(),
        include_turns: false,
    })
    .await?;
let ThreadReadResponse { thread: read, .. } =
    timeout(DEFAULT_READ_TIMEOUT, mcp.read_response(read_id)).await??;

assert_eq!(read.id, thread.id);
assert_eq!(read.path, Some(thread_path));
assert!(read.preview.is_empty());
assert_eq!(read.turns.len(), 0);
assert_eq!(read.status, ThreadStatus::Idle);

Ok(())

它证明“有路径、无文件、摘要仍成功”是正常组合,状态可以是 Idle,并不要求回到 NotLoaded。该断言不能证明完整历史也能读取。

同一文件中的 thread_read_include_turns_rejects_unmaterialized_loaded_thread 把开关改为 true,读的是错误响应。

源码文件:codex-rs/app-server/tests/suite/v2/thread_read.rs

相关函数/类型:thread_read_include_turns_rejects_unmaterialized_loaded_thread。

rust
let read_id = mcp
    .send_thread_read_request(ThreadReadParams {
        thread_id: thread.id.clone(),
        include_turns: true,
    })
    .await?;
let read_err: JSONRPCError = timeout(
    DEFAULT_READ_TIMEOUT,
    mcp.read_stream_until_error_message(RequestId::Integer(read_id)),
)
.await??;

assert!(
    read_err
        .error
        .message
        .contains("includeTurns is unavailable before first user message"),
    "unexpected error: {}",
    read_err.error.message
);

两项测试应配对看:失败来自历史要求升级,不是 ID 无效或 Thread 创建失败。做客户端时,可以先显示摘要,再按实际持久化时机刷新历史,不能因为暂时读取失败而创建替代 Thread。

9.2 无路径历史 ​

thread_read_loaded_include_turns_reads_store_history_without_rollout_path 通过 experimental_thread_store 选择共享 ID 的 InMemory Store,启动 in-process server 后直接向 Store 追加 fixture。

源码文件:codex-rs/app-server/tests/suite/v2/thread_read.rs

相关函数/类型:thread_read_loaded_include_turns_reads_store_history_without_rollout_path。

rust
let ThreadStartResponse { thread, .. } = serde_json::from_value(result)?;
assert_eq!(thread.path, None);

let thread_id = codex_protocol::ThreadId::from_string(&thread.id)?;
store
    .append_items(AppendThreadItemsParams {
        thread_id,
        items: store_history_items(),
    })
    .await?;

let result = client
    .request(ClientRequest::ThreadRead {
        request_id: RequestId::Integer(2),
        params: ThreadReadParams {
            thread_id: thread.id,
            include_turns: true,
        },
    })
    .await?
    .expect("thread/read should succeed");
let ThreadReadResponse { thread, .. } = serde_json::from_value(result)?;

assert_eq!(turn_user_texts(&thread.turns), vec!["history from store"]);
let [ThreadItem::UserMessage { content, .. }] = thread.turns[0].items.as_slice() else {
    panic!("expected one user message item");
};

这个测试在 path=None 的前提下仍返回 “history from store”。它反向证明 App Server 使用 Store 抽象,而不是在 Thread.path 为空时一律报错;但 InMemory 的成功不证明本地 rollout 的写盘和损坏恢复行为。

9.3 投影与全量 ​

paginated_history_lists_and_legacy_reads_use_projected_turns_and_items 在 Local Store 中显式创建 Paginated Thread,写入两轮记录,并更新同一条用户条目的 client_id。第一个轮次有完成事件,第二个只开始而未完成;测试先读未加载的 Thread,再走旧客户端的完整 resume,最后请求一条记录的初始页。

源码文件:codex-rs/app-server/tests/suite/v2/thread_read.rs

相关函数/类型:paginated_history_lists_and_legacy_reads_use_projected_turns_and_items。

rust
let read_id = mcp
    .send_thread_read_request(ThreadReadParams {
        thread_id: thread_id.to_string(),
        include_turns: true,
    })
    .await?;
let ThreadReadResponse {
    thread: unloaded_thread,
} = timeout(DEFAULT_READ_TIMEOUT, mcp.read_response(read_id)).await??;
assert_eq!(unloaded_thread.turns, expected_full_turns);

let legacy_resume_id = mcp
    .send_thread_resume_request(ThreadResumeParams {
        thread_id: thread_id.to_string(),
        ..Default::default()
    })
    .await?;
let ThreadResumeResponse {
    thread: legacy_thread,
    ..
} = timeout(DEFAULT_READ_TIMEOUT, mcp.read_response(legacy_resume_id)).await??;
assert_eq!(legacy_thread.turns, expected_full_turns);

两次都用整个 expected_full_turns 比较,能同时检查顺序、条目更新和残尾终态,避免只断言数组长度而漏掉错误内容。该 fixture 的两个轮次很小,不会触发全量补取函数的最大页大小;它证明投影路由和完整条目语义,不能单独证明大历史跨页时 cursor 始终前进。跨页行为要结合两层循环及其 repeated-cursor 错误分支阅读。

9.4 轮次与残尾 ​

thread_read_can_include_turns 构造带文本元素的 Legacy rollout,检查消息、显示类型与未加载状态。

源码文件:codex-rs/app-server/tests/suite/v2/thread_read.rs

相关函数/类型:thread_read_can_include_turns。

rust
assert_eq!(thread.turns.len(), 1);
let turn = &thread.turns[0];
assert_eq!(turn.status, TurnStatus::Completed);
assert_eq!(turn.items_view, TurnItemsView::Full);
assert_eq!(turn.items.len(), 1, "expected user message item");
match &turn.items[0] {
    ThreadItem::UserMessage { content, .. } => {
        assert_eq!(
            content,
            &vec![UserInput::Text {
                text: preview.to_string(),
                text_elements: text_elements.clone().into_iter().map(Into::into).collect(),
            }]
        );
    }
    other => panic!("expected user message item, got {other:?}"),
}
assert_eq!(thread.status, ThreadStatus::NotLoaded);

它证明用户文本及其字节范围标记能进入 UserInput::Text,返回 items_view=Full,且读出历史不等于加载运行对象。缺少事件边界的 fixture 被归约成 Completed,也不能据此认定每个 rollout 都天然有完整生命周期。

另一个测试人为追加 TurnStarted 与消息,故意不追加完成事件,再分别 resume、再次 resume 和 read。

源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs

相关函数/类型:thread_resume_and_read_interrupt_incomplete_rollout_turn_when_thread_is_idle。

rust
let ThreadReadResponse {
    thread: read_thread,
    ..
} = timeout(DEFAULT_READ_TIMEOUT, mcp.read_response(read_id)).await??;

assert_eq!(read_thread.status, ThreadStatus::Idle);
assert_eq!(read_thread.turns.len(), 2);
assert_eq!(read_thread.turns[1].id, turn_id);
assert_eq!(read_thread.turns[1].status, TurnStatus::Interrupted);

Ok(())

断言旧轮次仍为 Completed,残尾为 Interrupted,而 Thread 是 Idle。它验证响应中的残尾修正与重复读取的稳定性;不会证明原文件已补写一个中断事件。

在 Codex 仓库根目录执行下列上游测试,可同时覆盖读取与相邻恢复路径。just test 使用仓库配置的 nextest;若只想定位某个分支,可把过滤表达式缩小到完整测试名。

bash
just test --locked -p codex-app-server --test all \
  -E 'test(suite::v2::thread_read::)'
just test --locked -p codex-app-server --test all \
  -E 'test(thread_resume_and_read_interrupt_incomplete_rollout_turn_when_thread_is_idle)'
rg -n 'read_thread_view|apply_thread_read_store_fields|paginated_thread_full_turns' \
  codex-rs/app-server/src/request_processors/thread_processor.rs

最后用一个具体故障检验这条主线:客户端读到 path=null、turns=[],同时 status=Active。先确认请求有没有 includeTurns,再确认 Store 类型,最后核对 watch 与 Core Running;这三处证据齐全前,不能断言历史丢失。反过来,若开启完整历史后响应失败,沿“来源裁决 → 模式选择 → Store 操作 → 错误映射”逐段检查,就能把缺失、未物化、后端不支持与分页故障分开。