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。
#[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。
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。
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 的响应出口。
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 的来源选择。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
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。
// 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 的状态收尾。
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。
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。
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。
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。
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 includeTurns | live 配置快照 | 可以读摘要,不能走该完整历史契约 |
includeTurns is unavailable before first user message | live 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。
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。
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。
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。
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。
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。
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;若只想定位某个分支,可把过滤表达式缩小到完整测试名。
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 操作 → 错误映射”逐段检查,就能把缺失、未物化、后端不支持与分页故障分开。
