Skip to content

ThreadManager分叉

追踪 ThreadManager 如何截取源 Thread 历史、生成新身份,并以复制或引用方式建立可恢复的分叉。

基于rust-v0.150.0
CodexRustThreadManagerFork

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 ​

维度ResumeFork
ThreadId保留源 ID分配新 UUIDv7
root SessionId从源 metadata 恢复新 root 通常使用新 ThreadId
历史载体InitialHistory::ResumedInitialHistory::Forked
rollout重开源 writer创建 child rollout
lineage保留既有 lineageforked_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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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截止位置关键拒绝条件
Latestsource 当前 durable positionlineage、projection state 或 writer materialization 失败
ThroughTurn(id)包含最新可见的该 TurnTurn 仍 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

rust
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);
result

reservation 只阻止删除,不阻止 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 DBUnsupported { operation: "prepare_fork" }
named boundaryinvisible、in-progress 或无持久化位置InvalidRequest,不猜测近似位置
lineagecycle、缺文件、ID 不一致、越过祖先 cutoff拒绝引用,避免 child 指向错误物理范围
child Session 初始化persistence/Hook/MCP/network 等失败LiveThreadInitGuard discard child rollout
referenced child 创建child reference 尚未 durablePreparedFork reservation 继续阻止 source 删除
fork future 取消lineage materialization 已移入 detached taskreservation 保持到 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 使用错误 environmentenvironment 不从 rollout 自动恢复,应检查本次 Config/selections
paginated fork 找不到 TurnTurn 是否可见、是否 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

rust
// :: 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

rust
// :: 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

rust
// :: 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. 分叉边界验证 ​

  1. 对比Resume与Fork:分别说明ThreadId、history信封、writer动作和lineage结果。
  2. 给定包含环境contextual wrapper、两条真实用户消息和未完成工具调用的history,说明Nth-user截断应忽略 什么、保留什么,以及何时追加Interrupted终态。
  3. 解释Copied child为何可以独立读取却占用更多存储,Referenced child为何需要source reservation和 history_base。
  4. 使用只读命令核对截断、marker和lineage测试:
bash
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创建。