ThreadManager分叉
Fork 会继承一段源历史,但它不是 resume 的别名。resume 保留原 ThreadId 并重新打开同一 Thread;fork 把选中的历史前缀转换为 InitialHistory::Forked,创建新的 ThreadId,再用 forked_from_id 记录直接来源。
真正困难的部分是“历史前缀”并不总在稳定 Turn 边界上。源 rollout 可能停在一轮模型采样或工具调用 中间,也可能由多段 paginated lineage 组成。Codex 因此把 fork 拆成三项独立决策:选择快照边界、 归一化未完成 Turn、选择复制或引用持久化。
本文面向已经读过 ThreadManager恢复 的读者。Resume保留原 ThreadId并重开运行时,Fork则必须生成新身份;本文集中解释快照边界、Interrupted补全和Copied/Referenced 持久化,不重复展开通用 Session 创建事务。
读完后,应能从 fork_thread() 追踪到 ForkSnapshot、用户边界截断、InitialHistory::Forked、新 lineage 和 child持久化,并能解释为什么“截取第N条用户消息”不能按所有 user-role message简单计数。
1. Fork与Resume
| 维度 | Resume | Fork |
|---|---|---|
| ThreadId | 保留源 ID | 分配新 UUIDv7 |
| root SessionId | 从源 metadata 恢复 | 新 root 通常使用新 ThreadId |
| 历史载体 | InitialHistory::Resumed | InitialHistory::Forked |
| rollout | 重开源 writer | 创建 child rollout |
| lineage | 保留既有 lineage | forked_from_id 指向直接 source |
| active handle | 可能直接复用 | 必定创建新 handle |
| 当前运行状态 | 不可跨进程恢复 | 也不复制 token、channel、waiter 或 task |
用户发起的普通 fork 通常只有 forked_from_id,没有 agent-control 意义上的 parent_thread_id。subagent 创建路径可以同时设置 parent 与 fork source:前者描述控制关系,后者描述历史从哪里复制,不能把两个 字段当成同义词。
边界选择发生在新 Session 创建之前。无论历史来自路径、内存还是 paginated preparation,最终都进入 fork_thread_with_initial_history(),再沿 ThreadManager创建 展开的统一创建事务生成 child。
2. ForkSnapshot
ForkSnapshot 不是 UI 参数的简单枚举。每个变体都定义了越界或 mid-turn 时的确定行为。
源码位置:codex-rs/core/src/thread_manager.rs :: ForkSnapshot
pub enum ForkSnapshot {
// n 是从 0 开始的用户消息边界;结果严格结束在第 n 条用户消息之前。
TruncateBeforeNthUserMessage(usize),
// 以“现在中断源 Turn”的语义生成一致快照,但不实际中断 live source。
Interrupted,
}
impl From<usize> for ForkSnapshot {
fn from(value: usize) -> Self {
// 旧 usize 调用点继续映射到历史截断模式,保持源码兼容。
Self::TruncateBeforeNthUserMessage(value)
}
}TruncateBeforeNthUserMessage(0) 表示第一个用户消息之前;n 位于范围内时,历史在对应 user boundary 前截断。n 越界并不总是错误:源已位于稳定边界时保留全部 committed history;源停在 mid-turn 时, 退回到活动 Turn 的 opening boundary,丢掉未完成后缀。
Interrupted 则保留当前 persisted prefix。若 prefix 已在 TurnComplete/TurnAborted 边界,不增加任何 内容;若停在 Turn 中间,追加与 live interrupt 相同的 marker 和 TurnAborted event,使 child 的历史 不再看起来像仍有一个运行中的旧 task。
当前源码还保留了旧调用点的兼容映射:impl From<usize> for ForkSnapshot 把整数转换为 TruncateBeforeNthUserMessage,但新的调用方应显式选择 snapshot 语义。对于 paginated lineage, ThreadStore 先生成 PreparedFork,再由 reference-backed fork 入口传入 ForkPersistence::Referenced; 这条路径不会把所有祖先 item 复制进 child rollout。
3. 路径与预加载
fork_thread() 从 rollout path 读取 InitialHistory;fork_thread_from_history() 接收上层已经加载的 历史。后者默认使用 ForkPersistence::Copied。
源码位置:codex-rs/core/src/thread_manager.rs :: fork_thread, fork_thread_from_history
pub async fn fork_thread<S>(
&self,
snapshot: S,
config: Config,
path: PathBuf,
thread_source: Option<ThreadSource>,
parent_trace: Option<W3cTraceContext>,
) -> CodexResult<NewThread>
where
S: Into<ForkSnapshot>,
{
let snapshot = snapshot.into();
// 路径读取仍经过 ThreadStore;fork 不直接打开并解析任意 JSONL。
let history = self.initial_history_from_rollout_path(path).await?;
self.fork_thread_from_history(
snapshot,
config,
history,
thread_source,
parent_trace,
ClientMcpExtensions::default(),
)
.await
}
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,
) -> CodexResult<NewThread>
where
S: Into<ForkSnapshot>,
{
self.fork_thread_with_initial_history(
config,
ForkHistory {
snapshot: snapshot.into(),
initial_history: history,
// 普通 eager history fork 把 surviving prefix 写入 child rollout。
persistence: ForkPersistence::Copied,
},
thread_source,
parent_trace,
client_mcp_extensions,
)
.await
}路径入口与 ThreadManager恢复 的 resume path 读取共享相同 store 约束:允许 archived、要求完整 eager history, 因此不能直接承担 paginated lineage 的完整读取。paginated fork 先由 store 生成 PreparedFork,再走 专门的 reference-backed 入口。
4. Turn 快照边界
源码用 SnapshotTurnState 保存四项结果:是否 mid-turn、显式 Turn ID、开始时间,以及活动 Turn 在 rollout 中的起始位置。
源码位置:codex-rs/core/src/thread_manager.rs :: SnapshotTurnState, snapshot_turn_state
struct SnapshotTurnState {
ends_mid_turn: bool,
active_turn_id: Option<String>,
active_turn_started_at: Option<i64>,
active_turn_start_index: Option<usize>,
}
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)
{
// reducer 已经得到 terminal 状态时,不因残留 active 结构误判 mid-turn。
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(),
};
}
let Some(last_user_position) = truncation::user_message_positions_in_rollout(rollout_items)
.last()
.copied()
else {
return SnapshotTurnState {
ends_mid_turn: false,
active_turn_id: None,
active_turn_started_at: None,
active_turn_start_index: None,
};
};
// legacy/synthetic history 没有显式 lifecycle 时,用最后用户消息后的 terminal event 判断。
SnapshotTurnState {
ends_mid_turn: !rollout_items[last_user_position + 1..].iter().any(|item| {
matches!(
item,
RolloutItem::EventMsg(EventMsg::TurnComplete(_) | EventMsg::TurnAborted(_))
)
}),
active_turn_id: None,
active_turn_started_at: None,
active_turn_start_index: None,
}
}有显式 TurnStarted 的现代历史由 ThreadHistoryBuilder 判断;只有缺少 lifecycle 的 synthetic/legacy 历史才采用“最后用户消息之后是否存在 terminal event”的后备规则。后备路径不会凭空发明 Turn ID。
判断结果只描述 persisted snapshot,不会读取 live task 的 CancellationToken。因此即使源 Thread 已从 manager map 移除,只要 rollout 停在 mid-turn,fork 仍能生成闭合历史。
5. 按用户消息截断
截断函数先把 Resumed/Forked history 统一成可变 items,再计算所有 user-message position。
源码位置:codex-rs/core/src/thread_manager.rs :: truncate_before_nth_user_message
fn truncate_before_nth_user_message(
history: InitialHistory,
n: usize,
snapshot_state: &SnapshotTurnState,
) -> InitialHistory {
let mut items = match history {
InitialHistory::New | InitialHistory::Cleared => Vec::new(),
InitialHistory::Resumed(resumed) => Arc::unwrap_or_clone(resumed.history),
InitialHistory::Forked(items) => items,
};
let user_positions = truncation::user_message_positions_in_rollout(&items);
let rolled = if snapshot_state.ends_mid_turn && n >= user_positions.len() {
// 越界但源处于 mid-turn 时,优先切到显式 TurnStarted;legacy 才退到最后用户消息。
if let Some(cut_idx) = snapshot_state
.active_turn_start_index
.or_else(|| user_positions.last().copied())
{
items.truncate(cut_idx);
items
} else {
items
}
} else {
truncation::truncate_rollout_before_nth_user_message_from_start(items, n)
};
if rolled.is_empty() {
InitialHistory::New
} else {
InitialHistory::Forked(rolled)
}
}| 输入情况 | 结果 |
|---|---|
n 命中第 N 条用户消息 | 严格保留该消息之前的 committed prefix |
n 等于用户消息数,且已在 Turn 边界 | 保留完整历史 |
n 越界,且历史停在 mid-turn | 删除活动 Turn 的 opening boundary 及其后缀 |
| 没有用户消息且没有可识别活动边界 | 保留已有非用户历史;空结果转为 InitialHistory::New |
测试中的 n=1 会保留第一条用户消息及其 assistant 响应,严格停在第二条用户消息之前;n=2 等于 消息数且 source 已闭合时保持完整。越界不是数组错误,而是显式定义的快照语义。
6. Interrupted
fork_history_from_snapshot() 先把 resumed history 转成 fork-owned vector。只有 snapshot_state.ends_mid_turn 为 true 时,才调用 append_interrupted_boundary()。
源码位置:codex-rs/core/src/thread_manager.rs :: fork_history_from_snapshot, append_interrupted_boundary
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 {
// 只闭合 child snapshot,不向源 Session 提交 Interrupt。
append_interrupted_boundary(
history,
snapshot_state.active_turn_id,
snapshot_state.active_turn_started_at,
interrupted_marker,
)
} else {
history
}
}
}
}
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));
}
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));
}
// marker 可禁用,但 TurnAborted 始终作为协议终态追加。
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));
}
history.push(aborted_event);
InitialHistory::Forked(history)
}
}
}三个分支都返回 InitialHistory::Forked:空历史从新 vector 开始,已有 Forked history 原地追加, Resumed history 则在必要时通过 Arc::unwrap_or_clone 取得 child-owned vector。
marker 由 InterruptedTurnHistoryMarker::from_config_and_version() 决定:可能是 contextual user message、 Multi-agent V2 使用的 developer message,也可以禁用。无论 marker 是否存在,TurnAborted 都会记录 Interrupted reason;显式 Turn ID 存在时原样保留,legacy history 不会伪造 ID。
已经包含该终态的快照再次 fork 时,snapshot_turn_state() 判定它不再 mid-turn,因此不会重复追加 marker 或 abort event。测试专门断言 re-fork 后两者仍各只有一份。
7. 新身份与lineage
边界归一化完成后,manager 计算可用的 source lineage ID,并把它写入 ThreadSpawnRequest.forked_from_thread_id。路径读取形成的 Resumed history 能给出直接 source ID;调用方 若只提供脱离具体 Thread 的 Forked history,则只能保留其中已有 lineage。Session 最终收到 InitialHistory::Forked,因此会分配新的 ThreadId。
源码位置:codex-rs/core/src/thread_manager.rs :: fork_thread_with_initial_history
let ForkHistory {
snapshot,
initial_history: history,
persistence: fork_persistence,
} = fork_history;
let source_thread_id = match &history {
// 从 resume history 分叉时,直接来源就是被恢复的 conversation。
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,
..StartThreadOptions::new(config)
};
let mut request =
ThreadSpawnRequest::new(options, Arc::clone(&self.state.auth_manager), agent_control);
// 历史 lineage 与持久化策略在 Session 创建前写入,不能在注册后补记。
request.forked_from_thread_id = source_thread_id;
request.fork_persistence = fork_persistence;
Box::pin(self.state.spawn_thread(request)).await对普通 path/resume source,forked_from_id 是直接来源,即使源本身已有祖先;对没有当前 Thread 身份的 detached Forked history,源码不会伪造一层来源。新 root fork 不继承旧 root SessionId;subagent fork 若使用 non-root source,则由传入的 AgentControl 保持 agent-tree SessionId。
历史中的 selected capability roots 会进入 child extension init;thread environment selections 则不会 从 rollout 自动恢复,child 使用调用方 Config 或新请求提供的 environments。能力 lineage 与执行环境 生命周期不同,不能整体复制。
8. Copied fork
普通 fork 使用 ForkPersistence::Copied。Session 先用全部 fork items 重建内存 history,再把复制前缀 和当前 child settings 写入自己的 rollout。
源码位置:codex-rs/core/src/session/mod.rs :: ForkPersistence, record_initial_history fork branch
pub(crate) enum ForkPersistence {
Copied,
Referenced {
history_base: Option<HistoryPosition>,
inherited_item_count: usize,
},
}
let thread_settings_applied =
RolloutItem::EventMsg(handlers::thread_settings_applied_event(self).await);
match &self.fork_persistence {
ForkPersistence::Referenced {
inherited_item_count,
..
} => {
// 引用前缀属于 ancestor,只把 child 本地 settings 与边界写入新 rollout。
rollout_items.drain(..*inherited_item_count);
rollout_items.insert(0, thread_settings_applied);
}
ForkPersistence::Copied if is_paginated_subagent => {
rollout_items.clear();
rollout_items.push(thread_settings_applied);
}
ForkPersistence::Copied => {
// 普通 copied fork 将前缀与 child effective settings 作为同一次 append 持久化。
rollout_items.push(thread_settings_applied);
}
}
self.persist_rollout_items(&rollout_items).await;
self.ensure_rollout_materialized().await;复制会增加 child rollout 体积,但 child 可以脱离 source 独立恢复和删除。ThreadSettingsApplied 放在 复制前缀之后,确保 cold resume 读取到的是 child 当前 effective settings,而不是祖先最后一次设置。
9. Referencedfork
Paginated history 可能跨多个物理 rollout,复制整条祖先链既昂贵又破坏分页模型。Local store 因此先 prepare_fork(),返回冻结的 history_base 与 bounded model context。
源码位置:codex-rs/thread-store/src/types.rs :: ForkBoundary, PreparedFork
pub enum ForkBoundary {
Latest,
ThroughTurn(String),
BeforeTurn(String),
}
pub struct PreparedFork {
pub source_thread_id: ThreadId,
// position 可以指向 source 的祖先物理 rollout,不一定等于 source ThreadId。
pub history_base: Option<HistoryPosition>,
pub model_context: Arc<Vec<RolloutItem>>,
// reservation 阻止 source 在 child reference durable 前被删除。
_source_reservation: Box<dyn std::fmt::Debug + Send>,
}三种 named boundary 的语义如下:
| Boundary | 截止位置 | 关键拒绝条件 |
|---|---|---|
Latest | source 当前 durable position | lineage、projection state 或 writer materialization 失败 |
ThroughTurn(id) | 包含最新可见的该 Turn | Turn 仍 inProgress,或缺失 persisted end position |
BeforeTurn(id) | 位于该 Turn 原始 start 之前 | 没有 persisted start boundary |
prepare 需要 state DB 来定位可见 Turn 和物理 ordinal。它先持有 source lifecycle reservation,持久化 pending writer,解析完整 RolloutLineage,再将涉及的物理段 materialize 到 SQLite。lineage cycle、 缺失 rollout、metadata ThreadId 不匹配或非 paginated segment 都会拒绝准备。
history_base 是 HistoryPosition { thread_id, end_ordinal_exclusive, end_byte_offset }。若 boundary 正好 位于某 segment 的起始 ordinal,代码会回退到前一个 segment 的 end,避免产生指向空前缀的无效引用。 随后 model_context::load_for_fork() 只加载该冻结位置之前的上下文。
源码位置:codex-rs/core/src/thread_manager.rs :: fork_prepared_thread
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 {
// Latest 仍可能落在 mid-turn,统一使用 Interrupted 归一化未完成后缀。
snapshot: ForkSnapshot::Interrupted,
initial_history: history,
persistence: fork_persistence,
},
thread_source,
parent_trace,
client_mcp_extensions,
)
.await;
// child reference 已成功 materialize或创建失败后,才释放 source reservation。
drop(prepared);
resultreservation 只阻止删除,不阻止 source 继续追加。测试在 prepare 后向 source 写入新 Turn,并确认 child model context 不包含该消息:fork 看到的是冻结边界,而不是创建期间不断移动的最新状态。
10. 失败与清理边界
| 阶段 | 失败或边界条件 | 行为 |
|---|---|---|
| rollout path 读取 | 文件/metadata/history 无效 | 尚未创建 child,直接返回读取错误 |
| Nth 截断 | n 越界 | 稳定历史保持完整;mid-turn 删除未完成后缀,不报数组错误 |
| Interrupted | 已存在 terminal boundary | 不重复追加 marker 或 TurnAborted |
| paginated prepare | 无 state DB | Unsupported { operation: "prepare_fork" } |
| named boundary | invisible、in-progress 或无持久化位置 | InvalidRequest,不猜测近似位置 |
| lineage | cycle、缺文件、ID 不一致、越过祖先 cutoff | 拒绝引用,避免 child 指向错误物理范围 |
| child Session 初始化 | persistence/Hook/MCP/network 等失败 | LiveThreadInitGuard discard child rollout |
| referenced child 创建 | child reference 尚未 durable | PreparedFork reservation 继续阻止 source 删除 |
| fork future 取消 | lineage materialization 已移入 detached task | reservation 保持到 materialization 结束,再允许删除 |
普通 copied fork 失败不会改变 source;Interrupted 也只修改 child snapshot,不向 source 发送控制 Op。 Referenced fork 的额外风险是 source 生命周期:必须先让 child 的 history_base durable,才能允许删除 source,否则 child 的逻辑历史会断链。
11. Fork问题定位
| 现象 | 优先检查 |
|---|---|
| child 与 source ThreadId 相同 | 是否误走 resume;fork 的 Session 必须从 InitialHistory::Forked 分配新 ID |
| child 多出一条中断消息 | source snapshot 是否 mid-turn、marker mode 是否启用 |
| re-fork 重复中断 marker | 第一次 fork 是否持久化了 TurnAborted,snapshot_turn_state 是否识别 terminal |
| 第 N 条消息截断多保留一轮 | n 是 0-based 且语义为“严格在第 n 条用户消息之前” |
| 越界截断丢掉最后一轮 | source 当时处于 mid-turn,这是删除 unfinished suffix 的定义行为 |
| child 使用错误 environment | environment 不从 rollout 自动恢复,应检查本次 Config/selections |
| paginated fork 找不到 Turn | Turn 是否可见、是否 completed、projection 是否有 start/end ordinal |
| source 删除一直等待 | 是否仍有 PreparedFork reservation,child reference 是否已经 materialize |
| child 恢复时缺祖先历史 | history_base、rollout lineage 和 ancestor segment 是否仍有效 |
验证 fork 不能只比较历史文本。至少还要断言:新旧 ThreadId 不同、forked_from_id 正确、显式 Turn ID 在 abort event 中保留、disabled marker 只追加 event、re-fork 不重复边界、copied/reference 两种模式在 cold resume 后都得到同一模型上下文。
12. Fork语义测试
12.1 用户消息边界
测试构造 u1 → a1 → a2 → u2 → a3 → reasoning → tool → a4。n=1 表示在第二条真实用户消息 u2 之前截断,结果必须保留 u1/a1/a2。这可以防止 reasoning或tool item被误当成新的用户轮次。
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: truncates_before_requested_user_message(节选)
let items = [
user_msg("u1"),
assistant_msg("a1"),
assistant_msg("a2"),
user_msg("u2"),
assistant_msg("a3"),
ResponseItem::Reasoning {
id: Some(ResponseItemId::with_suffix("rs", "1")),
summary: vec![ReasoningItemReasoningSummary::SummaryText {
text: "s".to_string(),
}],
content: None,
encrypted_content: None,
internal_chat_message_metadata_passthrough: None,
},
ResponseItem::FunctionCall {
id: None,
call_id: "c1".to_string(),
name: "tool".to_string(),
namespace: None,
arguments: "{}".to_string(),
encrypted_function_args: None,
internal_chat_message_metadata_passthrough: None,
},
assistant_msg("a4"),
];
let initial = items
.iter()
.cloned()
.map(RolloutItem::ResponseItem)
.collect::<Vec<_>>();
let truncated = truncate_before_nth_user_message(
InitialHistory::Forked(initial),
/*n*/ 1,
&SnapshotTurnState {
ends_mid_turn: false,
active_turn_id: None,
active_turn_started_at: None,
active_turn_start_index: None,
},
);
let expected_items = vec![
RolloutItem::ResponseItem(items[0].clone()),
RolloutItem::ResponseItem(items[1].clone()),
RolloutItem::ResponseItem(items[2].clone()),
];
// 第二条真实用户消息及其后所有item都从child前缀排除。
assert_eq!(
serde_json::to_value(truncated.get_rollout_items()).unwrap(),
serde_json::to_value(&expected_items).unwrap(),
);完整测试还验证 n=2 保留原历史。它证明稳定历史的Nth-user规则;mid-turn越界时只删除unfinished suffix, 由 out_of_range_truncation_drops_only_unfinished_suffix_mid_turn 单独验证。
12.2 未完成Turn
append_interrupted_boundary() 不修改 source,而是在 child snapshot尾部追加可选模型可见marker和 TurnAborted(Interrupted)。测试覆盖已有历史和空历史,保证marker位于协议终态之前。
源码位置:codex-rs/core/src/thread_manager_tests.rs
// :: interrupted_fork_snapshot_appends_interrupt_boundary(断言节选)
let committed_history =
InitialHistory::Forked(vec![RolloutItem::ResponseItem(user_msg("hello"))]);
assert_eq!(
serde_json::to_value(
append_interrupted_boundary(
committed_history,
/*turn_id*/ None,
/*started_at*/ None,
InterruptedTurnHistoryMarker::ContextualUser,
)
.get_rollout_items(),
)
.expect("serialize interrupted fork history"),
serde_json::to_value(vec![
RolloutItem::ResponseItem(user_msg("hello")),
// 模型可见marker先于TurnAborted,恢复后模型不会把partial Turn当作正常结束。
RolloutItem::ResponseItem(contextual_user_interrupted_marker()),
RolloutItem::EventMsg(EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: None,
started_at: None,
reason: TurnAbortReason::Interrupted,
completed_at: None,
duration_ms: None,
})),
])
.expect("serialize expected interrupted fork history"),
);另一个测试在marker配置关闭时只追加 TurnAborted;MultiAgentV2测试要求developer-role marker。这些断言 证明child snapshot闭合,不证明source live Task被取消——Fork本来就不应中断source。
12.3 Referenced边界
Paginated source可能已经继承ancestor。测试建立ancestor→child lineage并追加child Turn,随后枚举Latest、 ThroughTurn和BeforeTurn,断言每个boundary对应的 history_base,并检查model context是否包含ancestor和 child消息。
源码位置:codex-rs/thread-store/src/local/thread_history_materialization_tests.rs
// :: referenced_paginated_rollout_projects_inherited_ordinal_range(边界断言节选)
let latest_history_base = prepare_paginated_fork(
&store,
child_id,
ForkBoundary::Latest,
)
.await
.history_base;
for (boundary, expected_base) in [
(ForkBoundary::Latest, latest_history_base),
(
ForkBoundary::ThroughTurn("source-turn".to_string()),
Some(history_base),
),
(
ForkBoundary::BeforeTurn("child-turn".to_string()),
Some(history_base),
),
(
ForkBoundary::ThroughTurn("child-turn".to_string()),
latest_history_base,
),
(ForkBoundary::BeforeTurn("source-turn".to_string()), None),
] {
let prepared = prepare_paginated_fork(&store, child_id, boundary).await;
// history_base冻结child未来读取source lineage的物理上界。
assert_eq!(prepared.history_base, expected_base);
assert_eq!(
contains_user_message(&prepared.model_context, "source message"),
expected_base.is_some(),
);
assert_eq!(
contains_user_message(&prepared.model_context, "child message"),
expected_base == latest_history_base,
);
}该测试证明Referenced模式的边界与model context一致。Copied模式没有 history_base:迁移测试 migration_preserves_copied_user_fork_history_without_creating_a_history_base 验证 copied parent ResponseItems仍在child物理rollout中,同时child metadata的 history_base.is_none()。
12.4 取消时序
Referenced fork准备可能触发压缩source的lineage materialization。测试用ancestor writer lock阻塞该工作, 启动 prepare_fork() 后主动abort task;随后并发删除source必须继续等待。只有释放ancestor lock、detached materialization结束后,delete才可继续。
这证明Future取消不会释放仍保护lineage的reservation。另一个测试进一步证明child reference持久化前 delete等待,持久化后source因已有durable child reference而拒绝删除。这两项验证Store lineage安全,不是 manager新Thread注册。
13. 分叉边界验证
- 对比Resume与Fork:分别说明ThreadId、history信封、writer动作和lineage结果。
- 给定包含环境contextual wrapper、两条真实用户消息和未完成工具调用的history,说明Nth-user截断应忽略 什么、保留什么,以及何时追加Interrupted终态。
- 解释Copied child为何可以独立读取却占用更多存储,Referenced child为何需要source reservation和
history_base。 - 使用只读命令核对截断、marker和lineage测试:
rg -n "truncate_before_nth_user_message|append_interrupted_boundary" \
codex-rs/core/src/thread_manager.rs codex-rs/core/src/thread_manager_tests.rs
rg -n "referenced_paginated_rollout|cancelled_fork_keeps_source_reserved|copied_user_fork" \
codex-rs/thread-store/src/local继续学习时,阅读 ThreadManager管理,理解source/child handle在 remove和shutdown后的可见性;需要回看通用创建事务时,返回 ThreadManager创建。
