ThreadFork处理流程
假设一个对话已经完成两轮,第三轮正在执行工具。用户想从第二轮另开一条路线:新对话应该看见前两轮,源对话应继续工作;如果用户改成“从现在分叉”,新对话又不能把源线程正在运行的工具当作自己的任务继续等待。更难的是,新对话可能并没有复制那两轮的文本,却仍能在界面和模型输入中看到它们。
这篇文章面向掌握 Rust Result、Arc、异步调用和锁的读者。先读ThreadStart处理流程了解新运行对象的装配,再读ThreadResume处理流程理解恢复已有执行者的路径。这里从 thread/fork 的请求字段进入,跟踪历史边界如何变成新 Thread 的可运行状态。Core 通用的按用户消息序号截断等算法见ThreadManager分叉;本文着重解释 App Server 的具名 Turn 边界、Store 引用、配置选择和对外交付。
先区分四个概念。ThreadId 是对话的逻辑身份;rollout 是记录事件和模型条目的持久日志;Turn 是协议中的任务轮次;HistoryPosition 是某个物理 rollout 文件的排他性结束位置。一次 fork 会创建新的 Thread 和执行状态,历史则可能采用复制或引用。它不复制操作系统进程、工具回调或 Git 工作树,也不是 shell 环境快照。下文源码均为实际实现的摘录,跨段省略使用 // ... 标明。
1. 输入与约束
1.1 API 入口
公开方法首先登记实验字段检查与请求串行化范围。
源码文件:codex-rs/app-server-protocol/src/protocol/common.rs
相关函数/类型:ClientRequest::ThreadFork;行号:526-531。
ThreadFork => "thread/fork" {
params: v2::ThreadForkParams,
inspect_params: true,
serialization: thread_or_path(params.thread_id, params.path),
response: v2::ThreadForkResponse,
},inspect_params 允许逐字段检查实验能力;thread_or_path 选择请求队列键。它不等于历史来源最终解析出的身份,尤其 path 与 ID 可能是同一记录的两种入口。具体排队规则已在恢复文章中展开,本篇继续追踪业务分派。
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数/类型:MessageProcessor::process_request;行号:1150-1160。
ClientRequest::ThreadFork { params, .. } => {
self.thread_processor
.thread_fork(
request_id.clone(),
params,
app_server_client_name.clone(),
client_version.clone(),
client_mcp_extensions.clone(),
)
.await
}客户端名称、版本与 MCP 扩展传给分叉处理器;创建出的 Core 对象需要这些连接上下文。分派不直接把原 Thread 的句柄克隆成一个新 Thread。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork;行号:551-568。
pub(crate) async fn thread_fork(
&self,
request_id: ConnectionRequestId,
params: ThreadForkParams,
app_server_client_name: Option<String>,
app_server_client_version: Option<String>,
client_mcp_extensions: ClientMcpExtensions,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
self.thread_fork_inner(
request_id,
params,
app_server_client_name,
app_server_client_version,
client_mcp_extensions,
)
.await
.map(|()| None)
}外层把内部成功映射为 None,因为分叉分支会自行排队响应和 thread/started 通知。若它返回 Err,通用分派出口才发送错误。读代码时要把“函数返回完成”和“响应已由该分支接管”分开,否则容易误以为客户端收到一个空结果。
1.2 截断与返回
选择分叉位置使用的是 Turn ID,不是任意 ThreadItem.id,也不是聊天数组的下标。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:ThreadForkParams 的来源与边界;行号:519-543。
pub struct ThreadForkParams {
pub thread_id: String,
/// Optional last turn id to fork through, inclusive.
///
/// When specified, turns after `last_turn_id` are omitted from the fork.
/// The referenced turn cannot be in progress.
#[ts(optional = nullable)]
pub last_turn_id: Option<String>,
/// Optional turn id to fork before, excluding that turn and all later turns.
/// Cannot be combined with `last_turn_id`.
#[experimental("thread/fork.beforeTurnId")]
#[ts(optional = nullable)]
pub before_turn_id: Option<String>,
/// [UNSTABLE] Specify the rollout path to fork from.
/// If specified, the thread_id param will be ignored.
#[experimental("thread/fork.path")]
#[serde(
default,
deserialize_with = "crate::protocol::serde_helpers::deserialize_empty_path_as_none"
)]
#[ts(optional = nullable)]
pub path: Option<PathBuf>,last_turn_id 包含目标轮次,但目标不能仍在执行;before_turn_id 排除目标轮次及后续记录,属于实验字段。非空 path 会选择源记录,空字符串按缺省处理。与 resume 已加载 ID 的重连语义不同,fork 没有先复用源执行者的分支;它通过 Store 读取指定来源。
配置字段决定新执行者的运行条件,和选择历史前缀是两条输入链。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:ThreadForkParams 配置字段;行号:545-578。
/// Configuration overrides for the forked thread, if any.
#[ts(optional = nullable)]
pub model: Option<String>,
#[ts(optional = nullable)]
pub model_provider: Option<String>,
#[serde(
default,
deserialize_with = "crate::protocol::serde_helpers::deserialize_double_option",
serialize_with = "crate::protocol::serde_helpers::serialize_double_option",
skip_serializing_if = "Option::is_none"
)]
#[ts(optional = nullable)]
pub service_tier: Option<Option<String>>,
#[ts(optional = nullable)]
pub cwd: Option<String>,
/// Replace the thread's runtime workspace roots. Paths must be absolute.
#[experimental("thread/fork.runtimeWorkspaceRoots")]
#[ts(optional = nullable)]
pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>,
#[experimental(nested)]
#[ts(optional = nullable)]
pub approval_policy: Option<AskForApproval>,
/// Override where approval requests are routed for review on this thread
/// and subsequent turns.
#[ts(optional = nullable)]
pub approvals_reviewer: Option<ApprovalsReviewer>,
#[ts(optional = nullable)]
pub sandbox: Option<SandboxMode>,
/// Named profile id for the forked thread. Cannot be combined with
/// `sandbox`.
#[experimental("thread/fork.permissions")]
#[ts(optional = nullable)]
pub permissions: Option<String>,
#[ts(optional = nullable)]model/model_provider 选择模型与提供方;cwd 和必须为绝对路径的 runtime_workspace_roots 决定目录背景与运行工作区。service_tier 的双层 Option 区分缺省、不使用指定档位的显式 null、具体档位字符串。审批策略决定何时请求审批,reviewer 决定交给谁审阅;sandbox 与具名 permissions 二选一。缺省项交给后文的配置装配,不表示所有字段自动复制源 Thread。
返回历史与生命周期还有独立开关。
源码文件:codex-rs/app-server-protocol/src/protocol/v2/thread.rs
相关函数/类型:ThreadForkParams 的响应与生命周期选项;行号:579-601。
pub config: Option<HashMap<String, serde_json::Value>>,
#[ts(optional = nullable)]
pub base_instructions: Option<String>,
#[ts(optional = nullable)]
pub developer_instructions: Option<String>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub ephemeral: bool,
/// Optional client-supplied analytics source classification for this forked thread.
#[ts(optional = nullable)]
pub thread_source: Option<ThreadSource>,
/// When true, return only thread metadata and live fork state without
/// populating `thread.turns`. This is useful when the client plans to call
/// `thread/turns/list` immediately after forking.
#[experimental("thread/fork.excludeTurns")]
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub exclude_turns: bool,
/// When true, carry the source thread's current goal into the fork without
/// starting its initial automatic continuation. The next explicit turn owns
/// the goal lifecycle, and normal automatic continuation resumes after it.
#[experimental("thread/fork.deferGoalContinuation")]
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub defer_goal_continuation: bool,
}exclude_turns 只减少响应中的展示历史,不会让模型失去继承上下文。ephemeral 决定新 Thread 是否持久化。defer_goal_continuation 不只是一个显示选项:它要求继承源 goal,并延后新 Thread 的首次自动继续;该动作依赖持久状态,所以不能随意与 ephemeral 组合。
1.3 前置拒绝
先看源记录读取与参数检查的实际顺序。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的来源与约束;行号:4634-4663。
let include_turns = !exclude_turns;
if sandbox.is_some() && permissions.is_some() {
return Err(invalid_request(
"`permissions` cannot be combined with `sandbox`",
));
}
let source_thread = self
.read_stored_thread_for_resume(
&thread_id,
path.as_ref(),
/*include_history*/ false,
)
.await?;
let paginated_source = matches!(source_thread.history_mode, ThreadHistoryMode::Paginated);
if last_turn_id.is_some() && before_turn_id.is_some() {
return Err(invalid_request(
"`beforeTurnId` cannot be combined with `lastTurnId`",
));
}
if ephemeral && defer_goal_continuation {
return Err(invalid_request(
"`deferGoalContinuation` cannot be combined with `ephemeral`",
));
}
if paginated_source && ephemeral && include_turns {
return Err(invalid_request(
"ephemeral paginated thread/fork requires `excludeTurns: true`",
));
}
let source_thread_id = source_thread.thread_id;sandbox/permissions 的冲突先于 Store 读取拒绝;边界互斥和 ephemeral 组合检查在源摘要加载后执行。读取复用 read_stored_thread_for_resume,因此未物化、归档、Paginated 旧路径等错误仍可能发生,不能因为操作叫 fork 就跳过来源有效性检查。
特别注意 ephemeral && paginated_source && include_turns:这里明确要求 excludeTurns=true。它不是“所有临时分叉都没有历史”,Legacy 的临时分叉仍可在响应中展示复制的轮次。三个开关决定的是不同层面的能力。
2. 具名截断
2.1 Legacy 包含边界
Legacy 先取得记录数组,再按请求截取。lastTurnId 要先在有效历史视图中找到轮次,然后回到原始记录寻找稳定边界。
源码文件:codex-rs/core/src/thread_rollout_truncation.rs
相关函数/类型:truncate_rollout_after_turn_id;行号:164-209。
pub fn truncate_rollout_after_turn_id(
mut items: Vec<RolloutItem>,
last_turn_id: &str,
) -> CodexResult<Vec<RolloutItem>> {
let turns = build_turns_from_rollout_items(&items);
let turn = turns
.iter()
.find(|turn| turn.id == last_turn_id)
.ok_or_else(|| {
CodexErr::InvalidRequest(format!(
"lastTurnId '{last_turn_id}' was not found in the source thread"
))
})?;
let target_start_index = items
.iter()
.position(|item| {
matches!(
item,
RolloutItem::EventMsg(EventMsg::TurnStarted(event))
if event.turn_id == last_turn_id
)
})
.ok_or_else(|| {
CodexErr::InvalidRequest(format!(
"lastTurnId '{last_turn_id}' is not a persisted canonical turn in the source thread"
))
})?;
if matches!(turn.status, TurnStatus::InProgress) {
return Err(CodexErr::InvalidRequest(format!(
"lastTurnId '{last_turn_id}' identifies an in-progress turn"
)));
}
let cut_index = items
.iter()
.enumerate()
.skip(target_start_index.saturating_add(1))
.find_map(|(index, item)| {
matches!(item, RolloutItem::EventMsg(EventMsg::TurnStarted(_))).then_some(index)
})
.unwrap_or(items.len());
items.truncate(cut_index);
Ok(items)
}这段算法有三个必要步骤。第一,build_turns_from_rollout_items 将回滚等事件作用到有效历史,避免选中已经被回滚的轮次。第二,原始数组里必须有该 ID 的 TurnStarted;某些旧日志能在 UI 中合成轮次 ID,但合成 ID 不能充当稳定的原始日志切点。第三,只有终态轮次可被包含,InProgress 会被拒绝。
实际 cut_index 是目标开始之后的下一个 TurnStarted,没有下一轮则取数组长度。它并非简单找到 TurnComplete 后立即截断:完成事件之后、下一轮开始之前的附属记录也可能留在前缀里。这里返回的是 rollout 前缀,不能用“返回前 N 个 UI item”替代算法。
2.2 Legacy 排除边界
beforeTurnId 的问题不同:只要能找到目标的开始位置,就能把整个目标轮次排除,不需要它已经完成。
源码文件:codex-rs/core/src/thread_rollout_truncation.rs
相关函数/类型:truncate_rollout_before_turn_id;行号:212-255。
pub fn truncate_rollout_before_turn_id(
mut items: Vec<RolloutItem>,
before_turn_id: &str,
) -> CodexResult<Vec<RolloutItem>> {
let cut_index = items.iter().position(|item| {
matches!(
item,
RolloutItem::EventMsg(EventMsg::TurnStarted(event))
if event.turn_id == before_turn_id
)
});
let Some(cut_index) = cut_index else {
// Older rollouts can expose generated turn IDs without a TurnStarted item to fork at.
if build_turns_from_rollout_items(&items)
.iter()
.any(|turn| turn.id == before_turn_id)
{
return Err(CodexErr::InvalidRequest(format!(
"beforeTurnId '{before_turn_id}' is not a persisted canonical turn in the source thread"
)));
}
return Err(CodexErr::InvalidRequest(format!(
"beforeTurnId '{before_turn_id}' was not found in the source thread"
)));
};
// A persisted turn boundary proves the turn exists unless a later rollback removes it.
if items[cut_index + 1..]
.iter()
.any(|item| matches!(item, RolloutItem::EventMsg(EventMsg::ThreadRolledBack(_))))
&& !build_turns_from_rollout_items(&items)
.iter()
.any(|turn| turn.id == before_turn_id)
{
return Err(CodexErr::InvalidRequest(format!(
"beforeTurnId '{before_turn_id}' was not found in the source thread"
)));
}
items.truncate(cut_index);
Ok(items)
}先找真实 TurnStarted。找不到时,再区分“UI 中存在合成 ID”和“目标完全不存在”,给出不同错误。有真实开始边界时,只有后面出现过 rollback,才需要重新投影确认目标没有被撤销。最后 truncate(cut_index) 不包含目标开始记录,所以可以从未完成轮次之前分叉。
两条路径不能合并成一个“找到 ID 后截取”的辅助函数:包含目标需要验证终态,排除目标则允许切掉未完成的工作;两者对旧日志与回滚的检查成本也不同。
2.3 三种选择
对前两轮完成、第三轮进行中的同一来源,边界语义如下。
| 请求 | 保留历史 | 对第三轮的处理 |
|---|---|---|
lastTurnId=第二轮 | 包含第二轮的有效前缀 | 不进入新 Thread |
beforeTurnId=第三轮 | 第三轮开始之前 | 不要求第三轮先结束 |
| 不传边界,即 Latest | 当前可取得的历史快照 | 在新 Thread 的快照中收口为中断 |
lastTurnId=第三轮 | 无结果 | 拒绝包含 InProgress 轮次 |
Latest 的收口在 Core 进行,后文会展开;它不会向源 Thread 提交 Interrupt。Paginated 则先把这些请求转换成 Store 的边界类型。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的 PreparedFork 分流;行号:4668-4699。
let prepared_fork = if paginated_source {
let boundary = match (last_turn_id.as_deref(), before_turn_id.as_deref()) {
(Some(turn_id), None) => {
codex_thread_store::ForkBoundary::ThroughTurn(turn_id.to_string())
}
(None, Some(turn_id)) => {
codex_thread_store::ForkBoundary::BeforeTurn(turn_id.to_string())
}
(None, None) => codex_thread_store::ForkBoundary::Latest,
(Some(_), Some(_)) => unreachable!("fork boundaries are mutually exclusive"),
};
Some(
self.thread_store
.prepare_fork(codex_thread_store::PrepareForkParams {
thread_id: source_thread_id,
boundary,
})
.await
.map_err(|err| match err {
ThreadStoreError::InvalidRequest { message } => invalid_request(message),
ThreadStoreError::ThreadNotFound { thread_id } => {
invalid_request(format!("no rollout found for thread id {thread_id}"))
}
ThreadStoreError::Unsupported { .. } => {
method_not_found("paginated_threads is not supported yet")
}
err => internal_error(format!("failed to prepare paginated fork: {err}")),
})?,
)
} else {
None
};PrepareForkParams.thread_id 使用解析后的 source_thread_id,而非继续相信请求里的 ID。后端不支持 prepare 时映射为方法不可用;持久来源找不到、非法边界和内部准备失败也保持不同错误类别。Legacy 不经过 prepare_fork,由数组截断负责边界。
上游测试创建 first/second/third 三轮,再按第二轮 ID fork,检查整个返回轮次 ID 列表和源文件内容。
源码文件:codex-rs/app-server/tests/suite/v2/thread_fork.rs
相关函数/类型:assert_thread_fork_at_named_boundary_keeps_only_terminal_prefix;行号:621-657,659-673。
}
let original_contents = std::fs::read_to_string(source_path.as_path())?;
let fork_id = mcp
.send_thread_fork_request(ThreadForkParams {
thread_id: source_thread_id.clone(),
last_turn_id: Some(turn_ids[1].clone()),
..Default::default()
})
.await?;
let ThreadForkResponse {
thread: forked_thread,
..
} = timeout(DEFAULT_READ_TIMEOUT, mcp.read_response(fork_id)).await??;
assert_eq!(
forked_thread
.turns
.iter()
.map(|turn| turn.id.clone())
.collect::<Vec<_>>(),
turn_ids[..2]
);
assert!(
forked_thread
.turns
.iter()
.all(|turn| turn.status == TurnStatus::Completed)
);
assert_eq!(forked_thread.forked_from_id, Some(source_thread_id.clone()));
if history_mode == ThreadHistoryMode::Legacy {
assert_eq!(forked_thread.preview, "first");
}
assert_eq!(
std::fs::read_to_string(source_path.as_path())?,
original_contents,
"forking at a turn must not mutate the source rollout"
// ...
let forked_path = forked_thread.path.clone().expect("forked thread path");
let forked_contents = std::fs::read_to_string(forked_path.as_path())?;
if history_mode == ThreadHistoryMode::Paginated {
assert!(
read_session_meta_line(forked_path.as_path())
.await?
.meta
.history_base
.is_some()
);
assert!(!forked_contents.contains(turn_ids[1].as_str()));
} else {
assert!(forked_contents.contains(turn_ids[1].as_str()));
}预期新 Thread 只含前两轮且都是 Completed,forked_from_id 指向源 Thread,源 rollout 的内容不变。对 Paginated,目标轮次文本不会被复制进新文件,文件中却有 history_base;对 Legacy,该目标轮次会出现在复制的文件中。因此“新文件里找不到旧 Turn ID”本身不是丢历史的证据。
3. 固定历史前缀
3.1 三类身份
Paginated 把继承范围保存为物理位置,而不是把任意长度的旧历史重复写进子文件。
源码文件:codex-rs/protocol/src/protocol.rs
相关函数/类型:HistoryPosition;行号:2863-2877。
/// Exclusive position in another rollout's paginated history.
#[derive(Serialize, Deserialize, Clone, Copy, Debug, PartialEq, Eq, JsonSchema, TS)]
pub struct HistoryPosition {
/// Rollout ID for the immutable prefix file.
///
/// `HistoryPosition` predates `thread/revert`, so this field is named `thread_id`. Treat its
/// value as a `rollout_id`: ordinary rollouts use the thread ID as their rollout ID, while a
/// reverted thread's filename carries a distinct rollout ID. It is not necessarily
/// [`SessionMeta::id`], which remains the stable thread ID across revert.
pub thread_id: ThreadId,
/// First rollout ordinal not included from the prefix file.
pub end_ordinal_exclusive: u64,
/// Byte offset immediately after the last included JSONL record from the prefix file.
pub end_byte_offset: u64,
}end_ordinal_exclusive 是第一条不包含的记录序号;end_byte_offset 是最后一条被包含的 JSONL 记录之后的字节位置。二者同时保存,分别服务投影与文件扫描。名为 thread_id 的字段实际上指 rollout ID;普通记录二者相同,revert 后一个逻辑 Thread 可以对应新的物理 rollout ID,不能继续把它当成逻辑主键。
用于准备分叉的返回对象还保留了“直接来源”的身份。
源码文件:codex-rs/thread-store/src/types.rs
相关函数/类型:ForkBoundary / PreparedFork;行号:182-191,211-222。
/// Requested boundary for inheriting a paginated thread's history.
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum ForkBoundary {
/// Inherit the source thread's latest durable state.
Latest,
/// Inherit history through the newest visible occurrence of this turn.
ThroughTurn(String),
/// Inherit history preceding the original visible occurrence of this turn.
BeforeTurn(String),
}
// ...
/// Frozen source history and model context for a reference-backed fork.
#[derive(Debug)]
pub struct PreparedFork {
/// Immediate source thread, even when the normalized history base names an ancestor.
pub source_thread_id: ThreadId,
/// Frozen physical rollout prefix inherited by the child.
pub history_base: Option<HistoryPosition>,
/// Bounded model context selected by the requested fork boundary.
pub model_context: Arc<Vec<RolloutItem>>,
/// Blocks source deletion until the child's history reference is durable.
_source_reservation: Box<dyn std::fmt::Debug + Send>,
}source_thread_id 指本次 fork 的直接源 Thread;history_base.thread_id 可能指更早的祖先 rollout;model_context 只保存启动模型所需的上下文。_source_reservation 则是暂时防止来源被删除的资源,不是持久引用本身。这四个字段分别回答“从谁分叉”“继承哪一段”“模型读什么”“准备时谁不能被清理”。
例如从 B 分叉 C,但边界落在 B 继承的 A 前缀中,可以有如下关系。
图中直接来源仍是 B,不能为了让指针一致把 forked_from_id 改成 A;来源元数据继承与物理历史边界承担不同职责。parent_thread_id 还属于 agent 的所有权关系,也不能拿 fork lineage 代替它。
3.2 双位置计算
Latest 取当前投影已经处理到的 next_ordinal/next_byte_offset;包含具名轮次则取该轮次最后记录的位置。
源码文件:codex-rs/thread-store/src/local/paginated_fork.rs
相关函数/类型:history_base_at_boundary 的 ThroughTurn;行号:105-135。
let latest_position = HistoryPosition {
thread_id: source_segment.rollout_id(),
end_ordinal_exclusive: latest_projection_state.next_ordinal,
end_byte_offset: latest_projection_state.next_byte_offset,
};
let pool = store.thread_history_db().await?;
let position = match boundary {
ForkBoundary::Latest => latest_position,
ForkBoundary::ThroughTurn(turn_id) => {
let row = find_visible_turn(pool, lineage, turn_id.as_str()).await?;
if row.status == "inProgress" {
return Err(ThreadStoreError::InvalidRequest {
message: format!("lastTurnId '{turn_id}' identifies an in-progress turn"),
});
}
let rollout_end_ordinal = row
.rollout_end_ordinal
.ok_or_else(|| missing_turn_position(turn_id.as_str()))?;
let rollout_end_byte_offset = row
.rollout_end_byte_offset
.ok_or_else(|| missing_turn_position(turn_id.as_str()))?;
HistoryPosition {
thread_id: row.rollout_id,
end_ordinal_exclusive: u64::try_from(rollout_end_ordinal)
.map_err(|_| invalid_turn_position(turn_id.as_str()))?
.checked_add(1)
.ok_or_else(|| invalid_turn_position(turn_id.as_str()))?,
end_byte_offset: u64::try_from(rollout_end_byte_offset)
.map_err(|_| invalid_turn_position(turn_id.as_str()))?,
}
}包含边界将结束 ordinal 加一变成排他性上界,字节位置已经是记录结尾,不再加一。整数转换与 checked_add 防止非法负值和溢出。目标状态为 inProgress 时拒绝,与 Legacy 的 API 语义一致。
排除具名轮次使用开始位置。
源码文件:codex-rs/thread-store/src/local/paginated_fork.rs
相关函数/类型:history_base_at_boundary 的 BeforeTurn;行号:136-154。
ForkBoundary::BeforeTurn(turn_id) => {
let row = find_source_turn(pool, lineage, turn_id.as_str()).await?;
if row.rollout_end_ordinal == Some(row.rollout_ordinal) {
return Err(ThreadStoreError::InvalidRequest {
message: format!("turn {turn_id} does not have a persisted start boundary"),
});
}
let rollout_byte_offset = row
.rollout_byte_offset
.ok_or_else(|| missing_turn_position(turn_id.as_str()))?;
HistoryPosition {
thread_id: row.rollout_id,
end_ordinal_exclusive: u64::try_from(row.rollout_ordinal)
.map_err(|_| invalid_turn_position(turn_id.as_str()))?,
end_byte_offset: u64::try_from(rollout_byte_offset)
.map_err(|_| invalid_turn_position(turn_id.as_str()))?,
}
}
};BeforeTurn 不做终态要求,但必须存在可持久定位的开始边界。若 rollout_end_ordinal == rollout_ordinal,实现把它识别成没有可用持久开始边界的轮次并拒绝;缺少字节位置同样不能猜测偏移。
还有一个不应忽略的差别:两种边界查找祖先链的方向相反。
源码文件:codex-rs/thread-store/src/local/thread_history/turn_lookup.rs
相关函数/类型:find_source_turn / find_visible_turn;行号:21-35。
pub(in crate::local) async fn find_source_turn(
pool: &sqlx::SqlitePool,
lineage: &RolloutLineage,
turn_id: &str,
) -> ThreadStoreResult<TurnRow> {
find_turn(pool, lineage.segments().iter(), turn_id).await
}
pub(in crate::local) async fn find_visible_turn(
pool: &sqlx::SqlitePool,
lineage: &RolloutLineage,
turn_id: &str,
) -> ThreadStoreResult<TurnRow> {
find_turn(pool, lineage.segments().iter().rev(), turn_id).await
}find_source_turn 从祖先到当前段扫描,定位原始出现;find_visible_turn 从当前段反向扫描,取最新可见出现。同一个 Turn ID 可能在继承前缀和子段的状态覆盖中出现,包含它应保留最新可见终态,排除它则要回到最初开始处。若两者都从链尾搜索,BeforeTurn 可能只切掉子段补写的终态,把原 Turn 主体留在继承历史中。
3.3 祖先上界
找到 Turn 行还不够,位置必须落在来源实际继承的范围内。
源码文件:codex-rs/thread-store/src/local/paginated_fork.rs
相关函数/类型:history_base_at_boundary 的祖先上界;行号:155-178。
let segment_index = lineage
.segments()
.iter()
.position(|segment| segment.rollout_id() == position.thread_id)
.ok_or_else(|| ThreadStoreError::Internal {
message: "fork position is outside the source lineage".to_string(),
})?;
if lineage.segments()[segment_index].end.is_some_and(|end| {
position.end_ordinal_exclusive > end.end_ordinal_exclusive
|| position.end_byte_offset > end.end_byte_offset
}) {
return Err(ThreadStoreError::InvalidRequest {
message: "fork boundary exceeds inherited source history".to_string(),
});
}
let history_base =
if position.end_ordinal_exclusive == lineage.segments()[segment_index].start_ordinal() {
segment_index
.checked_sub(1)
.and_then(|index| lineage.segments()[index].end)
} else {
Some(position)
};
Ok(history_base)先找到位置所属的 segment,再同时检查 ordinal 与 byte offset 不超过该段的冻结上界。数据库中“存在某轮次”不等于当前 Thread 有权继承到它;来源早先截断过的祖先后缀不能在再次 fork 时重新冒出来。
如果切点恰好等于当前段的起点,就退回前一段的结束位置;已经没有前段时返回 None。None 在这里表示无继承前缀,不是没有源 Thread,更不是 prepare 失败。新 Thread 仍可使用来源元数据建立空前缀分叉。
3.4 模型消费者
固定引用范围之后,模型输入从被截断的 lineage 中提取。
源码文件:codex-rs/thread-store/src/local/model_context.rs
相关函数/类型:load_for_fork;行号:87-113。
pub(super) async fn load_for_fork(
lineage: RolloutLineage,
history_base: Option<HistoryPosition>,
) -> ThreadStoreResult<Vec<RolloutItem>> {
let source_path = lineage
.segments()
.last()
.map(|segment| segment.rollout_path.as_path())
.ok_or_else(|| ThreadStoreError::Internal {
message: "fork lineage has no source segment".to_string(),
})?;
let session_meta = codex_rollout::read_session_meta_line(source_path)
.await
.map_err(|err| ThreadStoreError::Internal {
message: format!(
"failed to read session metadata {}: {err}",
source_path.display()
),
})?;
match history_base {
Some(history_base) => {
let lineage = lineage.truncate_at(history_base).await?;
scan_model_context_from_lineage(lineage, session_meta).await
}
None => Ok(vec![RolloutItem::SessionMeta(session_meta)]),
}
}有 history base 时先 truncate_at,再扫描模型上下文;没有前缀时仍保留源的 SessionMeta。模型扫描可利用压缩产生的 replacement-history checkpoint,因而 model_context.len() 不能当成完整 UI 历史长度。完整展示历史则沿引用前缀与子段投影读取,正如读取文章中的 Paginated 路径。
4. 引用的存活期
4.1 准备与物化
准备分叉时,源的文件和祖先链不能被其他生命周期操作移走。实现先取得 lifecycle reservation,再把持久化和 lineage 解析放进一个独立 task。
源码文件:codex-rs/thread-store/src/local/paginated_fork.rs
相关函数/类型:prepare 的生命周期保留;行号:15-40。
pub(super) async fn prepare(
store: &LocalThreadStore,
params: PrepareForkParams,
) -> ThreadStoreResult<PreparedFork> {
let PrepareForkParams {
thread_id,
boundary,
} = params;
let source_reservation = store.live_writer_locks.reserve_lifecycle(thread_id).await;
// Keep the source reserved until persistence and lineage materialization finish, even if the
// caller cancels fork preparation.
let lineage_store = store.clone();
let (lineage, source_reservation) = tokio::spawn(async move {
match live_writer::persist_thread(&lineage_store, thread_id).await {
Ok(()) | Err(ThreadStoreError::ThreadNotFound { .. }) => {}
Err(err) => return Err(err),
}
let lineage = lineage_store
.resolve_rollout_lineage_for_reference(thread_id)
.await?;
Ok::<_, ThreadStoreError>((lineage, source_reservation))
})
.await
.map_err(|err| ThreadStoreError::Internal {
message: format!("failed to resolve fork lineage: {err}"),
})??;reservation 被移进 tokio::spawn。如果调用者取消等待,已经启动的任务仍持有它,直到持久化与 lineage 物化完成或出错。吞掉的只有这里预期可出现的 ThreadNotFound;其他持久化错误直接传播,join 失败再映射为内部错误。这不是无条件忽略 Store 失败。
接下来将需要的投影追平,在 writer 锁保护下捕获边界,再释放 writer 锁读取固定前缀。
源码文件:codex-rs/thread-store/src/local/paginated_fork.rs
相关函数/类型:prepare 的物化与锁释放;行号:41-84。
let source_segment = lineage
.segments()
.last()
.ok_or_else(|| ThreadStoreError::Internal {
message: "fork lineage has no source segment".to_string(),
})?;
if store.state_db.is_none() {
return Err(ThreadStoreError::Unsupported {
operation: "prepare_fork",
});
}
if !matches!(boundary, ForkBoundary::Latest) {
for segment in lineage
.segments()
.iter()
.take(lineage.segments().len().saturating_sub(1))
{
let _ancestor_writer_guard = store.live_writer_locks.lock(segment.rollout_id()).await;
super::thread_history_materialization::materialize_to_sqlite(
store,
segment.rollout_id(),
segment.rollout_path.as_path(),
)
.await?;
}
}
let source_writer_guard = store.live_writer_locks.lock(thread_id).await;
super::thread_history_materialization::materialize_to_sqlite(
store,
source_segment.rollout_id(),
source_segment.rollout_path.as_path(),
)
.await?;
let history_base = history_base_at_boundary(store, thread_id, boundary, &lineage).await?;
drop(source_writer_guard);
let model_context = Arc::new(model_context::load_for_fork(lineage, history_base).await?);
Ok(PreparedFork::new(
thread_id,
history_base,
model_context,
source_reservation,
))非 Latest 边界需要查祖先 Turn,因此先把祖先段物化到 SQLite;来源段也必须追平。state_db 缺失直接 Unsupported。关键时机是 drop(source_writer_guard):边界一旦固定,后面的模型扫描和新 Thread 初始化不必一直阻止源追加记录。
4.2 两种锁
“阻止删除”和“阻止继续写入”是两种不同需求,Store 用两把锁表达。
源码文件:codex-rs/thread-store/src/local/mod.rs
相关函数/类型:ThreadCoordination / reserve_lifecycle;行号:155-164,185-201。
#[derive(Default)]
struct ThreadCoordination {
// Serialize writes and capture consistent fork snapshots.
writer: Arc<Mutex<()>>,
// Forks hold a shared lease until their child reference is durable; deletion, archive, and
// unarchive require exclusive access. Keeping this separate from `writer` lets the source
// accept writes during child initialization, including MCP startup that can take 30 seconds.
// Operations that need both locks must acquire `lifecycle` before `writer`.
lifecycle: Arc<RwLock<()>>,
}
// ...
async fn reserve_lifecycle(&self, thread_id: ThreadId) -> OwnedRwLockReadGuard<()> {
self.coordination(thread_id)
.await
.lifecycle
.clone()
.read_owned()
.await
}
async fn lock_lifecycle(&self, thread_id: ThreadId) -> OwnedRwLockWriteGuard<()> {
self.coordination(thread_id)
.await
.lifecycle
.clone()
.write_owned()
.await
}writer mutex 保护写入与一致的分叉边界捕获;lifecycle 的共享读 guard 阻止删除、归档等需要独占访问的操作。需要两把锁时先 lifecycle 后 writer。源仍可在 child 初始化期间继续写,避免一次可能等待 MCP 初始化的 fork 长时间卡住源任务。
下面把返回对象与两类资源分开画出。
图中的 OwnedRwLockReadGuard 是 Tokio 的实际 guard 类型;HistoryPosition.thread_id 指向物理 rollout。不要把 PreparedFork 的值拷贝当成锁续期:其私有 reservation 必须随原对象保持存活。
4.3 取消与删除
测试不仅检查 fork 能成功,还故意让删除与准备过程竞争。
源码文件:codex-rs/thread-store/src/local/thread_history_materialization_tests.rs
相关函数/类型:prepared_fork_reserves_source_until_child_reference_is_durable;行号:1165-1194,1208-1216。
let mut delete = Box::pin(store.delete_thread(DeleteThreadParams {
thread_id: source_thread_id,
}));
tokio::select! {
biased;
result = &mut delete => {
panic!("source deletion completed before its child reference was durable: {result:?}")
}
_ = tokio::task::yield_now() => {}
}
tokio::time::timeout(
Duration::from_secs(10),
store.append_items(AppendThreadItemsParams {
thread_id: source_thread_id,
items: vec![
turn_started("later-turn"),
user_message("later source message"),
turn_completed("later-turn"),
],
}),
)
.await
.expect("source write should not wait behind deletion")
.expect("source writes remain available during fork preparation");
assert_eq!(prepared.history_base, Some(history_base));
assert!(!contains_user_message(
prepared.model_context.as_slice(),
"later source message"
));
// ...
let error = delete
.await
.expect_err("durable child reference protects its source");
assert!(
error
.to_string()
.contains("forked history still references")
);持有 prepared 时,删除不能完成;同时源 append 仍能在超时上限内完成,prepared 中的 boundary 和 model context 都没有包含后来新增的用户消息。子引用持久化后释放 prepared,删除请求被“仍被分叉历史引用”拒绝。保护机制由临时 guard 交接为持久引用关系,不能理解成 guard 一释放源就必然可以删除。
另一个测试持有祖先 writer 锁,等到 detached lineage task 开始物化来源后,取消外层 prepare。
源码文件:codex-rs/thread-store/src/local/thread_history_materialization_tests.rs
相关函数/类型:cancelled_fork_keeps_source_reserved_until_lineage_materialization_finishes;行号:1114-1136。
.expect("detached lineage task should materialize the source");
preparation.abort();
assert!(
preparation
.await
.expect_err("fork preparation should be cancelled")
.is_cancelled()
);
let mut delete = Box::pin(store.delete_thread(DeleteThreadParams {
thread_id: source_thread_id,
}));
tokio::select! {
biased;
result = &mut delete => {
panic!("source deletion completed while lineage materialization was active: {result:?}")
}
_ = tokio::task::yield_now() => {}
}
drop(ancestor_writer_guard);
delete
.await
.expect("delete source after lineage materialization finishes");取消完成后删除仍不能立即越过未完成的 lineage 任务;释放祖先锁后删除才继续。这里证明的是准备阶段的特定取消窗口,不是“任何 fork 阶段取消都会创建 child”,也不保证整个 API 跨所有 Store 都有一个原子事务。
5. 未完成轮次
5.1 快照终态
Latest 可以碰到半个 Turn。Core 的处理不是把源线程停下来,而是让新 Thread 的历史表示“继承到这里时已经中断”。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:fork_history_from_snapshot;行号:2229-2260。
fn fork_history_from_snapshot(
snapshot: ForkSnapshot,
history: InitialHistory,
interrupted_marker: InterruptedTurnHistoryMarker,
) -> InitialHistory {
let snapshot_state = snapshot_turn_state(&history);
match snapshot {
ForkSnapshot::TruncateBeforeNthUserMessage(nth_user_message) => {
truncate_before_nth_user_message(history, nth_user_message, &snapshot_state)
}
ForkSnapshot::Interrupted => {
let history = match history {
InitialHistory::New => InitialHistory::New,
InitialHistory::Cleared => InitialHistory::Cleared,
InitialHistory::Forked(history) => InitialHistory::Forked(history),
InitialHistory::Resumed(resumed) => {
InitialHistory::Forked(Arc::unwrap_or_clone(resumed.history))
}
};
if snapshot_state.ends_mid_turn {
append_interrupted_boundary(
history,
snapshot_state.active_turn_id,
snapshot_state.active_turn_started_at,
interrupted_marker,
)
} else {
history
}
}
}
}ForkSnapshot::Interrupted 先把 Resumed 历史转换成 Forked,必要时追加中断边界。App Server 的两条分叉入口都使用这种快照语义;Core 另外支持按第 N 个用户消息截断,属于通用 API,本篇的 lastTurnId 不是直接传给该序号参数。
如何判断半个 Turn,需要看聚合后的状态而非只搜索最后一条消息。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:snapshot_turn_state 的显式边界;行号:2172-2199。
fn snapshot_turn_state(history: &InitialHistory) -> SnapshotTurnState {
let rollout_items = history.get_rollout_items();
let mut builder = ThreadHistoryBuilder::new();
for item in rollout_items {
builder.handle_rollout_item(item);
}
let active_turn_id = builder.active_turn_id_if_explicit();
if builder.has_active_turn() && active_turn_id.is_some() {
let active_turn_snapshot = builder.active_turn_snapshot();
if active_turn_snapshot
.as_ref()
.is_some_and(|turn| turn.status != TurnStatus::InProgress)
{
return SnapshotTurnState {
ends_mid_turn: false,
active_turn_id: None,
active_turn_started_at: None,
active_turn_start_index: None,
};
}
return SnapshotTurnState {
ends_mid_turn: true,
active_turn_id,
active_turn_started_at: active_turn_snapshot.and_then(|turn| turn.started_at),
active_turn_start_index: builder.active_turn_start_index(),
};
}存在显式 active Turn 时,终态快照不再当作进行中;真正 InProgress 才保留原 turn ID、开始时间和索引。没有显式事件的旧记录还有单独回退:从最后一个真实用户消息往后查找完成或中断事件。这个回退不会凭空生成一个稳定的 Turn ID,相关细节可沿前置的 Core 分叉文章继续追踪。
5.2 中断标记
收口同时涉及供 UI 归约的事件和供模型理解的提示条目。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:append_interrupted_boundary;行号:2265-2304。
fn append_interrupted_boundary(
history: InitialHistory,
turn_id: Option<String>,
started_at: Option<i64>,
interrupted_marker: InterruptedTurnHistoryMarker,
) -> InitialHistory {
let aborted_event = RolloutItem::EventMsg(EventMsg::TurnAborted(TurnAbortedEvent {
turn_id,
reason: TurnAbortReason::Interrupted,
started_at,
completed_at: None,
duration_ms: None,
}));
match history {
InitialHistory::New | InitialHistory::Cleared => {
let mut history = Vec::new();
if let Some(marker) = interrupted_turn_history_marker(interrupted_marker) {
history.push(RolloutItem::ResponseItem(marker.into()));
}
history.push(aborted_event);
InitialHistory::Forked(history)
}
InitialHistory::Forked(mut history) => {
if let Some(marker) = interrupted_turn_history_marker(interrupted_marker) {
history.push(RolloutItem::ResponseItem(marker.into()));
}
history.push(aborted_event);
InitialHistory::Forked(history)
}
InitialHistory::Resumed(resumed) => {
let mut history = Arc::unwrap_or_clone(resumed.history);
if let Some(marker) = interrupted_turn_history_marker(interrupted_marker) {
history.push(RolloutItem::ResponseItem(marker.into()));
}
history.push(aborted_event);
InitialHistory::Forked(history)
}
}
}TurnAborted 保留可用的源 turn ID 与 started_at,但 completed_at/duration 不伪造运行测量。模型提示条目是否写入由 marker 配置决定,之后仍会写中断事件。因此“关闭中断提示文案”和“不产生中断终态”不是一回事。
源码文件:codex-rs/core/src/tasks/mod.rs
相关函数/类型:InterruptedTurnHistoryMarker::from_config_and_version;行号:86-98。
pub(crate) fn from_config_and_version(
config: &Config,
multi_agent_version: MultiAgentVersion,
) -> Self {
if !config.agent_interrupt_message_enabled {
return Self::Disabled;
}
if multi_agent_version == MultiAgentVersion::V2 {
Self::Developer
} else {
Self::ContextualUser
}
}禁用时不添加模型提示;V2 multi-agent 使用 developer 角色,否则使用 contextual user 形式。它是模型上下文的角色差异,不能仅在 UI 的 Turn.status 中观察到,也不能把某一角色硬写成所有 fork 的行为。
5.3 源与分支隔离
集成测试先构造带 active-turn 和 before-fork 内容的来源。显式 ThroughTurn 被拒绝,默认 Latest 则成功建立分支。
源码文件:codex-rs/app-server/tests/suite/v2/thread_fork.rs
相关函数/类型:assert_thread_fork_freezes_active_paginated_turn_as_interrupted;行号:1602-1617,1638-1659。
let invalid_fork_id = mcp
.send_thread_fork_request(ThreadForkParams {
thread_id: source_thread_id.clone(),
last_turn_id: Some("active-turn".to_string()),
..Default::default()
})
.await?;
let invalid_fork = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_error_message(RequestId::Integer(invalid_fork_id)),
)
.await??;
assert_eq!(
invalid_fork.error.message,
"lastTurnId 'active-turn' identifies an in-progress turn"
);
// ...
assert!(matches!(
child_rollout.as_slice(),
[
RolloutLine { item: RolloutItem::SessionMeta(_), .. },
RolloutLine {
item: RolloutItem::EventMsg(EventMsg::ThreadSettingsApplied(_)),
..
},
RolloutLine {
item: RolloutItem::ResponseItem(response_item),
..
},
RolloutLine {
item: RolloutItem::EventMsg(EventMsg::TurnAborted(aborted)),
..
},
] if matches!(
&response_item.item,
codex_protocol::models::ResponseItem::Message { role, .. }
if role == expected_marker_role
) && aborted.turn_id.as_deref() == Some("active-turn")
));Paginated 子文件只含自己的 SessionMeta、实际设置、新增中断提示和 TurnAborted,旧正文留在 history base。测试随后继续向源追加 after-fork,在分支下一轮模型请求中检查 before-fork 存在、after-fork 不存在,并核对中断提示角色。这同时约束固定前缀和模型消费者,不能只断言 forked_thread.status == Idle。
这类快照没有把源正在执行的 shell、MCP 请求或审批 callback 复制到 child。旧工具在原执行者那里继续;child 的历史中断边界防止模型把继承来的半个任务误当作仍有一个本地工具等待返回。
6. 配置的时间点
6.1 当前设置与旧历史
“从旧 Turn 分叉”不等于“回到那一刻的权限”。App Server 为未显式覆盖的设置寻找源线程的最新值。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的设置来源判定;行号:4759-4809。
/*personality*/ None,
);
typesafe_overrides.ephemeral = ephemeral.then_some(true);
let restore_approval_policy = typesafe_overrides.approval_policy.is_none();
let restore_approvals_reviewer = typesafe_overrides.approvals_reviewer.is_none()
&& !request_overrides
.as_ref()
.is_some_and(|overrides| overrides.contains_key("approvals_reviewer"));
let restore_permission_profile =
!has_permission_override(request_overrides.as_ref(), &typesafe_overrides);
let needs_latest_settings =
restore_approval_policy || restore_approvals_reviewer || restore_permission_profile;
let loaded_parent_settings = if paginated_source && needs_latest_settings {
if let Ok(parent) = self.thread_manager.get_thread(source_thread_id).await {
let snapshot = parent.thread_settings_snapshot().await;
Some(PersistedResumeSettings {
approval_policy: snapshot.approval_policy,
approvals_reviewer: Some(snapshot.approvals_reviewer),
active_permission_profile: snapshot.active_permission_profile,
})
} else {
None
}
} else {
None
};
let latest_context = if paginated_source
&& needs_latest_settings
&& loaded_parent_settings.is_none()
&& (last_turn_id.is_some() || before_turn_id.is_some())
{
Some(
self.thread_store
.load_latest_model_context(StoreLoadThreadHistoryParams {
thread_id: source_thread_id,
include_archived: true,
})
.await
.map_err(thread_store_resume_read_error)?
.items,
)
} else {
None
};
let persisted_settings = loaded_parent_settings.or_else(|| {
latest_persisted_resume_settings(
latest_context
.as_deref()
.unwrap_or_else(|| source_history_items.as_ref()),
)
});先分别判断 approval policy、reviewer、permission profile 是否需要继承。Paginated 来源已加载时,直接取 live thread_settings_snapshot;若未加载、请求具名旧边界且仍需要设置,就另外加载最新模型上下文。这样被截断的模型历史不会迫使新 Thread 恢复过时的安全配置。
Legacy 也不是先截断再找设置:它先从完整 source_history_items 读取最近设置,之后才执行具名历史截断。因此两个存储模式都不能简单概括成“所有字段随选定历史一起回滚”。不过两者是否取 live snapshot、是否多读一次最新上下文,仍有明确分支差异。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的设置应用;行号:4810-4829。
if let Some(persisted_settings) = persisted_settings {
if restore_approval_policy {
typesafe_overrides.approval_policy = Some(persisted_settings.approval_policy);
}
if restore_approvals_reviewer {
typesafe_overrides.approvals_reviewer = persisted_settings.approvals_reviewer;
}
if restore_permission_profile {
typesafe_overrides.persisted_permission_profile_id = persisted_settings
.active_permission_profile
.map(|profile| profile.id);
}
}
// Derive a Config using the same logic as new conversation, honoring overrides if provided.
let config = self
.config_manager
.load_for_cwd(request_overrides, typesafe_overrides, history_cwd)
.await
.map_err(|err| config_load_error(&err))?;
let goals_enabled = config.features.enabled(Feature::Goals);显式覆盖存在时,不用持久值替换它;继承权限时记录的是 profile ID,交给当前 Config 加载器重新解释。最后的 load_for_cwd 以源摘要 cwd 作为恢复背景,继续应用请求覆盖与配置文件规则。命名 profile 的内容可能已改变,得到的 sandbox 要以最终响应为准。
图中两条时间线分别决定“模型继承哪些记录”和“新执行者如何运行”。
这里没有复制 resume 的全部模型持久元数据合并流程。fork 的入口构造自己的 Config,读者应按该函数的实际调用判断哪些字段继承,不能把上一篇中 merge_persisted_resume_metadata 的行为自动套过来。thread_fork_preserves_persisted_permission_profile_and_honors_overrides 会分别检查继承 profile、显式 sandbox 与显式 permission profile;reviewer 的配对测试同时覆盖 Legacy/Paginated,以及旧边界分叉仍采用最新设置。
6.2 平台覆盖
Windows 在构造请求覆盖时还有独立分支。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的 Windows 配置继承;行号:4727-4748。
match WindowsSandboxLevel::from_config(&self.config) {
WindowsSandboxLevel::Elevated => {
cli_overrides
.insert("windows.sandbox".to_string(), serde_json::json!("elevated"));
}
WindowsSandboxLevel::RestrictedToken => {
cli_overrides.insert(
"windows.sandbox".to_string(),
serde_json::json!("unelevated"),
);
}
WindowsSandboxLevel::Disabled => {}
}
}
let request_overrides = if cli_overrides.is_empty() {
None
} else {
Some(cli_overrides)
};
let runtime_workspace_roots = runtime_workspace_roots.map(resolve_runtime_workspace_roots);
let mut typesafe_overrides = self.build_thread_config_overrides(
model,服务器当前 WindowsSandboxLevel 为 Elevated 或 RestrictedToken 时,处理器写入 windows.sandbox 对应字符串;Disabled 不插入。它发生在 Config 加载之前,不能描述为启动完成后才临时改变 OS token。其他平台跳过这段 cfg!(windows);历史边界与引用算法本身没有因此变成另一套。
本机执行某个跨平台测试通过,并不能证明 Windows 沙箱实际建立了相应权限边界。阅读本篇时可以定位配置来源与生效阶段;具体 OS 强制效果需要对应平台的执行与安全专题。
7. 新身份与落盘
7.1 项目归属预置
持久分支继承源的 project ID。这个字段不能等创建完成后才随意补进 UI,否则创建时消费持久元数据的路径可能看见不完整归属。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的 project 元数据预置;行号:4871-4883。
let inherited_project_id = source_thread.project_id.clone();
let reserved_thread_id = if config.ephemeral {
None
} else {
stage_pending_project_metadata(
self.thread_manager.as_ref(),
self.thread_store.as_ref(),
inherited_project_id.as_deref(),
"thread/fork",
)
.await?
};只有持久配置才尝试 stage metadata。helper 在有 project ID 时预留新 Thread ID,并将 metadata patch 写入 Store。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:stage_pending_project_metadata;行号:31-58。
async fn stage_pending_project_metadata(
thread_manager: &ThreadManager,
thread_store: &dyn ThreadStore,
project_id: Option<&str>,
operation: &'static str,
) -> Result<Option<ThreadId>, JSONRPCErrorError> {
let Some(project_id) = project_id else {
return Ok(None);
};
let thread_id = thread_manager.reserve_thread_id();
thread_store
.stage_pending_thread_metadata(
thread_id,
StoreThreadMetadataPatch {
project_id: Some(Some(project_id.to_string())),
..Default::default()
},
)
.await
.map_err(|error| match error {
ThreadStoreError::Unsupported { .. } => {
method_not_found(format!("{operation} is unavailable without sqlite state"))
}
ThreadStoreError::InvalidRequest { message } => invalid_request(message),
error => internal_error(format!("failed to stage {operation} metadata: {error}")),
})?;
Ok(Some(thread_id))
}预留身份与创建调用使用同一个 ID;没有 project 则无需预留。后端不支持这项操作时明确返回方法不可用,不能创建一个丢了 project 的“部分成功”响应。这个 pending metadata 属于创建准备阶段,与最终 Thread 的持久记录有不同清理时机。
7.2 两条 Core 入口
App Server 根据 prepared fork 是否存在选择 Core 调用。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的 Core 创建与错误清理;行号:4884-4931。
let new_thread = if let Some(prepared_fork) = prepared_fork {
self.thread_manager
.fork_prepared_thread(
config,
prepared_fork,
thread_source,
parent_trace,
client_mcp_extensions,
reserved_thread_id,
)
.await
} else {
self.thread_manager
.fork_thread_from_history(
ForkSnapshot::Interrupted,
config,
InitialHistory::Resumed(ResumedHistory {
conversation_id: source_thread_id,
history: history_items,
rollout_path: source_thread.rollout_path.clone(),
}),
thread_source,
parent_trace,
client_mcp_extensions,
reserved_thread_id,
)
.await
};
let NewThread {
thread_id,
thread: forked_thread,
session_configured,
..
} = match new_thread {
Ok(new_thread) => new_thread,
Err(err) => {
remove_pending_project_metadata(self.thread_store.as_ref(), reserved_thread_id)
.await;
return Err(match err.details() {
CodexErrorDetails::Io(_) | CodexErrorDetails::Json(_) => {
invalid_request(format!("failed to load thread {source_thread_id}: {err}"))
}
CodexErrorDetails::InvalidRequest(message) => invalid_request(message.clone()),
_ => internal_error(format!("error forking thread: {err}")),
});
}
};Paginated 将 PreparedFork 整体交给 manager;Legacy 则把截取后的历史以 Resumed 形态携带源 ID,再指定 ForkSnapshot::Interrupted。后一形态只是为了传递来源身份,并不会让 child 复用源 Thread ID。
Core 失败时先清理预置的 project metadata,再区分 I/O/JSON、InvalidRequest 与其他错误;它不是用统一内部错误隐藏全部原因。清理 helper 自身失败只记 warning,因此还要把业务错误和清理错误分别看待。
Legacy 选择复制持久化。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:fork_thread_from_history;行号:1255-1281。
pub async fn fork_thread_from_history<S>(
&self,
snapshot: S,
config: Config,
history: InitialHistory,
thread_source: Option<ThreadSource>,
parent_trace: Option<W3cTraceContext>,
client_mcp_extensions: ClientMcpExtensions,
reserved_thread_id: Option<ThreadId>,
) -> CodexResult<NewThread>
where
S: Into<ForkSnapshot>,
{
self.fork_thread_with_initial_history(
config,
ForkHistory {
snapshot: snapshot.into(),
initial_history: history,
persistence: ForkPersistence::Copied,
},
thread_source,
parent_trace,
client_mcp_extensions,
reserved_thread_id,
)
.await
}ForkPersistence::Copied 是这条入口的明确选择,不是根据文件是否存在临时推断。来源历史随后会转为 Forked,交给新 Session 建立自己的记录。
Paginated 则把引用位置和继承 item 数一起传入。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:fork_prepared_thread;行号:1284-1318。
pub async fn fork_prepared_thread(
&self,
config: Config,
prepared: PreparedFork,
thread_source: Option<ThreadSource>,
parent_trace: Option<W3cTraceContext>,
client_mcp_extensions: ClientMcpExtensions,
reserved_thread_id: Option<ThreadId>,
) -> CodexResult<NewThread> {
let history = InitialHistory::Resumed(ResumedHistory {
conversation_id: prepared.source_thread_id,
history: Arc::clone(&prepared.model_context),
rollout_path: None,
});
let fork_persistence = ForkPersistence::Referenced {
history_base: prepared.history_base,
inherited_item_count: prepared.model_context.len(),
};
let result = self
.fork_thread_with_initial_history(
config,
ForkHistory {
snapshot: ForkSnapshot::Interrupted,
initial_history: history,
persistence: fork_persistence,
},
thread_source,
parent_trace,
client_mcp_extensions,
reserved_thread_id,
)
.await;
drop(prepared);
result
}inherited_item_count 用于区分模型重建所需的旧上下文与 child 应真正追加到本地文件的新记录。drop(prepared) 位于 await 返回之后,使 reservation 覆盖正常的 child 初始化过程;它不是在函数一进入就丢掉来源保护。
共同路径记录直接来源,再生成新执行者。
源码文件:codex-rs/core/src/thread_manager.rs
相关函数/类型:fork_thread_with_initial_history;行号:1320-1368。
async fn fork_thread_with_initial_history(
&self,
config: Config,
fork_history: ForkHistory,
thread_source: Option<ThreadSource>,
parent_trace: Option<W3cTraceContext>,
client_mcp_extensions: ClientMcpExtensions,
reserved_thread_id: Option<ThreadId>,
) -> CodexResult<NewThread> {
let ForkHistory {
snapshot,
initial_history: history,
persistence: fork_persistence,
} = fork_history;
// `forked_from_id()` describes this history's existing lineage. When
// forking a resumed thread, the child copies the resumed thread itself.
let source_thread_id = match &history {
InitialHistory::Resumed(resumed) => Some(resumed.conversation_id),
InitialHistory::Forked(_) => history.forked_from_id(),
InitialHistory::New | InitialHistory::Cleared => None,
};
let multi_agent_version = self
.state
.effective_multi_agent_version_for_spawn(
&history,
/*session_source*/ None,
/*parent_thread_id*/ None,
source_thread_id,
&config,
)
.await;
let interrupted_marker =
InterruptedTurnHistoryMarker::from_config_and_version(&config, multi_agent_version);
let history = fork_history_from_snapshot(snapshot, history, interrupted_marker);
let agent_control = self.agent_control_for_config(&config);
let options = StartThreadOptions {
initial_history: history,
thread_source,
parent_trace,
client_mcp_extensions,
reserved_thread_id,
..StartThreadOptions::new(config)
};
let mut request =
ThreadSpawnRequest::new(options, Arc::clone(&self.state.auth_manager), agent_control);
request.forked_from_thread_id = source_thread_id;
request.fork_persistence = fork_persistence;
Box::pin(self.state.spawn_thread(request)).await
}Resumed 的 conversation_id 在此被当作 immediate source;不能使用它过去的 forked_from_id(),否则从 B 再 fork 会错误地把 C 的直接来源写成 A。request.forked_from_thread_id 与 parent_thread_id 分离,thread_source 也由本次公开请求决定。初始历史变成 Forked 后,正常创建管线使用新 ID 或预留 ID,完成 Session 装配与 manager 注册。
7.3 本地增量
新 Session 首先重建继承上下文,随后决定哪些记录要真正写入 child 的 rollout。
源码文件:codex-rs/core/src/session/mod.rs
相关函数/类型:Session::record_initial_history 的 Forked 持久化;行号:1423-1461。
let thread_settings_applied =
RolloutItem::EventMsg(thread_settings::applied_event(self).await);
match &self.fork_persistence {
ForkPersistence::Referenced {
inherited_item_count,
..
} => {
// Ancestor records remain behind history_base; only effective child
// settings and boundaries synthesized by snapshot processing are local.
rollout_items.drain(..*inherited_item_count);
rollout_items.insert(0, thread_settings_applied);
}
ForkPersistence::Copied if is_paginated_subagent => {
// Paginated subagents already persist inherited context when their live
// thread is created.
rollout_items.clear();
rollout_items.push(thread_settings_applied);
}
ForkPersistence::Copied => {
// Keep the copied prefix and effective child settings in one append so a
// cold resume cannot observe inherited settings as the latest value.
rollout_items.push(thread_settings_applied);
}
}
self.persist_rollout_items(&rollout_items).await;
// Forked threads should remain file-backed immediately after startup.
self.ensure_rollout_materialized(PersistContext::Standard)
.await;
// Flush after seeding history and any persisted rollout copy.
if !is_subagent {
let _ = self.flush_rollout().await;
}
Some(turn_context)
}
};
if let Some(turn_context) = turn_context
&& turn_context.config.memories.disable_on_external_contextReferenced 会从待写数组中 drain 掉继承的 inherited_item_count 个旧条目,只在本地留下实际 child 设置以及快照算法新增的中断边界。Copied 则保留前缀,并把 child 设置放在其后;这样重新冷恢复 child 时,不会把复制前缀里的源设置误认成最后生效值。
持久 fork 会立即物化自己的记录。Referenced 在 Session 启动返回前还有一处可失败的严格持久化调用。
源码文件:codex-rs/core/src/session/session.rs
相关函数/类型:Session::new 的引用持久化屏障;行号:1595-1599。
if matches!(&sess.fork_persistence, ForkPersistence::Referenced { .. }) {
// Keep the source reserved until the child's history reference is durable.
sess.try_ensure_rollout_materialized(PersistContext::Standard)
.await?;
}前面一些 Session 记录路径以内部方式处理持久化错误,这里的 try_ensure_rollout_materialized(...).await? 则让失败向创建调用传播。它服务于“正常创建完成之前,child 引用已持久化”的不变量。不要把 ensure 与 try_ensure 的错误传播方式看成一样。
以下时序图只表示正常成功路径;取消时源码已明确保护的范围,要结合上一节测试判断,不能据图假定所有阶段都自动回滚。
测试 thread_fork_creates_reference_backed_paginated_thread 检查 child 文件不含原用户正文,却通过 history base 和下一次模型请求保留该内容。
源码文件:codex-rs/app-server/tests/suite/v2/thread_fork.rs
相关函数/类型:thread_fork_creates_reference_backed_paginated_thread;行号:1410-1422,1439-1455。
assert_eq!(forked_thread.forked_from_id, Some(conversation_id.clone()));
assert_eq!(forked_thread.turns.len(), 1);
let forked_thread_id = forked_thread.id.clone();
let forked_path = forked_thread.path.expect("forked rollout path");
assert!(!std::fs::read_to_string(forked_path.as_path())?.contains("Saved user message"));
let meta = read_session_meta_line(forked_path.as_path()).await?;
let history_base = meta.meta.history_base.expect("history base");
assert_eq!(
history_base.thread_id,
ThreadId::from_string(conversation_id.as_str())?
);
let turn_id = mcp
// ...
let response_request = requests
.iter()
.find(|request| request.url.path().ends_with("/responses"))
.expect("forked turn response request");
let request_body = response_request.body_json::<Value>()?;
let model_input = request_body["input"]
.as_array()
.expect("response input array");
let model_input = serde_json::to_string(model_input)?;
assert!(model_input.contains("Saved user message"));
assert!(model_input.contains("Continue from the fork"));
// excludeTurns only controls response hydration; it must not change the inherited prefix.
let exclude_id = mcp
.send_thread_fork_request(ThreadForkParams {
thread_id: conversation_id,
exclude_turns: true,前一组断言验证存储形态,后一组检查 /responses 的真实 input:旧消息与新输入都必须存在。它比仅检查 UI turns 更有力度,因为模型读取路径和界面投影可能各自出错。该测试还比较 excludeTurns 分支的 history base,验证减少响应内容不会缩短继承范围。
8. 目标继承
8.1 当前快照
goal 不是每次 fork 都复制。入口要求显式 deferGoalContinuation,新 Thread 必须有持久路径,最终 Config 还需启用 Goals。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的 goal 继承;行号:4953-4980。
let inherited_goal = if defer_goal_continuation
&& session_configured.rollout_path.is_some()
&& goals_enabled
{
if let Some(state_db) = forked_thread.state_db().or_else(|| self.state_db.clone()) {
self.thread_goal_processor
.flush_goal_progress_for_fork(source_thread_id)
.await
.map_err(|err| {
internal_error(format!("failed to flush source thread goal: {err}"))
})?;
inherit_thread_goal_snapshot(&state_db, source_thread_id, thread_id)
.await
.map_err(|err| {
internal_error(format!("failed to inherit source thread goal: {err}"))
})?
} else {
false
}
} else {
false
};
if inherited_goal {
self.thread_goal_processor
.restore_inherited_goal_runtime(thread_id)
.await;
}先刷新源 goal 的使用量,再把快照复制给 child;源没有 goal、状态数据库不可用或开关不满足时不会建立继承目标。刷新或写入失败会传播错误,恢复 goal runtime 的辅助调用则有自己的 warning 路径。
源码文件:codex-rs/app-server/src/request_processors/thread_fork_goal.rs
相关函数/类型:inherit_thread_goal_snapshot;行号:5-28。
pub(super) async fn inherit_thread_goal_snapshot(
state_db: &StateRuntime,
source_thread_id: ThreadId,
target_thread_id: ThreadId,
) -> anyhow::Result<bool> {
let Some(mut goal) = state_db
.thread_goals()
.get_thread_goal(source_thread_id)
.await?
else {
return Ok(false);
};
if let Err(err) = validate_thread_goal_objective(&goal.objective) {
tracing::warn!(%source_thread_id, "skipping invalid inherited thread goal: {err}");
return Ok(false);
}
goal.thread_id = target_thread_id;
state_db
.thread_goals()
.replace_thread_goal_snapshot(&goal)
.await?;
Ok(true)
}复制前验证 objective;无效目标跳过。成功时只修改 goal.thread_id,其余字段沿用快照,所以它不是“同样 objective 的全新预算”。goal ID、status、token budget、已用 token/时间和原始时间戳都可能被保留。对话历史即使截成空前缀,目标仍是 fork 当下的目标状态,而非那个旧历史位置的状态。
8.2 同一事务
复制快照和阻止首次自动继续必须一起持久化。
源码文件:codex-rs/state/src/runtime/goals.rs
相关函数/类型:replace_thread_goal_snapshot;行号:68-123。
pub async fn replace_thread_goal_snapshot(
&self,
goal: &crate::ThreadGoal,
) -> anyhow::Result<()> {
let mut transaction = self.pool.begin().await?;
sqlx::query(
r#"
INSERT INTO thread_goals (
thread_id,
goal_id,
objective,
status,
token_budget,
tokens_used,
time_used_seconds,
created_at_ms,
updated_at_ms
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(thread_id) DO UPDATE SET
goal_id = excluded.goal_id,
objective = excluded.objective,
status = excluded.status,
token_budget = excluded.token_budget,
tokens_used = excluded.tokens_used,
time_used_seconds = excluded.time_used_seconds,
created_at_ms = excluded.created_at_ms,
updated_at_ms = excluded.updated_at_ms
"#,
)
.bind(goal.thread_id.to_string())
.bind(&goal.goal_id)
.bind(&goal.objective)
.bind(goal.status.as_str())
.bind(goal.token_budget)
.bind(goal.tokens_used)
.bind(goal.time_used_seconds)
.bind(datetime_to_epoch_millis(goal.created_at))
.bind(datetime_to_epoch_millis(goal.updated_at))
.execute(&mut *transaction)
.await?;
sqlx::query(
r#"
INSERT INTO thread_goal_continuation_deferrals (thread_id)
VALUES (?)
ON CONFLICT(thread_id) DO NOTHING
"#,
)
.bind(goal.thread_id.to_string())
.execute(&mut *transaction)
.await?;
transaction.commit().await?;
Ok(())
}事务首先 upsert thread_goals,再插入 thread_goal_continuation_deferrals,最后 commit。如果只存 goal 而没有 defer 标记,child 被 idle 生命周期观察到时就可能立即开始下一轮;如果只存 defer 而没有目标,又会留下无意义阻塞。表中的目标预算与使用量直接绑定原快照,不能在 fork 时悄悄归零。
idle 消费者在开始自动继续之前读取标记。
源码文件:codex-rs/ext/goal/src/runtime.rs
相关函数/类型:GoalRuntime::continue_if_idle 的延后检查;行号:361-380。
pub(crate) async fn continue_if_idle(&self) -> Result<(), String> {
if !self.tools_visible() {
self.inner.accounting_state.clear_active_goal();
return Ok(());
}
// Hold this through the read/start window so external set/clear cannot
// change the goal after we read it but before the continuation launches.
let _goal_state_permit = self.goal_state_permit().await?;
if self
.inner
.state_dbs
.thread_goals()
.has_thread_goal_continuation_deferral(self.thread_id())
.await
.map_err(|err| err.to_string())?
{
return Ok(());
}标记存在就返回,因此跨进程重新 resume 也能保持延后状态。它不是仅在当前 Rust 对象上设置一个 bool。开始新的显式 Turn 时,goal 扩展清除这个标记,再启动该轮计费。
源码文件:codex-rs/ext/goal/src/extension.rs
相关函数/类型:GoalExtension::on_turn_start 的延后清除;行号:196-213。
fn on_turn_start<'a>(&'a self, input: TurnStartInput<'a>) -> ExtensionFuture<'a, ()> {
Box::pin(async move {
let Some(runtime) = goal_runtime_handle(input.thread_store) else {
return;
};
if !runtime.is_enabled() {
return;
}
if let Err(err) = self
.state_dbs
.thread_goals()
.clear_thread_goal_continuation_deferral(runtime.thread_id())
.await
{
tracing::warn!("failed to clear deferred goal continuation: {err}");
}清除失败被记录为 warning,不能据一条普通 Turn 开始通知就断言数据库写入一定成功。这里展示了完整消费链:fork 写入 defer,idle 检查 defer,下一次 Turn 消费 defer。
8.3 预算连续性
集成测试用预算 150、已用 37 token 和 11 秒的 Active goal,分别创建完整历史、单轮前缀和空前缀分支。
源码文件:codex-rs/app-server/tests/suite/v2/thread_fork.rs
相关函数/类型:thread_fork_defers_inherited_active_goal_until_next_turn;行号:862-884,893-912,948-956。
} = timeout(DEFAULT_READ_TIMEOUT, mcp.read_response(fork_id)).await??;
let forked_thread_id = ThreadId::from_string(&forked_thread.id)?;
assert_eq!(forked_thread.turns.len(), expected_turn_count);
let mut expected_goal = source_goal.clone();
expected_goal.thread_id = forked_thread_id;
assert_eq!(
state_db
.thread_goals()
.get_thread_goal(forked_thread_id)
.await?,
Some(expected_goal)
);
assert!(
state_db
.thread_goals()
.has_thread_goal_continuation_deferral(forked_thread_id)
.await?
);
forked_threads.push(forked_thread);
}
assert_eq!(
state_db
// ...
.any(|method| method == "turn/started"),
"deferred goal should not start a turn while forking"
);
assert_eq!(
server
.received_requests()
.await
.expect("wiremock requests")
.iter()
.filter(|request| request.url.path().ends_with("/responses"))
.count(),
2,
"deferred goal should not issue a model request while forking"
);
let forked_thread = forked_threads.pop().expect("empty-prefix fork");
let forked_thread_id = ThreadId::from_string(&forked_thread.id)?;
drop(mcp);
let mut mcp = TestAppServer::builder()
.with_codex_home(codex_home.path())
// ...
.thread_goals()
.has_thread_goal_continuation_deferral(forked_thread_id)
.await?,
"first explicit turn should consume the deferred-goal marker"
);
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/started"),
)它要求快照除 thread ID 外相同,fork 不触发新的 turn/started 或模型请求,重启服务后仍保持延后。显式执行下一轮后,defer 标记消失,正常自动继续恢复;测试最终检查 child 用量达到 157 并进入 BudgetLimited,而 source goal 保持原值。这个例子说明历史前缀长度与剩余预算之间没有“按比例重算”的关系。
9. 结果与再执行
9.1 临时分支
临时 fork 不建立持久 writer,却仍需要模型可用的继承历史和恰当的响应摘要。处理器在历史被 Core 消费之前准备展示字段。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的临时历史视图;行号:4854-4870。
let ephemeral_preview = if ephemeral {
if paginated_source && last_turn_id.is_none() && before_turn_id.is_none() {
source_thread.preview.clone()
} else {
preview_from_rollout_items(&history_items)
}
} else {
String::new()
};
let ephemeral_turns = if ephemeral && include_turns {
build_legacy_api_turns_from_rollout_items(&history_items)
} else {
Vec::new()
};
let ephemeral_token_usage_turn_id = (ephemeral && include_turns)
.then(|| restored_token_usage_turn_id(&history_items, ephemeral_turns.as_slice()));
let token_usage_history_items = paginated_source.then(|| Arc::clone(&history_items));Legacy 可以提前归约 ephemeral turns;Paginated 临时 fork 强制 excludeTurns,因此不在这里装配完整展示数组。Paginated Latest 的 preview 可用源摘要,具名截断则从被选中的历史推导 preview,避免展示已排除的后续内容。token usage 的回放也随是否提供展示 turns 决定,不能为了显示状态行再次偷偷做全量历史读取。
无 rollout path 时,响应直接从配置快照构造。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的无路径响应;行号:5027-5039。
} else {
let mut thread = build_thread_from_snapshot(
thread_id,
session_configured.session_id.to_string(),
forked_thread.multi_agent_version(),
&config_snapshot,
/*path*/ None,
);
thread.preview = ephemeral_preview;
thread.forked_from_id = Some(source_thread_id.to_string());
thread.turns = ephemeral_turns;
(thread, ephemeral_token_usage_turn_id)
};forked_from_id、preview 和已准备的 turns 明确赋入;project 等摘要字段随后再补齐。pathless 表示没有自己的持久记录,不表示这是一个无效 Thread,也不表示它没有模型上下文。
源码文件:codex-rs/app-server/tests/suite/v2/thread_fork.rs
相关函数/类型:assert_thread_fork_ephemeral_remains_pathless_and_omits_listing;行号:2092-2104,2180-2190。
assert!(
thread.ephemeral,
"ephemeral forks should be marked explicitly"
);
assert_eq!(
thread.path, None,
"ephemeral forks should not expose a path"
);
assert_eq!(thread.preview, preview);
assert_eq!(thread.status, ThreadStatus::Idle);
assert_eq!(thread.name, None);
if history_mode == ThreadHistoryMode::Paginated {
assert!(thread.turns.is_empty());
// ...
assert!(
data.iter().all(|candidate| candidate.id != fork_thread_id),
"ephemeral forks should not appear in thread/list"
);
assert!(
data.iter().any(|candidate| candidate.id == conversation_id),
"persistent source thread should remain listed"
);
let turn_id = mcp
.send_turn_start_request(TurnStartParams {Legacy/Paginated 两个测试都检查 ephemeral=true、path=None、Idle,并要求它不出现在 thread/list,源仍在列表中。之后还向临时 Thread 发起真实 Turn,验证它能继续执行。这个结果不能外推成进程退出后仍可按 ID 冷恢复它;那是持久化承诺,和进程内可执行性不同。
9.2 名称与通知
名称继承发生在 Core 创建成功之后;有持久路径时更新新 Thread 的 metadata。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的名称继承;行号:4938-4952。
if session_configured.rollout_path.is_some()
&& let Some(name) = source_thread_name.clone()
{
self.thread_manager
.update_thread_metadata(
thread_id,
StoreThreadMetadataPatch {
name: Some(Some(name)),
..Default::default()
},
/*include_archived*/ true,
)
.await
.map_err(|err| core_thread_write_error("inherit source thread name", err))?;
}名称来源先经过 normalize,只有非空合法结果才复制。持久写入失败在这里会使 fork RPC 失败,而此时 Core 新对象已经存在;客户端不能把所有错误理解成“从未创建过分支”。这与前面的 Core 创建失败时清理 pending metadata 是不同阶段。
watch 状态采用 silent upsert。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的静默状态装配;行号:5062-5077。
self.thread_watch_manager
.upsert_thread_silently(&thread.id)
.await;
thread.status = resolve_thread_status(
self.thread_watch_manager
.loaded_status_for_thread(&thread.id)
.await,
/*has_in_progress_turn*/ false,
);
let sandbox = config_snapshot.sandbox_policy().into();
let active_permission_profile =
thread_response_active_permission_profile(config_snapshot.active_permission_profile);
let thread_originator = config_snapshot.originator.clone();
let response = ThreadForkResponse {
thread: thread.clone(),这让 fork 通过 thread/started 正式介绍新对象,避免先冒出一个客户端还不认识的 ID 的 thread/status/changed。Thread 状态取 watch 结果且不假装有正在执行的新 Turn;继承历史里的 Interrupted 是旧轮次的状态,两者可以同时存在。
响应、usage、started 与 goal 快照按下面的顺序安排。
源码文件:codex-rs/app-server/src/request_processors/thread_processor.rs
相关函数/类型:thread_fork_inner 的响应与通知顺序;行号:5093-5120。
let connection_id = request_id.connection_id;
self.outgoing
.send_response_with_thread_originator(request_id, response, thread_originator)
.await;
// `excludeTurns` is the cheap fork path, so skip restored usage replay
// instead of rebuilding history only to attribute a historical update.
if let Some(token_usage_turn_id) = token_usage_turn_id {
// Mirror the resume contract for forks: the new thread is usable as soon
// as the response arrives, so restored usage must follow immediately.
send_thread_token_usage_update_to_connection(
&self.outgoing,
connection_id,
thread_id,
forked_thread.as_ref(),
token_usage_turn_id,
)
.await;
}
self.outgoing
.send_server_notification(ServerNotification::ThreadStarted(notif))
.await;
if inherited_goal {
self.thread_goal_processor
.emit_thread_goal_snapshot(thread_id)
.await;
}
Ok(())ThreadForkResponse 可以包含完整 turns,而 thread_started_notification 会生成适合广播的介绍视图。usage 回放在响应之后,started 之后才发继承 goal 快照;不能把这几种消息合并成“fork 成功事件”。出站队列顺序也不等于所有客户端已经渲染完成。
临时分叉测试会一路读取通知,若在其 started 之前见到属于该 ID 的 status/changed 就失败;同时断言 started 的 turns 为空。这是协议消费者对静默 upsert 与通知投影的反向约束。
响应前的创建过程没有覆盖所有步骤的总事务。下面按真实失败点区分“尚未创建”和“已经有新执行者”,用于定位错误后残留的对象。
remove_pending_project_metadata 只处理 Core 创建失败分支的预置元数据;名称、goal 或后续历史读取失败没有在此统一删除新 Thread。客户端排障时应先辨认失败阶段,再检查已加载对象与持久记录,不能仅凭 RPC 错误推断所有副作用都已回滚。
9.3 两次分叉
现在回到开头的场景,按源码推演两次请求:A 完成 T1、T2,T3 进行中;第一次 lastTurnId=T2 得到 B,第二次从 B 默认分叉得到 C。
- 第一次选择 T2 的终态边界,B 不继承 T3;Paginated 的 history base 固定 A 的前缀。B 的权限可能采用 A 的最新设置,而非 T2 时的旧值。
- 第二次的直接来源是 B,因此 C 的
forked_from_id=B。如果归一化后的前缀恰好退回 A,C 的 history base 可以指 A 的 rollout;这不矛盾。 - A 继续写入 T3 不会扩大 B/C 已固定的前缀。若 B 又执行自己的新 Turn,再分叉时才会选到 B 的新增段。
- 只有明确请求 goal 延后继承且条件满足时,child 才得到当前 goal 快照;目标使用量不会随旧历史截断自动回退。
遇到“新文件没有旧消息”“从旧轮次分叉却得到新权限”“fork 报错后 loaded 列表里已有新 ID”时,分别检查 history base、设置来源与最后完成的创建阶段。不要通过删除源文件或重新复制全量历史来掩盖问题;那会破坏当前实现的引用关系。
以下命令从公开集成测试进入,再分别检查两种边界算法和生命周期保护。在 Codex 仓库根目录执行。
just test --locked -p codex-app-server --test all \
-E 'test(suite::v2::thread_fork::)'
just test --locked -p codex-core --lib \
-E 'test(thread_rollout_truncation::tests) or test(fork_snapshot) or test(interrupted_snapshot)'
just test --locked -p codex-thread-store --lib \
-E 'test(fork) or test(referenced_paginated_rollout_projects_inherited_ordinal_range) or test(active_turn_stores_only_its_start_position)'
rg -n 'thread_fork_inner|prepare_fork|fork_prepared_thread' \
codex-rs/app-server/src/request_processors/thread_processor.rs \
codex-rs/core/src/thread_manager.rs尝试独立复述两条消费者链:历史如何从请求边界变成固定位置、再进入 child 的模型输入;配置与 goal 如何从当前源状态变成 child 的实际运行条件。能把这两条链与源记录的保留周期连起来,才能判断一次 fork 是选错了前缀、引用尚未落盘,还是新执行者的配置与预期不同。
