ThreadResume处理流程
重新打开一个正在生成答案的对话,客户端需要同时得到已有内容和后续事件;重新打开一个已经退出进程的对话,则需要从持久记录重建运行对象。它们都使用 thread/resume,却不能采用同一套“读文件、创建会话、发送历史”的流程。前者如果再创建一个 Core Thread,就会出现两个执行者;如果先读历史、稍后才订阅,又可能漏掉两者之间产生的事件。
本文面向掌握 Rust Arc、Mutex、channel 与 async/await 的读者。先读ThreadStart处理流程了解 Core 对象与监听任务如何建立,再读ThreadRead历史加载区分持久历史和返回给界面的历史视图。本文回答的是:一次 resume 怎样选定执行者,把配置、历史、订阅和仍待处理的请求交给它;底层逐条重建算法可继续读Rollout重建与恢复。
这里的“冷恢复”指需要创建新的进程内运行对象,“热重连”指复用已加载的 CodexThread。已加载 Thread 可以处于 Idle,不一定正在执行 Turn;源码中 resume_running_thread 的 running 首先指对象还在 manager 中。Thread ID 是持久对话身份,Session 是本次承载状态的 Core 对象,Turn 是一次任务,listener 是消费 Core 事件的 App Server 任务。冷恢复创建对象也不意味着所有协议 ID 都重新生成。下文摘录来自实际实现,跨段省略用 // ... 表示。
1. 恢复入口
1.1 请求分派
协议宏同时登记方法名、实验字段检查和请求串行化范围。resume 的公开名称带斜杠,不能从内部 Rust 函数名推成 thread.resume。
源码文件:codex-rs/app-server-protocol/src/protocol/common.rs
相关函数/类型:ClientRequest::ThreadResume 注册。
ThreadResume => "thread/resume" {
params: v2::ThreadResumeParams,
inspect_params: true,
serialization: thread_or_path(params.thread_id, params.path),
response: v2::ThreadResumeResponse,
},inspect_params 允许按字段决定实验能力要求;它不表示整个方法都只能实验客户端使用。thread_or_path 将请求纳入线程或路径维度的串行队列,但随后仍需 Core 注册表与 writer 所有权检查:同一资源可能有 ID 和路径两种表示,请求排队不能代替持久层排他性。
源码文件:codex-rs/app-server-protocol/src/protocol/common.rs
相关函数/类型:serialization_scope_expr 的 thread_or_path 分支。
($actual_params:ident, thread_or_path($params:ident . $thread_field:ident, $params2:ident . $path_field:ident)) => {
if !$actual_params.$thread_field.is_empty() {
Some(ClientRequestSerializationScope::Thread {
thread_id: $actual_params.$thread_field.clone(),
})
} else if let Some(path) = $actual_params.$path_field.clone() {
Some(ClientRequestSerializationScope::ThreadPath { path })
} else {
Some(ClientRequestSerializationScope::Thread {
thread_id: $actual_params.$thread_field.clone(),
})
}
};只要 thread_id 非空,排队键就优先用该字符串;只有空 ID 才改用 path。因此不能把冷恢复的“history > path > ID”来源优先级照搬成排队规则。路径解析出真实 ID 后,仍需后续所有者检查。
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数/类型:MessageProcessor::process_request。
ClientRequest::ThreadResume { params, .. } => {
self.thread_processor
.thread_resume(
request_id.clone(),
params,
app_server_client_name.clone(),
client_version.clone(),
client_mcp_extensions.clone(),
)
.await
}请求 ID、客户端名称与版本、MCP 扩展能力一起传入业务处理器。客户端名称稍后参与返回载荷裁剪;MCP 扩展则参与冷恢复装配,二者都不是模型提示词中的自由文本。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:ThreadRequestProcessor::thread_resume。
pub(crate) async fn thread_resume(
&self,
request_id: ConnectionRequestId,
params: ThreadResumeParams,
app_server_client_name: Option<String>,
app_server_client_version: Option<String>,
client_mcp_extensions: ClientMcpExtensions,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
self.thread_resume_inner(
request_id,
params,
app_server_client_name,
app_server_client_version,
client_mcp_extensions,
)
.await
.map(|()| None)
}这里成功返回的是 None。冷恢复自己发送响应,热重连把发送工作交给 listener;通用分派层不应再补发一个空响应。None 表示业务分支接管响应,而不是客户端应收到 JSON null。
1.2 输入契约
history 与 path 都是实验字段。先读其反序列化规则,才能区分“不传路径”和“传了一个空路径”。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:ThreadResumeParams 的身份字段。
pub struct ThreadResumeParams {
pub thread_id: String,
/// [UNSTABLE] FOR CODEX CLOUD - DO NOT USE.
/// If specified, the thread will be resumed with the provided history
/// instead of loaded from disk.
#[experimental("thread/resume.history")]
#[ts(optional = nullable)]
pub history: Option<Vec<ResponseItem>>,
/// [UNSTABLE] Specify the rollout path to resume from.
/// If specified for a non-running thread, the thread_id param will be ignored.
/// If thread_id identifies a running thread, the path must match the active
/// rollout path.
#[experimental("thread/resume.path")]
#[serde(
default,
deserialize_with = "crate::protocol::serde_helpers::deserialize_empty_path_as_none"
)]
#[ts(optional = nullable)]
pub path: Option<PathBuf>,
// ...
}deserialize_empty_path_as_none 把空字符串路径视为缺省。非空 path 对未加载线程可以决定恢复来源;对已加载 ID 却变成一致性检查。history 是模型的 ResponseItem 数组,不是界面返回的 Turn[],不能把 read 的响应直接回填进这个字段。
以下两个选项只控制响应中的历史呈现。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:ThreadResumeParams 的历史选项。
/// When true, return only thread metadata and live-resume state without
/// populating `thread.turns`. This is useful when the client plans to call
/// `thread/turns/list` immediately after resuming.
#[experimental("thread/resume.excludeTurns")]
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub exclude_turns: bool,
/// When present, include a `thread/turns/list` page in the resume response
/// so clients can bootstrap recent turns without a second request.
#[experimental("thread/resume.initialTurnsPage")]
#[ts(optional = nullable)]
pub initial_turns_page: Option<ThreadResumeInitialTurnsPageParams>,exclude_turns 缺省 false,所以传统客户端仍收到完整 thread.turns;initial_turns_page 可以同时请求一个初始页。冷恢复即使排除响应轮次,也必须重建模型上下文。因此“少返回历史”和“不恢复历史”是两种完全不同的操作。
2. 身份与来源
2.1 入口拒绝
已进入卸载阶段的 Thread 不能被重新挂上连接。处理器在读取大量历史前先检查这一状态,并拒绝互斥的权限输入。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的前置检查。
if let Ok(thread_id) = ThreadId::from_string(¶ms.thread_id)
&& self
.pending_thread_unloads
.lock()
.await
.contains(&thread_id)
{
self.outgoing
.send_error(
request_id,
invalid_request(format!(
"thread {thread_id} is closing; retry thread/resume after the thread is closed"
)),
)
.await;
return Ok(());
}
if params.sandbox.is_some() && params.permissions.is_some() {
self.outgoing
.send_error(
request_id,
invalid_request("`permissions` cannot be combined with `sandbox`"),
)
.await;
return Ok(());
}pending_thread_unloads 锁保护正在关闭的 ID 集合,不覆盖整个恢复流程;错误提示要求等关闭结束再重试。sandbox 与 permissions 同时出现直接报错,避免旧枚举与命名权限配置争夺同一个来源。此时尚未启动替代 Thread。
接着取得列表状态 permit,再探测已加载对象。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的运行探测。
let _thread_list_state_permit = match self.acquire_thread_list_state_permit().await {
Ok(permit) => permit,
Err(error) => {
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
};
let stored_thread_from_running_probe = match self
.resume_running_thread(
&request_id,
¶ms,
app_server_client_name.clone(),
app_server_client_version.clone(),
/*cold_resume_history*/ None,
)
.await
{
Ok(RunningThreadResumeResult::Handled) => return Ok(()),
Ok(RunningThreadResumeResult::NotRunning(stored_thread)) => stored_thread,
Err(error) => {
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
};permit 的局部变量一直活到当前函数返回;它与按请求资源排队是两层协调。返回 Handled 表示已处理或已交给 listener;NotRunning(Some(stored_thread)) 则携带刚读到的摘要,冷路径可以复用。不要把 Handled 解释成对端已经收到响应。
2.2 来源优先级
运行对象探测先处理 history,再尝试已加载 ID,最后才根据存储解析出来的真实 ID 查 manager。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:resume_running_thread 的身份探测。
let running_thread = if params.history.is_some() {
if let Ok(existing_thread_id) = ThreadId::from_string(¶ms.thread_id)
&& self
.thread_manager
.get_thread(existing_thread_id)
.await
.is_ok()
{
return Err(invalid_request(format!(
"cannot resume thread {existing_thread_id} with history while it is already running"
)));
}
None
} else if let Ok(existing_thread_id) = ThreadId::from_string(¶ms.thread_id)
&& let Ok(existing_thread) = self.thread_manager.get_thread(existing_thread_id).await
{
let source_thread = self
.read_stored_thread_for_resume(
¶ms.thread_id,
/*path*/ None,
/*include_history*/ false,
)
.await?;
Some((existing_thread_id, existing_thread, source_thread))
} else {
let source_thread = self
.read_stored_thread_for_resume(
¶ms.thread_id,
params.path.as_ref(),
/*include_history*/ false,
)
.await?;
let existing_thread_id = source_thread.thread_id;
match self.thread_manager.get_thread(existing_thread_id).await {
Ok(existing_thread) => Some((existing_thread_id, existing_thread, source_thread)),
Err(_) => {
return Ok(RunningThreadResumeResult::NotRunning(Some(Box::new(
source_thread,
))));
}
}
};若传 history 且 ID 已加载,立即拒绝;若 ID 已加载,即使带 path,也先按这个 ID 读取摘要。只有没有命中已加载 ID 时,存储路径才能解析出另一个真实 Thread ID。源码把“由参数定位来源”与“按来源身份找执行者”分成两步,这正是别名路径与 ID 不应各自创建执行者的原因。
冷路径的来源顺序在另一段中明确写出。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的初始历史选择。
let resume_result = if let Some(history) = history {
self.resume_thread_from_history(history.as_slice())
.await
.map(|thread_history| (thread_history, None))
} else if let Some(stored_thread) = stored_thread_from_running_probe {
self.load_resume_initial_history_from_stored_thread(*stored_thread)
.await
.map(|(thread_history, stored_thread)| (thread_history, Some(stored_thread)))
} else {
match self
.read_stored_thread_for_resume(
&thread_id,
path.as_ref(),
/*include_history*/ false,
)
.await
{
Ok(stored_thread) => self
.load_resume_initial_history_from_stored_thread(stored_thread)
.await
.map(|(thread_history, stored_thread)| (thread_history, Some(stored_thread))),
Err(error) => Err(error),
}
};
let (thread_history, resume_source_thread) = match resume_result {
Ok(value) => value,
Err(error) => {
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
};
let paginated_thread_id = resume_source_thread.as_ref().and_then(|thread| {
matches!(thread.history_mode, ThreadHistoryMode::Paginated).then_some(thread.thread_id)
});
let paginated_resume = paginated_thread_id.is_some();优先 history,其次复用探测得到的摘要,否则重新按 path/ID 读取。重读并非全是多余 I/O:后面会看到,关闭一个空闲缓存对象可能刷新新记录,所以替换之前必须重新取持久来源。
下面用一张图归纳身份规则;它不把“已加载”和“当前 Turn 运行中”合成一个判断。
图中的 Forked 来自明确的实现转换,不是根据方法名猜测的恢复类型。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:resume_thread_from_history。
async fn resume_thread_from_history(
&self,
history: &[ResponseItem],
) -> Result<InitialHistory, JSONRPCErrorError> {
if history.is_empty() {
return Err(invalid_request("history must not be empty"));
}
Ok(InitialHistory::Forked(
history
.iter()
.cloned()
.map(|item| RolloutItem::ResponseItem(item.into()))
.collect(),
))
}空数组无效。非空数组转成 RolloutItem::ResponseItem 后走 InitialHistory::Forked,没有从旧 Thread 继承持久身份的承诺。对应地,Core 为没有预留 ID 的 Forked 生成新 ID,而 Resumed 使用记录里的 conversation ID。
源码文件:codex-rs/core/src/session/session.rs
相关函数/类型:Session::new 的 thread_id 选择。
let thread_id = match (&initial_history, reserved_thread_id) {
(
InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_),
Some(thread_id),
) => thread_id,
(InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_), None) => {
agent_control.generate_thread_id()
}
(InitialHistory::Resumed(resumed_history), None) => resumed_history.conversation_id,
(InitialHistory::Resumed(_), Some(_)) => {
return Err(anyhow::anyhow!(
"reserved thread ID cannot be used when resuming a thread"
));
}
};所以“API 名为 resume”并不足以断言返回 ID 必等于传入 ID。客户端应采用响应中的 thread.id。同一段还拒绝 Resumed 搭配 reserved ID,防止把旧记录恢复到一个任意新身份。
2.3 路径与归档
热重连比较请求路径与实际活动路径;不同就报 stale path,不会悄悄改接另一份文件。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:resume_running_thread 的路径一致性。
let paginated_resume =
matches!(source_thread.history_mode, ThreadHistoryMode::Paginated);
let existing_thread_rollout_path = existing_thread.rollout_path();
let active_path = existing_thread_rollout_path
.as_ref()
.or(source_thread.rollout_path.as_ref());
if let (Some(requested_path), Some(active_path)) = (params.path.as_ref(), active_path)
&& !path_utils::paths_match_after_normalization(requested_path, active_path)
{
return Err(invalid_request(format!(
"cannot resume running thread {existing_thread_id} with stale path: requested `{}`, active `{}`",
requested_path.display(),
active_path.display()
)));
}活动路径优先来自 Core,缺失时才用 Store 摘要。路径比较经过归一化,不能用请求字符串的字节相等替代这个规则。
冷路径的 Paginated 记录也要重新按 ID 找当前路径,因为旧窗口的 rollout 路径可能已过期。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:read_stored_thread_for_resume 的路径与归档检查。
let stored_thread = result.map_err(thread_store_resume_read_error)?;
if let Some(requested_path) = path
&& matches!(stored_thread.history_mode, ThreadHistoryMode::Paginated)
{
let current_thread = self
.thread_store
.read_thread(StoreReadThreadParams {
thread_id: stored_thread.thread_id,
include_archived: true,
include_history: false,
})
.await
.map_err(thread_store_resume_read_error)?;
if let Some(current_path) = current_thread.rollout_path.as_ref()
&& !path_utils::paths_match_after_normalization(
codex_rollout::plain_rollout_path(requested_path).as_path(),
codex_rollout::plain_rollout_path(current_path).as_path(),
)
{
return Err(invalid_request(format!(
"cannot resume paginated thread {} with stale path: requested {}, current {}; omit path and resume by thread id",
stored_thread.thread_id,
requested_path.display(),
current_path.display()
)));
}
}
if stored_thread.archived_at.is_some() {
let thread_id = stored_thread.thread_id;
return Err(invalid_request(format!(
"session {thread_id} is archived. Run `codex unarchive {thread_id}` to unarchive it first."
)));
}
Ok(stored_thread)这里先允许 Store 读到归档元数据,再显式拒绝恢复归档 Thread。与 read 能检查归档历史不同,resume 不会自动 unarchive。Paginated 的旧路径错误会指导调用者省略 path、按 ID 恢复;不能把路径当成永不变化的主键。
3. 配置生效
3.1 空闲缓存替换
热重连通常保留正在使用的配置。但无订阅者的 Idle Thread 可以只是缓存,用户带来的新配置不能永远被缓存挡住。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:resume_running_thread 的覆盖与替换。
let config_snapshot = existing_thread.config_snapshot().await;
let mismatch_details = collect_resume_override_mismatches(params, &config_snapshot);
if !mismatch_details.is_empty() {
let has_subscribers = !self
.thread_state_manager
.subscribed_connection_ids(existing_thread_id)
.await
.is_empty();
let loaded_status = self
.thread_watch_manager
.loaded_status_for_thread(&existing_thread_id.to_string())
.await;
let is_running =
matches!(existing_thread.agent_status().await, AgentStatus::Running);
// Parent-owned V2 children must not be rebuilt from public resume overrides.
if can_accept_direct_input(
existing_thread.multi_agent_version(),
&config_snapshot.session_source,
) && !has_subscribers
&& matches!(loaded_status, ThreadStatus::Idle)
&& !is_running
{
// A loaded idle thread is only a cache entry. Shut it down
// before removing it so cold resume cannot duplicate a
// thread that timed out during shutdown.
match wait_for_thread_shutdown(&existing_thread).await {
ThreadShutdownResult::Complete => {
self.thread_manager.remove_thread(&existing_thread_id).await;
self.finalize_thread_teardown(existing_thread_id).await;
// Shutdown can flush newer rollout items, so reload the
// stored thread before starting the replacement session.
return Ok(RunningThreadResumeResult::NotRunning(None));
}
ThreadShutdownResult::SubmitFailed => {
warn!("failed to submit Shutdown to thread {existing_thread_id}");
}
ThreadShutdownResult::TimedOut => {
warn!("thread {existing_thread_id} shutdown timed out");
}
}
}
// Preserve rejoin semantics when another client can still observe
// the loaded thread or shutdown did not complete.
tracing::warn!(
"thread/resume overrides ignored for loaded thread {}: {}",
existing_thread_id,
mismatch_details.join("; ")
);
}替换要同时满足四个条件:允许直接输入、没有订阅连接、watch 状态 Idle、Core 不在 Running。只有 shutdown_and_wait 真正完成,才从 manager 删除并清理 listener/watch,然后返回 NotRunning(None),迫使外层重读存储。若提交关闭失败或超时,旧对象保留,继续 rejoin,并记录覆盖被忽略的原因。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:wait_for_thread_shutdown。
pub(super) async fn wait_for_thread_shutdown(thread: &Arc<CodexThread>) -> ThreadShutdownResult {
match tokio::time::timeout(Duration::from_secs(10), thread.shutdown_and_wait()).await {
Ok(Ok(())) => ThreadShutdownResult::Complete,
Ok(Err(_)) => ThreadShutdownResult::SubmitFailed,
Err(_) => ThreadShutdownResult::TimedOut,
}
}10 秒是这次缓存关闭等待的超时,不是整个 resume RPC 的统一截止时间。不能在超时后直接删除旧句柄并创建新对象;旧执行者可能仍然持有 writer、工具或后台任务。
这个决策适合画成资源替换流程,而不是一个虚构的 Thread 生命周期枚举。
检查覆盖不只是比较 model。命名权限、通用 config、base/developer instructions 等字段的出现本身也会被记录为差异;即使某个 HashMap 中的值碰巧等于当前值,处理器也不会假装已经逐层验证其等价性。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:collect_resume_override_mismatches 的复杂覆盖。
if request.permissions.is_some() {
mismatch_details.push(format!(
"permissions override was provided and ignored while running; active={:?}",
config_snapshot.active_permission_profile
));
}
if let Some(requested_personality) = request.personality.as_ref()
&& config_snapshot.personality.as_ref() != Some(requested_personality)
{
mismatch_details.push(format!(
"personality requested={requested_personality:?} active={:?}",
config_snapshot.personality
));
}
if request.config.is_some() {
mismatch_details
.push("config overrides were provided and ignored while running".to_string());
}
if request.base_instructions.is_some() {
mismatch_details
.push("baseInstructions override was provided and ignored while running".to_string());
}
if request.developer_instructions.is_some() {
mismatch_details.push(
"developerInstructions override was provided and ignored while running".to_string(),
);
}
mismatch_details对可观察的已加载对象,响应里的 model、cwd、权限和 instruction sources 来自 live 配置快照。请求被接受不等于配置更新成功;判断生效应看响应字段,而不是只看 JSON-RPC 没有返回错误。
3.2 持久设置合并
冷恢复先构造类型化覆盖,再把 config.approval_policy 提升到类型化字段,随后继承持久设置。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的 approval_policy 提取。
let history_cwd = thread_history.session_cwd();
let runtime_workspace_roots = runtime_workspace_roots.map(resolve_runtime_workspace_roots);
let mut typesafe_overrides = self.build_thread_config_overrides(
model,
model_provider,
service_tier,
cwd,
runtime_workspace_roots,
approval_policy,
approvals_reviewer,
sandbox,
permissions,
base_instructions,
developer_instructions,
personality,
);
if typesafe_overrides.approval_policy.is_none()
&& let Some(value) = request_overrides
.as_mut()
.and_then(|overrides| overrides.remove("approval_policy"))
{
let approval_policy = match serde_json::from_value(value) {
Ok(approval_policy) => approval_policy,
Err(err) => {
self.outgoing
.send_error(
request_id,
invalid_params(format!(
"invalid `approval_policy` config override: {err}"
)),
)
.await;
return Ok(());
}
};
typesafe_overrides.approval_policy = Some(approval_policy);
}只有类型化 approval policy 缺省时才从 HashMap 移除并解析该键;解析失败是参数错误。这个操作让显式请求覆盖在后续合并中保持可见,避免持久 approval policy 把用户本次选择盖掉。
最近的安全设置并不只存在 SQLite:历史中最近的 TurnContext 或 ThreadSettingsApplied 也可能是来源。
源码文件:codex-rs/app-server/src/request_processors/persisted_resume_settings.rs
相关函数/类型:latest_persisted_resume_settings。
pub(super) fn latest_persisted_resume_settings(
history: &[RolloutItem],
) -> Option<PersistedResumeSettings> {
history
.iter()
.enumerate()
.rev()
.find_map(|(index, item)| match item {
RolloutItem::TurnContext(turn_context) => Some(PersistedResumeSettings {
approval_policy: turn_context.approval_policy,
approvals_reviewer: turn_context.approvals_reviewer.or_else(|| {
history[..index].iter().rev().find_map(|item| match item {
RolloutItem::TurnContext(turn_context) => turn_context.approvals_reviewer,
RolloutItem::EventMsg(EventMsg::ThreadSettingsApplied(event)) => {
Some(event.thread_settings.approvals_reviewer)
}
_ => None,
})
}),
active_permission_profile: turn_context.active_permission_profile.clone(),
}),
RolloutItem::EventMsg(EventMsg::ThreadSettingsApplied(event)) => {
Some(PersistedResumeSettings {
approval_policy: event.thread_settings.approval_policy,
approvals_reviewer: Some(event.thread_settings.approvals_reviewer),
active_permission_profile: event
.thread_settings
.active_permission_profile
.clone(),
})
}
_ => None,
})
}从后向前扫描得到最近设置;旧 TurnContext 缺少 reviewer 时,还会向更早记录寻找该字段。这里恢复的是字段来源,不是简单取最后一个 JSON 对象。随后处理器只在相应显式覆盖缺省时采用这些值。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:load_and_apply_persisted_resume_metadata。
async fn load_and_apply_persisted_resume_metadata(
&self,
thread_history: &InitialHistory,
request_overrides: &mut Option<HashMap<String, serde_json::Value>>,
typesafe_overrides: &mut ConfigOverrides,
) -> Option<ThreadMetadata> {
let InitialHistory::Resumed(resumed_history) = thread_history else {
return None;
};
if let Some(persisted_settings) = latest_persisted_resume_settings(&resumed_history.history)
{
if typesafe_overrides.approval_policy.is_none() {
typesafe_overrides.approval_policy = Some(persisted_settings.approval_policy);
}
if typesafe_overrides.approvals_reviewer.is_none()
&& !request_overrides
.as_ref()
.is_some_and(|overrides| overrides.contains_key("approvals_reviewer"))
{
typesafe_overrides.approvals_reviewer = persisted_settings.approvals_reviewer;
}
if !has_permission_override(request_overrides.as_ref(), typesafe_overrides) {
typesafe_overrides.persisted_permission_profile_id = persisted_settings
.active_permission_profile
.map(|profile| profile.id);
}
}
let state_db_ctx = self.state_db.clone()?;
let persisted_metadata = state_db_ctx
.get_thread(resumed_history.conversation_id)
.await
.ok()
.flatten()?;
merge_persisted_resume_metadata(request_overrides, typesafe_overrides, &persisted_metadata);
Some(persisted_metadata)
}权限恢复保留的是 profile ID,让当前配置重新解析;不是把旧 sandbox 的路径集合直接安装进新 Session。SQLite 缺失或读取失败在此被转换为 None,因此这是可选元数据补充,不是整个恢复必需的数据源。恢复历史本身失败仍会由前面的 Store 路径拒绝。
模型还有一条成组优先级:显式 model、provider 或 effort 相关覆盖出现时,不再把旧模型组合部分混入。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:merge_persisted_resume_metadata / has_model_resume_override。
fn merge_persisted_resume_metadata(
request_overrides: &mut Option<HashMap<String, serde_json::Value>>,
typesafe_overrides: &mut ConfigOverrides,
persisted_metadata: &ThreadMetadata,
) {
if has_model_resume_override(request_overrides.as_ref(), typesafe_overrides) {
return;
}
typesafe_overrides.model = persisted_metadata.model.clone();
typesafe_overrides.model_provider = Some(persisted_metadata.model_provider.clone());
if let Some(reasoning_effort) = persisted_metadata.reasoning_effort.as_ref() {
request_overrides.get_or_insert_with(HashMap::new).insert(
"model_reasoning_effort".to_string(),
serde_json::Value::String(reasoning_effort.to_string()),
);
}
}
// ...
fn has_model_resume_override(
request_overrides: Option<&HashMap<String, serde_json::Value>>,
typesafe_overrides: &ConfigOverrides,
) -> bool {
typesafe_overrides.model.is_some()
|| typesafe_overrides.model_provider.is_some()
|| request_overrides.is_some_and(|overrides| overrides.contains_key("model"))
|| request_overrides
.is_some_and(|overrides| overrides.contains_key("model_reasoning_effort"))
}函数没有把所有 config key 都当模型覆盖;例如 has_model_resume_override 明确检查类型化 model/provider,以及 map 中的 model、model_reasoning_effort。阅读时应按这组条件推演,不能扩写成“任何 config 覆盖都会屏蔽持久模型”。
最终仍要经过完整配置加载器。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的配置加载。
let has_explicit_model_resume_override =
has_model_resume_override(request_overrides.as_ref(), &typesafe_overrides);
let persisted_metadata = self
.load_and_apply_persisted_resume_metadata(
&thread_history,
&mut request_overrides,
&mut typesafe_overrides,
)
.await;
// Derive a Config using the same logic as new conversation, honoring overrides if provided.
let mut config = match self
.config_manager
.load_for_cwd(request_overrides, typesafe_overrides, history_cwd)
.await
{
Ok(config) => config,
Err(err) => {
let error = config_load_error(&err);
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
};
if !has_explicit_model_resume_override
&& persisted_metadata
.as_ref()
.is_some_and(|metadata| metadata.reasoning_effort.is_none())
{
config.model_reasoning_effort = None;
}history_cwd 作为恢复背景传入 load_for_cwd;显式字段与配置文件继续按加载器规则合成。若没有显式模型覆盖,且持久元数据明确缺少 effort,代码会把当前配置 effort 清成 None,而不是偷偷套用服务器默认 effort。响应要报告实际解析结果。
3.3 权限重解析
测试 cold_resume_reresolves_persisted_active_permission_profile 分别覆盖 Legacy 与 Paginated:先用 dev profile 运行并持久化,再把同名 profile 改为 read-only 后重启服务。关键断言如下。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:cold_resume_reresolves_persisted_active_permission_profile。
assert!(matches!(sandbox, AppSandboxPolicy::ReadOnly { .. }));
assert_eq!(
active_permission_profile,
Some(ActivePermissionProfile {
id: "dev".to_string(),
extends: Some(BUILT_IN_PERMISSION_PROFILE_READ_ONLY.to_string()),
})
);
assert!(
!runtime_workspace_roots.contains(&AbsolutePathBuf::from_absolute_path(
previous_workspace_root.path(),
)?)
);profile 身份仍是 dev,但权限变成当前的 read-only,旧 runtime workspace root 不再沿用。它证明冷恢复重新解释配置,不能外推到热重连会即时重载配置。
另一项测试删除旧 profile,并配置新的默认权限。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:cold_resume_with_removed_permission_profile_uses_configured_default。
MockResponsesConfig::new(&server.uri())
.with_root_config(&format!(
"default_permissions = \"{BUILT_IN_PERMISSION_PROFILE_DANGER_FULL_ACCESS}\""
))
.write(codex_home.path())?;
// ...
assert!(matches!(sandbox, AppSandboxPolicy::DangerFullAccess));
assert_eq!(
active_permission_profile,
Some(ActivePermissionProfile::new(
BUILT_IN_PERMISSION_PROFILE_DANGER_FULL_ACCESS,
))
);结果采用当前默认 profile。这里的 danger-full-access 是测试专门设置的输入,不是恢复机制的固定默认值,更不是遇到配置缺失时一律放宽权限的规则。修改恢复逻辑时,这两项测试必须一起看:一个覆盖“同名定义改变”,另一个覆盖“原 ID 消失”。
4. 历史双通道
4.1 模型恢复输入
Paginated Thread 的完整 UI 历史可能非常大,模型只需要最新上下文窗口。冷恢复按 history mode 选择模型输入。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:load_resume_initial_history_from_stored_thread。
async fn load_resume_initial_history_from_stored_thread(
&self,
stored_thread: StoredThread,
) -> Result<(InitialHistory, StoredThread), JSONRPCErrorError> {
if matches!(stored_thread.history_mode, ThreadHistoryMode::Paginated) {
let model_context = self
.thread_store
.load_latest_model_context(StoreLoadThreadHistoryParams {
thread_id: stored_thread.thread_id,
include_archived: true,
})
.await
.map_err(thread_store_resume_read_error)?;
let history = InitialHistory::Resumed(ResumedHistory {
conversation_id: model_context.thread_id,
history: Arc::new(model_context.items),
rollout_path: stored_thread.rollout_path.clone(),
});
return Ok((history, stored_thread));
}
let thread_id = stored_thread.thread_id.to_string();
let rollout_path = stored_thread.rollout_path.clone();
let mut stored_thread = self
.read_stored_thread_for_resume(
&thread_id,
rollout_path.as_ref(),
/*include_history*/ true,
)
.await?;
let history = self
.stored_thread_to_initial_history(&mut stored_thread)
.await?;
Ok((history, stored_thread))
}Paginated 使用 load_latest_model_context;Legacy 则重读 include_history=true,再转换为 InitialHistory::Resumed。不能把 load_latest_model_context 改成 UI 的 list_turns:后者包含展示投影,缺少模型恢复需要的上下文记录,且其条目结构不同。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:stored_thread_to_initial_history。
#[tracing::instrument(level = "trace", skip_all)]
async fn stored_thread_to_initial_history(
&self,
stored_thread: &mut StoredThread,
) -> Result<InitialHistory, JSONRPCErrorError> {
let thread_id = stored_thread.thread_id;
let history = stored_thread
.history
.take()
.map(|history| history.items)
.ok_or_else(|| {
internal_error(format!(
"thread {thread_id} did not include persisted history"
))
})?;
Ok(InitialHistory::Resumed(ResumedHistory {
conversation_id: thread_id,
history: Arc::new(history),
rollout_path: stored_thread.rollout_path.clone(),
}))
}history.take() 把 Vec 从摘要中移出,再由 Arc 承载恢复历史。后面 response_history = thread_history.clone() 可以与 Core 输入共享这份记录;它不意味着允许响应裁剪去修改底层模型历史。
图中两条路径会在客户端看来汇合到一次 resume,但消费者不同。
图中的“后续模型请求”不是 resume 函数立即发起的用户 Turn。恢复首先建立可继续工作的状态;新输入以及扩展生命周期是否触发后续工作,要看对应消费者。
4.2 Core 状态安装
Core 的恢复分支重建上下文,恢复部分状态与 token usage,并推迟写入新的初始上下文。
源码文件:codex-rs/core/src/session/mod.rs
相关函数/类型:Session::record_initial_history 的 Resumed 分支。
InitialHistory::Resumed(resumed_history) => {
let turn_context = self.new_default_turn().await;
let rollout_items = resumed_history.history;
if matches!(
rollout_items.iter().rev().find_map(|item| match item {
RolloutItem::EventMsg(event) => agent_status_from_event(event),
_ => None,
}),
Some(AgentStatus::Interrupted)
) {
self.agent_status.send_replace(AgentStatus::Interrupted);
}
let previous_turn_settings = self
.apply_rollout_reconstruction(&turn_context, &rollout_items)
.await;
// ...
// Seed usage info from the recorded rollout so UIs can show token counts
// immediately on resume/fork.
if let Some(info) = Self::last_token_info_from_rollout(&rollout_items) {
let mut state = self.state.lock().await;
state.set_token_info(Some(info));
}
// Defer seeding the session's initial context until the first turn starts so
// turn/start overrides can be merged before we write to the rollout.
if !is_subagent {
let _ = self.flush_rollout().await;
}
None
}agent_status 只在记录中最新可识别状态为 Interrupted 时特别恢复为中断;不会因为历史末尾有 TurnStarted,就重新创建原来的异步 Turn task。apply_rollout_reconstruction 返回 previous settings,供模型差异提示等逻辑使用;token info 进入 SessionState,让界面能够在下一次推理前显示已有消耗。
状态安装还有容易遗漏的媒体与元数据配对约束。
源码文件:codex-rs/core/src/session/mod.rs
相关函数/类型:apply_rollout_reconstruction 的媒体与状态安装。
let (mut prepared_history, metadata): (Vec<_>, Vec<_>) = history
.into_iter()
.map(|envelope| (envelope.item, envelope.metadata))
.unzip();
let _ = prepare_image_response_items(
&mut prepared_history,
ImagePreparationMode::DetailBased,
ImageResizeNoticeMode::Disabled,
);
prepare_audio_response_items(&mut prepared_history);
assert_eq!(
prepared_history.len(),
metadata.len(),
"replay media preparation must remain one-to-one when resize notices are disabled"
);
history = prepared_history
.into_iter()
.zip(metadata)
.map(|(item, metadata)| ResponseItemEnvelope { item, metadata })
.collect();
{
let mut state = self.state.lock().await;
state.replace_annotated_history(history, reference_context_item);
if let Some(world_state) = world_state_baseline {
state.history.set_world_state_baseline(world_state);
}
let fallback_ids = state.auto_compact_window_ids();
let window_id = window_id.unwrap_or(fallback_ids.window_id);
state.restore_auto_compact_window(
window_number,
AutoCompactWindowIds {
first_window_id: first_window_id.unwrap_or(window_id),
previous_window_id,
window_id,
},
);
state.set_previous_turn_settings(previous_turn_settings.clone());
}历史 envelope 被拆成 item 和 metadata 两个 Vec,媒体准备禁用额外 resize notice,然后检查长度不变,再 zip 回去。这样来源等旁路信息不会因多插一个提示条目而错位。最后持有 SessionState 锁替换 annotated history,并恢复 world-state baseline;这与 UI 的 ThreadHistoryBuilder 不是同一个归约器。
5. 冷恢复装配
5.1 身份与注册
App Server 把配置与初始历史交给 ThreadManager,并另留 response_history 供展示响应使用。下面省略成功分支的装配内容,先看调用与错误出口。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的 Core 调用。
let response_history = thread_history.clone();
match self
.thread_manager
.resume_thread_with_history(
config,
thread_history,
self.auth_manager.clone(),
self.request_trace_context(&request_id).await,
client_mcp_extensions,
)
.await
{
Ok(NewThread {
thread_id,
thread: codex_thread,
session_configured,
..
}) => {
// ...
}
Err(err) => {
let error = match err.details() {
CodexErrorDetails::InvalidRequest(message) => invalid_request(message.clone()),
_ => internal_error(format!("error resuming thread: {err}")),
};
self.outgoing.send_error(request_id, error).await;
}
}
Ok(())成功返回的是 NewThread 及其 session_configured;Core 的 InvalidRequest 保留请求错误性质,其余错误包装为内部恢复失败。下面继续进入这次调用,观察来源继承和注册,而非根据 resume 名称假定只是载入文本。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:ThreadManager::resume_thread_with_history。
pub async fn resume_thread_with_history(
&self,
config: Config,
initial_history: InitialHistory,
auth_manager: Arc<AuthManager>,
parent_trace: Option<W3cTraceContext>,
client_mcp_extensions: ClientMcpExtensions,
) -> CodexResult<NewThread> {
let agent_control = self.agent_control_for_config(&config);
let (session_source, thread_source) = initial_history
.get_resumed_session_sources()
.unwrap_or_else(|| (self.state.session_source.clone(), None));
if let InitialHistory::Resumed(resumed) = &initial_history
&& initial_history.get_multi_agent_version() == Some(MultiAgentVersion::V2)
&& !session_source.is_non_root_agent()
{
agent_control
.restore_v2_agent_metadata(&config, resumed.conversation_id)
.await;
}
let options = StartThreadOptions {
initial_history,
session_source: Some(session_source),
thread_source,
parent_trace,
client_mcp_extensions,
..StartThreadOptions::new(config)
};
Box::pin(self.state.spawn_thread(ThreadSpawnRequest::new(
options,
auth_manager,
agent_control,
)))
.await
}V2 根 agent 还会先恢复 agent metadata;普通恢复来源优先来自记录,而不是当前客户端名称。StartThreadOptions 把这些来源、trace 与 MCP 扩展传进 spawn_thread。这一步建立 Session 服务与 I/O,不是调用旧 Turn 的 Rust 栈继续执行。
Session ID 也要从实现确认。
源码文件:codex-rs/core/src/session/session.rs
相关函数/类型:Session::new 的 session_id 选择。
let resumed_session_id = match &initial_history {
InitialHistory::Resumed(resumed) => {
resumed.history.iter().find_map(|item| match item {
RolloutItem::SessionMeta(meta_line) => Some(meta_line.meta.session_id),
_ => None,
})
}
InitialHistory::New | InitialHistory::Cleared | InitialHistory::Forked(_) => None,
};
// Legacy subagent rollouts synthesize session_id from their own thread id.
let resumed_session_id = resumed_session_id.filter(|session_id| {
!session_configuration.session_source.is_non_root_agent()
|| *session_id != SessionId::from(thread_id)
});
let session_id = resumed_session_id.unwrap_or_else(|| {
if session_configuration.session_source.is_non_root_agent() {
agent_control.session_id()
} else {
SessionId::from(thread_id)
}
});Resumed 优先读 SessionMeta.session_id;兼容旧 sub-agent 的特定值会被过滤,再用 agent control 或 Thread ID 回退。热重连直接沿用现有对象的 session configured 值。所以“冷恢复必然换 Session ID”并不成立,应该区分新建内存对象与协议身份字段。
创建管线在启动成功后才完成 manager 注册。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:ThreadManagerState::finalize_thread_spawn。
async fn finalize_thread_spawn(
&self,
session: Arc<Session>,
io: SessionIo,
session_source: SessionSource,
) -> CodexResult<NewThread> {
let thread_id = session.thread_id();
let event = io.next_event().await?;
let session_configured = match event {
Event {
id,
msg: EventMsg::SessionConfigured(session_configured),
} if id == INITIAL_SUBMIT_ID => session_configured,
_ => {
return Err(CodexErr::SessionConfiguredNotFirstEvent);
}
};
{
let mut threads = self.threads.write().await;
if let std::collections::hash_map::Entry::Vacant(e) = threads.entry(thread_id) {
let thread = Arc::new(CodexThread::new(
session,
io,
session_configured.clone(),
session_configured.rollout_path.clone(),
session_source,
));
e.insert(thread.clone());
return Ok(NewThread {
thread_id,
thread,
session_configured,
});
}
}
if let Err(err) = io.shutdown_and_wait().await {
warn!("failed to shut down duplicate thread {thread_id}: {err}");
}
Err(CodexErr::InvalidRequest(format!(
"thread {thread_id} is already running"
)))
}首个事件必须是指定初始提交 ID 的 SessionConfigured。拿到后在 threads.write() 下做 Entry::Vacant 检查;若 ID 已被另一条路径注册,退出锁后关闭重复 I/O,再返回 already running。锁保护的是注册原子性,不会在持锁期间执行整个 Session 初始化或 shutdown。
5.2 响应装配
成功返回 NewThread 后,App Server 还有工作可能失败。尤其 Paginated 要让重新打开的 writer 把投影追到可读状态。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的冷恢复后处理。
let instruction_sources = codex_thread.legacy_instruction_sources().await;
let SessionConfiguredEvent { rollout_path, .. } = session_configured;
let Some(rollout_path) = rollout_path else {
let error =
internal_error(format!("rollout path missing for thread {thread_id}"));
self.outgoing.send_error(request_id, error).await;
return Ok(());
};
// Paginated JSONL is canonical, but its SQLite projection can lag after a
// previous write failure. Persist after reopening the live writer so legacy
// response hydration reads the latest durable turns and items.
if paginated_resume
&& let Err(error) = self
.thread_store
.persist_thread(thread_id, PersistContext::Standard)
.await
.map_err(thread_store_resume_read_error)
{
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
let materialized_turns = if paginated_resume && include_turns {
match self.paginated_thread_full_turns(thread_id).await {
Ok(turns) => Some(turns),
Err(error) => {
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
}
} else {
None
};必须有 rollout path;Paginated 的 persist_thread(Standard) 失败就发错误。接着按 include_turns 决定是否全量补取 UI 轮次。模型上下文已经恢复,并不自动意味着响应的全部展示历史已经准备好。
随后建立 listener,把连接与 Thread 关联起来。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的 listener 建立。
log_listener_attach_result(
self.ensure_conversation_listener(
thread_id,
request_id.connection_id,
/*raw_events_enabled*/ false,
)
.await,
thread_id,
request_id.connection_id,
"thread",
);这一调用的返回值交给 log_listener_attach_result,失败走日志处理,而不是在这里直接 ? 返回。不能据此承诺“resume 响应成功就绝对已经完成所有订阅工作”;应结合监听器状态和后续事件验证。摘要加载、watch upsert 与残尾修正完成后,配置字段从实际 session_configured 和快照填入响应。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的响应与尾声。
let thread_originator = config_snapshot.originator.clone();
let response = ThreadResumeResponse {
thread,
model: session_configured.model,
model_provider: session_configured.model_provider_id,
service_tier: session_configured.service_tier,
cwd: session_configured.cwd,
runtime_workspace_roots: config_snapshot.workspace_roots,
instruction_sources,
approval_policy: session_configured.approval_policy.into(),
approvals_reviewer: session_configured.approvals_reviewer.into(),
sandbox,
active_permission_profile,
reasoning_effort: session_configured.reasoning_effort,
multi_agent_mode: MultiAgentMode::ExplicitRequestOnly,
initial_turns_page,
turns_backwards_cursor,
items_backwards_cursor,
};
let connection_id = request_id.connection_id;
self.outgoing
.send_response_with_thread_originator(request_id, response, thread_originator)
.await;
// `excludeTurns` is explicitly the cheap resume path, so avoid
// rebuilding history only to attribute a replayed usage update.
if let Some(token_usage_turn_id) = token_usage_turn_id {
// The client needs restored usage before it starts another turn.
// Sending after the response preserves JSON-RPC request ordering while
// still filling the status line before the next turn lifecycle begins.
send_thread_token_usage_update_to_connection(
&self.outgoing,
connection_id,
thread_id,
codex_thread.as_ref(),
token_usage_turn_id,
)
.await;
}
self.thread_goal_processor
.emit_resume_goal_snapshot(thread_id)
.await;
codex_thread
.emit_thread_idle_lifecycle_if_idle(ThreadIdleCause::Completed)
.await;代码先把响应排队,再按可归属的 Turn ID 回放 token usage,然后发 goal 快照,最后让 idle 生命周期消费者反应。token usage 不只依赖 excludeTurns:冷 Paginated 可以从模型恢复历史确定归属;若无法确定有效轮次则不能随便借最后一个显示 Turn 的 ID。上述 await 保证该分支的排队顺序,不保证网络客户端已经渲染完画面。
6. 热重连交接
6.1 待发送请求
热重连直接从现有对象获取配置,按需要加载 Legacy 历史或 Paginated 投影,再确保 listener 存在。它不调用 resume_thread_with_history 来替换当前执行者。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:resume_running_thread 的监听准备。
let thread_state = self
.thread_state_manager
.thread_state(existing_thread_id)
.await;
self.ensure_listener_task_running(
existing_thread_id,
existing_thread.clone(),
thread_state.clone(),
)
.await?;
Self::set_app_server_client_info(
existing_thread.as_ref(),
app_server_client_name,
app_server_client_version,
)
.await?;真正的恢复响应被包装成一条 listener command。请求对象带着本次历史快照、配置和分页选项跨越任务边界。
源码文件:codex-rs/app-server/src/thread_state.rs
相关函数/类型:PendingThreadResumeRequest。
pub(crate) struct PendingThreadResumeRequest {
pub(crate) request_id: ConnectionRequestId,
pub(crate) history_items: Vec<RolloutItem>,
/// Usage attribution already resolved while cold-loading a paginated child.
pub(crate) cold_resume_token_usage_turn_id: Option<String>,
pub(crate) config_snapshot: ThreadConfigSnapshot,
pub(crate) instruction_sources: Vec<LegacyAppPathString>,
pub(crate) thread_summary: codex_app_server_protocol::Thread,
pub(crate) emit_thread_goal_update: bool,
pub(crate) thread_goal_state_db: Option<StateDbHandle>,
pub(crate) include_turns: bool,
pub(crate) initial_turns_page:
Option<codex_app_server_protocol::ThreadResumeInitialTurnsPageParams>,
pub(crate) paginated_turns: Option<Vec<Turn>>,
pub(crate) paginated_initial_turns_page: Option<codex_app_server_protocol::TurnsPage>,
pub(crate) paginated_initial_turns_page_with_active_slot:
Option<codex_app_server_protocol::TurnsPage>,
pub(crate) resume_cursor_store: Option<Arc<dyn codex_thread_store::ThreadStore>>,
pub(crate) redact_resume_payloads: bool,
}结构体没有接管 CodexThread 的所有权:运行对象仍由 manager/listener 持有。这里的 history_items 属于本次响应,config_snapshot 属于捕获时的运行配置,paginated_initial_turns_page_with_active_slot 则是为尚未落进历史页的实时 Turn 预留空间。不要把这些字段理解为另一份完整 Session。
以下关系图只画跨任务交接的数据所有权。
关键不在多一层包装,而在由同一个消费者决定“历史响应”和“后续事件”的交接顺序。发送端在 channel 已关闭时明确报错。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:resume_running_thread 的命令投递。
let resume_cursor_store = paginated_resume.then(|| Arc::clone(&self.thread_store));
let command = crate::thread_state::ThreadListenerCommand::SendThreadResumeResponse(
Box::new(crate::thread_state::PendingThreadResumeRequest {
request_id: request_id.clone(),
history_items,
cold_resume_token_usage_turn_id,
config_snapshot,
instruction_sources,
thread_summary,
emit_thread_goal_update,
thread_goal_state_db,
include_turns,
initial_turns_page: params.initial_turns_page.clone(),
paginated_turns,
paginated_initial_turns_page,
paginated_initial_turns_page_with_active_slot,
resume_cursor_store,
redact_resume_payloads,
}),
);
if listener_command_tx.send(command).is_err() {
return Err(internal_error(format!(
"failed to enqueue running thread resume for thread {existing_thread_id}: thread listener command channel is closed"
)));
}
return Ok(RunningThreadResumeResult::Handled);send(command) 成功之后,外层就返回 Handled。恢复响应尚未在这里发出;若 listener 消费前 Thread 进入卸载,消费者必须再次检查状态。
6.2 同一消费者
listener 在一个 tokio::select! 中接收取消、命令和 Core 事件,并在选中命令后 await 完整处理函数。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:ensure_listener_task_running 的 select 分支。
tokio::spawn(async move {
loop {
tokio::select! {
biased;
_ = &mut cancel_rx => {
// Listener was superseded or the thread is being torn down.
break;
}
listener_command = listener_command_rx.recv() => {
let Some(listener_command) = listener_command else {
break;
};
handle_thread_listener_command(
conversation_id,
&conversation,
codex_home.as_path(),
&thread_state_manager,
&thread_state,
&thread_watch_manager,
&outgoing_for_task,
&pending_thread_unloads,
listener_command,
)
.await;
}
event = conversation.next_event() => {
let event = match event {
Ok(event) => event,
Err(err) => {
tracing::warn!("thread.next_event() failed with: {err}");
break;
}
};
// ...
}
// ...
}
}
let mut thread_state = thread_state.lock().await;
if thread_state.listener_generation == listener_generation {
thread_state_manager.unregister_listener_command_tx(conversation_id);
thread_state.clear_listener();
}
});biased 使同时就绪时按分支顺序选择:取消在前,命令在事件之前。但这不等于 resume 能抢占一个已经执行中的事件处理;只有回到 select 的边界才重新选择。处理恢复命令期间,同一个 listener 不会又去消费下一个 Core 事件,这正是交接的局部顺序保证。
处理器此时才取 active turn 快照,并与之前读到的持久历史合并。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:handle_pending_thread_resume_request 的 live 合并。
let active_turn = {
let state = thread_state.lock().await;
state.active_turn_snapshot()
};
tracing::debug!(
thread_id = %conversation_id,
request_id = ?pending.request_id,
active_turn_present = active_turn.is_some(),
active_turn_id = ?active_turn.as_ref().map(|turn| turn.id.as_str()),
active_turn_status = ?active_turn.as_ref().map(|turn| &turn.status),
"composing running thread resume response"
);
let has_live_in_progress_turn =
matches!(conversation.agent_status().await, AgentStatus::Running)
|| active_turn
.as_ref()
.is_some_and(|turn| matches!(turn.status, TurnStatus::InProgress));
let request_id = pending.request_id;
let connection_id = request_id.connection_id;
let mut thread = pending.thread_summary;
if pending.include_turns {
if let Some(turns) = pending.paginated_turns.take() {
thread.turns = turns;
} else {
populate_thread_turns_from_history(
&mut thread,
&pending.history_items,
/*active_turn*/ None,
);
}
if let Some(active_turn) = active_turn.as_ref() {
merge_turn_history_with_active_turn(&mut thread.turns, active_turn.clone());
}
}
let thread_status = thread_watch_manager
.loaded_status_for_thread(&thread.id)
.await;
set_thread_status_and_interrupt_stale_turns(
&mut thread,
thread_status.clone(),
has_live_in_progress_turn,
);has_live_in_progress_turn 结合 Core Running 与 listener 的 InProgress 快照,避免仅因 watch 更新稍晚就把真实活动轮次标为中断。完整历史合并按 Turn ID 去重,再用最新快照替换。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:merge_turn_history_with_active_turn。
pub(super) fn merge_turn_history_with_active_turn(turns: &mut Vec<Turn>, active_turn: Turn) {
turns.retain(|turn| turn.id != active_turn.id);
turns.push(active_turn);
}这不是把每个 item delta 再与旧列表做拼接;active snapshot 已代表 listener 当前聚合的完整 Turn,整体替换避免同一个命令或用户条目出现两遍。代价是持久读取与快照捕获不是全局事务快照,必须把保证限定在同一 listener 的事件交接。
6.3 订阅与断连
响应组装完成后,listener 再检查 closing,并尝试加入连接。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:handle_pending_thread_resume_request 的连接接入。
{
let pending_thread_unloads = pending_thread_unloads.lock().await;
if pending_thread_unloads.contains(&conversation_id) {
drop(pending_thread_unloads);
outgoing
.send_error(
request_id,
invalid_request(format!(
"thread {conversation_id} is closing; retry thread/resume after the thread is closed"
)),
)
.await;
return;
}
if !thread_state_manager
.try_add_connection_to_thread(conversation_id, connection_id)
.await
{
tracing::debug!(
thread_id = %conversation_id,
connection_id = ?connection_id,
"skipping running thread resume for closed connection"
);
return;
}
}这次检查不能因为入口查过一次就删掉:排队和读历史都跨越 await,卸载可能在中途开始。若连接已经关闭,try_add_connection_to_thread 返回 false,函数不再发送这个恢复响应。关闭连接不会被当成应该重新创建 Session 的理由。
时序图强调 App Server 的 listener 串行边界;Core 仍可继续运行并把事件留在队列里。
thread_resume_keeps_in_flight_turn_streaming 先运行一轮产生历史,再启动第二轮,收到 turn/started 后调用 resume。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:thread_resume_keeps_in_flight_turn_streaming。
let resume_id = primary
.send_thread_resume_request(ThreadResumeParams {
thread_id: thread.id,
..Default::default()
})
.await?;
let ThreadResumeResponse {
thread: resumed_thread,
..
} = timeout(DEFAULT_READ_TIMEOUT, primary.read_response(resume_id)).await??;
assert_ne!(resumed_thread.status, ThreadStatus::NotLoaded);
timeout(
DEFAULT_READ_TIMEOUT,
primary.read_stream_until_notification_message("turn/completed"),
)
.await??;它断言恢复后不是 NotLoaded,并继续收到原 Turn 的 completed。这说明 resume 没有把进行中的任务替换掉;该测试不要求响应瞬间一定 Active,因为原 Turn 可能在命令真正执行前完成。把断言改成无条件 Active,反而会把正常并发时序误判为失败。
7. 分页接续
7.1 初始页与全量
初始页是 resume 附带的一次 thread/turns/list 结果,默认倒序和 Summary;它与 thread.turns 是否填充是独立选择。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:ThreadResumeInitialTurnsPageParams / TurnsPage。
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct ThreadResumeInitialTurnsPageParams {
/// Optional turn page size.
#[ts(optional = nullable)]
pub limit: Option<u32>,
/// Optional turn pagination direction; defaults to descending.
#[ts(optional = nullable)]
pub sort_direction: Option<SortDirection>,
/// How much item detail to include for each returned turn; defaults to summary.
#[ts(optional = nullable)]
pub items_view: Option<TurnItemsView>,
}
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(export_to = "v2/")]
pub struct TurnsPage {
pub data: Vec<Turn>,
pub next_cursor: Option<String>,
pub backwards_cursor: Option<String>,
}excludeTurns=true 搭配 initialTurnsPage 才能在减少全量展示历史的同时取得最近一页。Legacy 的初始页先归约历史,再按分页选项构造;Paginated 使用 Store 查询,不应把两条路径的 I/O 成本写成相同。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:build_thread_resume_initial_turns_page。
pub(super) fn build_thread_resume_initial_turns_page(
items: &[RolloutItem],
loaded_status: ThreadStatus,
has_live_running_thread: bool,
active_turn: Option<Turn>,
params: &ThreadResumeInitialTurnsPageParams,
) -> Result<codex_app_server_protocol::TurnsPage, JSONRPCErrorError> {
build_thread_turns_page_response(
items,
loaded_status,
has_live_running_thread,
active_turn,
ThreadTurnsPageOptions {
cursor: None,
limit: params.limit,
sort_direction: params.sort_direction.unwrap_or(SortDirection::Desc),
items_view: params.items_view.unwrap_or(TurnItemsView::Summary),
},
)
.map(Into::into)
}active_turn 是可选快照,冷路径传 None,热路径由 listener 提供。分页最终处理的对象已经是协议 Turn,不会反过来缩短模型的恢复上下文。
7.2 实时轮次占位
最棘手的场景是:请求倒序最近一页,limit=1,活动 Turn 还没有出现在 Store 返回的页中。若直接在已有一条记录后再插活动 Turn,会超出 limit;若直接丢弃旧记录并保留原 next cursor,又会跳过那条被挤掉的历史。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:paginated_resume_initial_turns_page_with_active_slot。
async fn paginated_resume_initial_turns_page_with_active_slot(
&self,
thread_id: ThreadId,
params: &ThreadResumeInitialTurnsPageParams,
) -> Result<codex_app_server_protocol::TurnsPage, JSONRPCErrorError> {
// A running resume overlays the newest live turn on this durable page.
// Reserve one row so the overlay keeps the requested limit and the
// durable next cursor still starts after the last returned stored turn.
let page_size = thread_turns_page_size(params.limit);
if page_size == 1 {
// ThreadStore does not accept an empty page. Use its backwards cursor as
// the next cursor so the omitted durable row is returned next.
let mut page = self
.paginated_resume_initial_turns_page(thread_id, params)
.await?;
page.next_cursor = page.backwards_cursor.clone();
page.data.clear();
return Ok(page);
}
let mut params = params.clone();
params.limit = Some((page_size - 1) as u32);
self.paginated_resume_initial_turns_page(thread_id, ¶ms)
.await
}limit 大于 1 时,备用页少取一条,为活动 Turn 留位置;limit 为 1 时,Store 不接受空页,所以仍取一条,再清空 data,并把 backwards cursor 用作 next cursor,使本次未返回的持久记录在下一页还能出现。cursor 不只是“最后看过哪条”,还必须和本次实际返回的集合一致。
listener 确认活动 Turn 不在倒序页内时,才换用这份备用页。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:handle_pending_thread_resume_request 的初始页选择。
let mut initial_turns_page = if let Some(mut page) = pending.paginated_initial_turns_page.take()
{
if let (Some(active_turn), Some(params)) =
(active_turn, pending.initial_turns_page.as_ref())
{
let sort_direction = params.sort_direction.unwrap_or(SortDirection::Desc);
let active_turn_is_in_page = page.data.iter().any(|turn| turn.id == active_turn.id);
if matches!(sort_direction, SortDirection::Desc)
&& !active_turn_is_in_page
&& let Some(page_with_active_slot) =
pending.paginated_initial_turns_page_with_active_slot.take()
{
page = page_with_active_slot;
}
merge_active_turn_into_page(&mut page, active_turn, params);
}
super::thread_processor::normalize_thread_turns_status(
&mut page.data,
thread_status,
has_live_in_progress_turn,
);
Some(page)
} else if let Some(params) = pending.initial_turns_page.as_ref() {
match super::thread_processor::build_thread_resume_initial_turns_page(
&pending.history_items,
thread.status.clone(),
has_live_in_progress_turn,
active_turn,
params,
) {
Ok(page) => Some(page),
Err(error) => {
outgoing.send_error(request_id, error).await;
return;
}
}
} else {
None
};如果活动 Turn 已经在页中,就直接替换该 Turn,不额外占一个位置。随后根据 items view 裁剪快照,避免请求 NotLoaded 却收到 live 的完整 items。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:merge_active_turn_into_page。
fn merge_active_turn_into_page(
page: &mut codex_app_server_protocol::TurnsPage,
mut active_turn: Turn,
params: &codex_app_server_protocol::ThreadResumeInitialTurnsPageParams,
) {
super::thread_processor::apply_thread_turns_items_view(
std::slice::from_mut(&mut active_turn),
params.items_view.unwrap_or(TurnItemsView::Summary),
);
let sort_direction = params.sort_direction.unwrap_or(SortDirection::Desc);
let page_size = super::thread_processor::thread_turns_page_size(params.limit);
let active_turn_is_in_page = page.data.iter().any(|turn| turn.id == active_turn.id);
page.data.retain(|turn| turn.id != active_turn.id);
match sort_direction {
SortDirection::Asc
if active_turn_is_in_page
|| (page.data.len() < page_size && page.next_cursor.is_none()) =>
{
page.data.push(active_turn);
}
SortDirection::Asc => {}
SortDirection::Desc => page.data.insert(0, active_turn),
}
}倒序总是把新活动轮次放在开头;升序只有活动轮次原本就在这一页,或已到末页且还有容量时才追加。它不会为了“让用户看到正在执行”而把最新轮次塞进任意一页早期历史,破坏排序语义。
该关系可以用具体的 limit=1 情形核对。
这里的 T1/T2 是帮助推演的轮次名。实现没有要求它们按字符串排序;真正排序与 cursor 由 Store 和分页参数决定。
thread_resume_rejoins_running_paginated_thread_with_initial_page 用受控制的模型响应保持 Turn 未完成,发送 excludeTurns=true、倒序 limit 1、NotLoaded 的请求。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:thread_resume_rejoins_running_paginated_thread_with_initial_page。
let initial_turns_page = initial_turns_page.expect("resume should include initial turns page");
assert_eq!(initial_turns_page.data.len(), 1);
let resumed_running_turn = initial_turns_page
.data
.first()
.expect("resume page should include the running turn");
assert_eq!(resumed_running_turn.id, running_turn.id);
assert_eq!(resumed_running_turn.items_view, TurnItemsView::NotLoaded);
assert!(resumed_running_turn.items.is_empty());
assert_eq!(resumed_running_turn.status, TurnStatus::InProgress);
assert!(initial_turns_page.backwards_cursor.is_some());
assert!(initial_turns_page.next_cursor.is_some());
// ...
let initial_turns_page = initial_turns_page.expect("resume should include initial turns page");
assert_eq!(initial_turns_page.data.len(), 1);
assert_eq!(initial_turns_page.data[0].id, seed_turn.id);
// The running-thread resume response is queued onto the thread listener task.
// If the in-flight turn completes before that queued command runs, the response
// can legitimately observe the thread as idle.
match &thread.status {
ThreadStatus::Active { active_flags } => assert!(active_flags.is_empty()),
ThreadStatus::Idle => {}
status => panic!("unexpected thread status after running resume: {status:?}"),
}倒序页必须只有一个活动 Turn、items 为空、状态 InProgress,并携带两个 cursor;改成升序后,第一页应是较早的 seed Turn。末尾允许 Active 或 Idle 的断言,再次提醒我们响应观测的是 listener 处理时刻,不是客户端开始发送请求的时刻。这个 fixture 验证活动页语义,并不证明任意外部 Store 的 cursor 实现都正确。
7.3 反向游标
响应顶层还有 turnsBackwardsCursor 与 itemsBackwardsCursor,供客户端从恢复时的存储位置向过去补取。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:paginated_resume_backwards_cursors。
pub(super) async fn paginated_resume_backwards_cursors(
thread_store: &dyn ThreadStore,
thread_id: ThreadId,
) -> Result<(Option<String>, Option<String>), JSONRPCErrorError> {
let turns_page = thread_store
.list_turns(StoreListTurnsParams {
thread_id,
include_archived: true,
cursor: None,
page_size: 1,
sort_direction: StoreSortDirection::Desc,
items_view: StoredTurnItemsView::NotLoaded,
})
.await
.map_err(paginated_history_list_error)?;
let items_page = thread_store
.list_items(StoreListItemsParams {
thread_id,
turn_id: None,
include_archived: true,
cursor: None,
page_size: 1,
sort_direction: StoreSortDirection::Desc,
sort_key: StoreItemSortKey::CreatedAtOrdinal,
after_updated_at_ordinal: None,
})
.await
.map_err(paginated_history_list_error)?;
Ok((turns_page.backwards_cursor, items_page.backwards_cursor))
}两次查询都只取倒序一条,但分别针对 Turn 和所有 item;它们不是同一种 cursor,不能互换。返回值取 Store 的 backwards cursor,协议注释规定按 Desc 使用时第一页包含游标指向的记录。即使 excludeTurns=true,Paginated resume 仍可能执行这些轻量查询,因此它只是省掉全量历史装配,并不等于零存储访问。
8. 待审批请求
8.1 响应之后重放
重新接入正在等待审批的 Thread,只回传一个状态为 InProgress 的 Turn 还不够。客户端需要再次收到原来的服务端请求,才能把批准或拒绝交回仍在等待的工具。
源码文件:codex-rs/app-server/src/request_processors/thread_lifecycle.rs
相关函数/类型:handle_pending_thread_resume_request 的发送尾声。
outgoing
.send_response_with_thread_originator(request_id, response, originator)
.await;
// Warm metadata-only resumes skip history reconstruction. Cold paginated children can
// replay usage using attribution captured before the listener was attached.
if let Some(token_usage_turn_id) = token_usage_turn_id {
// Rejoining a loaded thread has the same UI contract as a cold resume, but
// uses the live conversation state instead of reconstructing a new session.
send_thread_token_usage_update_to_connection(
outgoing,
connection_id,
conversation_id,
conversation.as_ref(),
token_usage_turn_id,
)
.await;
}
if pending.emit_thread_goal_update {
if let Some(state_db) = pending.thread_goal_state_db {
send_thread_goal_snapshot_notification(outgoing, conversation_id, &state_db).await;
} else {
tracing::warn!(
thread_id = %conversation_id,
"state db unavailable when reading thread goal for running thread resume"
);
}
}
outgoing
.replay_requests_to_connection_for_thread(connection_id, conversation_id)
.await;
// App-server owns resume response and snapshot ordering, so wait until
// replay completes before letting extensions react to the idle thread.
conversation
.emit_thread_idle_lifecycle_if_idle(ThreadIdleCause::Completed)
.await;顺序是 resume response、可归属的 token usage、goal 快照、pending server requests,最后才触发 idle 生命周期。goal 的读取失败会记录 warning;它不应使已经排队的恢复响应被“撤回”。这些尾声也解释了为什么不应在业务处理器刚读完文件时就直接发热重连响应。
源码文件:codex-rs/app-server/src/outgoing_message.rs
相关函数/类型:replay_requests_to_connection_for_thread。
pub(crate) async fn replay_requests_to_connection_for_thread(
&self,
connection_id: ConnectionId,
thread_id: ThreadId,
) {
let requests = self.pending_requests_for_thread(thread_id).await;
for request in requests {
if let Err(err) = self
.sender
.send(OutgoingEnvelope::ToConnection {
connection_id,
message: OutgoingMessage::Request(request),
write_complete_tx: None,
})
.await
{
warn!("failed to resend request to client: {err:?}");
}
}
}重放把 pending 表中已有的 request 发送给指定连接,沿用原请求 ID,不重新生成审批请求。发送失败只是 warning;write_complete_tx: None 也说明这一步没有等待对端读完或渲染确认。pending callback 的生命周期仍由原服务端请求管理。
源码文件:codex-rs/app-server/src/outgoing_message.rs
相关函数/类型:notify_client_response。
pub(crate) async fn notify_client_response(&self, id: RequestId, result: Result) {
let entry = self.take_request_callback(&id).await;
match entry {
Some((id, entry)) => {
let completed_at_ms = now_unix_timestamp_ms();
if let Ok(response) = entry.request.response_from_result(result.clone()) {
tracing::info!("<- response: {response:?}");
if !matches!(response, ServerResponse::PermissionsRequestApproval { .. }) {
self.analytics_events_client
.track_server_response(completed_at_ms, response);
}
}
if entry.callback.send(Ok(result)).is_err() {
warn!("could not notify callback for {id:?}: receiver dropped");
}
}
None => {
warn!("could not find callback for {id:?}");
}
}
}客户端回复后,take_request_callback 取走对应回调,再送入原等待者。重复回复没有第二个 callback 可唤醒。resume 的职责是恢复请求的可见性,不是恢复一个已不存在的操作系统进程,也不是重新执行工具命令。
8.2 原请求一致性
thread_resume_replays_pending_command_execution_request_approval 先取得原命令审批请求,在它未解决时 resume,再读取服务端重放请求。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:thread_resume_replays_pending_command_execution_request_approval。
let replayed_request = timeout(
DEFAULT_READ_TIMEOUT,
primary.read_stream_until_request_message(),
)
.await??;
pretty_assertions::assert_eq!(replayed_request, original_request);
let ServerRequest::CommandExecutionRequestApproval { request_id, .. } = replayed_request else {
panic!("expected CommandExecutionRequestApproval request");
};
primary
.send_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Accept,
})?,
)
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
primary.read_stream_until_notification_message("turn/completed"),
)
.await??;
wait_for_responses_request_count(&server, /*expected_count*/ 3).await?;测试比较整个 request,而不只比较 method;随后用重放请求中的同一 ID 批准,并等待原 Turn 完成。它证明 pending request 关联没有被 resume 破坏。文件变更审批有对应的 thread_resume_replays_pending_file_change_request_approval 测试。两者验证的是同一服务进程内的热重连,不能外推成跨进程冷恢复会重新启动旧工具并还原旧回调。
9. 载荷裁剪
某些远程客户端不适合接收完整 MCP 输出和图像生成载荷。这个处理位于响应层,必须与模型恢复历史分开。
源码文件:codex-rs/app-server/src/request_processors/thread_resume_redaction.rs
相关函数/类型:客户端选择与载荷裁剪。
// Temporary bandaid for remote clients: thread/resume can include large MCP and
// image-generation payloads. Keep this response-only so persisted rollout
// history, model resume history, and other APIs stay unchanged.
const REDACTED_PAYLOAD: &str = "[redacted]";
const CHATGPT_REMOTE_CLIENT_NAMES: &[&str] =
&["codex_chatgpt_android_remote", "codex_chatgpt_ios_remote"];
pub(super) fn should_redact_thread_resume_payloads(client_name: Option<&str>) -> bool {
client_name.is_some_and(|client_name| CHATGPT_REMOTE_CLIENT_NAMES.contains(&client_name))
}
pub(super) fn redact_thread_resume_payloads(turns: &mut [Turn]) {
for turn in turns {
turn.items.retain_mut(|item| match item {
ThreadItem::McpToolCall {
arguments,
result,
error,
..
} => {
*arguments = JsonValue::String(REDACTED_PAYLOAD.to_string());
if result.is_some() {
*result = Some(Box::new(redacted_mcp_tool_call_result()));
}
if let Some(error) = error {
error.message = REDACTED_PAYLOAD.to_string();
}
true
}
ThreadItem::ImageGeneration(_) => false,
_ => true,
});
}
}
fn redacted_mcp_tool_call_result() -> McpToolCallResult {
McpToolCallResult {
content: vec![serde_json::json!({
"type": "text",
"text": REDACTED_PAYLOAD,
})],
structured_content: None,
meta: None,
}
}匹配的是两个精确客户端名称。MCP 调用保留条目壳,但参数变成 [redacted];存在 result 才替换 result,存在 error 才覆盖 message;图像生成条目直接移除。app_context、只读提示等未在匹配分支中改动,不能把它描述成“清除 MCP 所有元数据”。
冷路径和热路径都会分别处理 thread.turns 与 initial_turns_page.data。调用发生在协议 Turn 值已经组装之后,既不修改 Store 中的 rollout,也不修改送进 Core 的 model context。它是兼容特定客户端载荷大小的措施,不构成全部 API 的保密边界;同一历史仍可由其他读取接口呈现。
测试对两个远程名字分别取完整 turns 和初始页,并对 MCP 内容与图像条目做断言。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:thread_resume_redacts_payloads_for_chatgpt_remote_clients。
assert_eq!(read_only_hint, &Some(false));
let result = result.as_ref().expect("redacted MCP result");
assert_eq!(
result.content,
vec![json!({
"type": "text",
"text": "[redacted]",
})]
);
assert_eq!(result.structured_content, None);
assert_eq!(result.meta, None);
assert_eq!(error, &None);
assert!(
!remote_turn
.items
.iter()
.any(|item| matches!(item, ThreadItem::ImageGeneration(_))),
"remote resume should drop image generation items for {client_name}"
);相反,some_other_client 的 fixture 保留原参数、structured content、meta 和图像数据。这个对照能检查裁剪没有意外扩大到普通客户端,但不能证明任意后续新增 ThreadItem 都自动受保护;新增条目落到 _ => true 会原样保留,需要逐项审查它是否应参加载荷裁剪。
10. 恢复故障
10.1 所有者边界
V2 的 parent-owned child 不允许通过公共 resume 的配置覆盖变成独立运行者。外层在重建配置前识别此情况,要求通过已加载父对象恢复它。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_resume_inner 的子线程所有者检查。
// Parent-owned V2 children must resume through their owner, not caller configuration.
if let InitialHistory::Resumed(resumed_history) = &thread_history
&& let Some((source, _)) = thread_history.get_resumed_session_sources()
&& !can_accept_direct_input(thread_history.get_multi_agent_version(), &source)
{
let child_thread_id = resumed_history.conversation_id;
self.thread_manager
.ensure_multi_agent_v2_child_loaded(child_thread_id)
.await
.map_err(|err| {
tracing::warn!(
thread_id = %child_thread_id,
error = %err,
"failed to resume a multi-agent v2 child through its parent"
);
invalid_request(
"cannot resume an unloaded multi-agent v2 sub-agent through its parent; resume the parent first, or use thread/read to inspect it",
)
})?;
let cold_resume_history = paginated_resume.then(|| thread_history.get_rollout_items());
// Attach to the resolved child with only the caller's history-paging preferences.
let attach_params = ThreadResumeParams {
thread_id: child_thread_id.to_string(),
exclude_turns,
initial_turns_page,
..Default::default()
};
return match self
.resume_running_thread(
&request_id,
&attach_params,
app_server_client_name,
app_server_client_version,
cold_resume_history,
)
.await?
{
RunningThreadResumeResult::Handled => Ok(()),
RunningThreadResumeResult::NotRunning(_) => Err(invalid_request(
"cannot resume an unloaded multi-agent v2 sub-agent through its parent; resume the parent first, or use thread/read to inspect it",
)),
};
}成功取得 child 后,只保留调用者的历史分页偏好,再走热重连;不把请求中的 model、sandbox 或 instructions 带给子线程重新建配置。父对象未加载时,错误提示指向先恢复 parent,或用 read 检查 child。这个边界取决于 multi-agent 版本与来源,不是所有 sub-agent 都一律适用。
进程间则由 writer 所有权防止并行恢复同一个持久线程。上游 thread_resume_rejects_legacy_writer_owned_by_another_process 和 Paginated 对应测试先启动主 App Server 并完成一轮,再让第二个 App Server 共享同一 Codex home、使用独立 SQLite 目录恢复该 ID,检查失败。manager 的 map 只能保护本进程,不能替代这层存储排他性。
源码文件:codex-rs/app-server/tests/suite/v2/thread_resume.rs
相关函数/类型:assert_thread_resume_rejects_writer_owned_by_another_process。
assert_eq!(error.error.code, -32600);
assert_eq!(
error.error.message,
format!("thread {} already has an active writer", thread.id)
);
timeout(DEFAULT_READ_TIMEOUT, primary.shutdown_gracefully()).await??;
let next_resume_id = secondary
.send_thread_resume_request(ThreadResumeParams {
thread_id: thread.id.clone(),
..Default::default()
})
.await?;
let _: ThreadResumeResponse = timeout(
DEFAULT_READ_TIMEOUT,
secondary.read_response(next_resume_id),
)
.await??;错误码为 -32600,消息明确指出 active writer;主进程正常关闭后,同一请求可以成功。独立 SQLite 目录避免把此失败误认为同一数据库连接冲突;它证明了 writer 的排他与释放路径,不覆盖进程被强杀时所有文件系统的锁恢复时序。
10.2 失败后的状态
恢复包含多个可失败阶段,应按最后完成的阶段判断资源,而不是把所有错误都理解成“什么都没发生”。
| 失败点 | 已有状态 | 排查方向 |
|---|---|---|
| closing、history 为空、互斥权限 | 替代对象尚未创建 | 参数与 pending unload 集合 |
| 未物化、归档、stale path、Store unsupported | 尚未创建新的 Core 对象 | 来源 ID、当前路径、history mode、后端能力 |
| 配置加载或 required MCP 初始化 | 新 Session 创建管线失败 | config error、服务初始化错误,不应把它归为历史 JSON 损坏 |
| 缓存 shutdown 超时 | 旧运行对象仍保留 | 不得删除缓存后强行再启动一个 writer |
| Core 成功后 persist、历史装配或 cursor 失败 | manager 可能已有 Thread | 先确认 loaded 状态,再决定重试;响应不是事务提交点 |
| 热重连命令队列已关闭 | 现有 Core 对象未由此代码销毁 | listener 生命周期、connection 状态与 channel |
| pending request 发送失败 | 原 callback 可能仍在等回复 | 出站连接与请求 ID,不能重跑工具代替回复 |
在 thread_resume_inner 的成功分支里,多处后处理失败只发送错误并 return,没有配对的 remove_thread。这意味着 RPC 错误不能证明 manager 中没有对象;客户端重试时可能从冷路径切换为热路径。也不能把热重连中的 closing 复查误写成全流程回滚机制。
历史解析器的容错范围同样需要收住:即使某种 rollout 读取允许跳过损坏行,也不代表能够完整还原丢失的工具输出或任务状态。本文的恢复契约是从可用记录构建当前运行状态,不是操作系统级检查点恢复。Legacy/Paginated、所选 Store、实验字段与 multi-agent 版本是本篇的主要条件;平台路径差异由路径归一化与存储实现处理,没有一套独立的 Windows resume 状态机。
10.3 定向复现
在 Codex 仓库根目录执行以下命令。第一条同时覆盖读取、来源、配置、分页、运行中重连、token usage 与审批重放;第二条只选 App Server 的恢复辅助逻辑。
just test --locked -p codex-app-server --test all \
-E 'test(suite::v2::thread_read::) or test(suite::v2::thread_resume::)'
just test --locked -p codex-app-server --lib \
-E 'test(resume) or test(persisted_resume_settings)'
rg -n 'resume_running_thread|load_resume_initial_history_from_stored_thread' \
codex-rs/app-server/src/request_processors/thread_processor.rs
rg -n 'handle_pending_thread_resume_request|replay_requests_to_connection_for_thread' \
codex-rs/app-server/src/request_processors/thread_lifecycle.rs \
codex-rs/app-server/src/outgoing_message.rs可以用三个场景检查自己是否能顺着源码定位。第一,传入新 model 但返回旧 model:先确认已加载对象是否有订阅者、是否在运行,以及缓存关闭是否完成。第二,倒序第一页只有活动 Turn,下一页却少一轮:检查是否选择了 active-slot 页,以及 limit=1 时 next cursor 是否指回被挤出的历史。第三,resume 后仍等待批准:检查原 pending request 是否在 response 后重放,客户端是否用同一个 request ID 回复。
最后尝试从 ClientRequest::ThreadResume 复述两条完整路径:冷恢复到 InitialHistory、Config、Session 注册和响应;热重连到 listener command、active snapshot、连接加入、响应和原审批 callback。能解释它们在哪里共享历史投影、在哪里保留不同的执行者与时序,才能判断一次恢复失败应该重试什么,而不是重复创建对话。
