Skip to content

V2 Thread协议

沿着 thread/start、resume、fork、read 和 archive,理解 V2 Thread 的身份、历史、订阅与关闭边界。

基于rust-v0.150.0
CodexRustAppServerThread

V2 Thread协议 ​

本文承接AppServer V2方法注册机制,面向已经理解 typed request、ThreadId 和 rollout 的读者。V2 Thread API 并不是简单 CRUD:thread/start 建立 live owner,thread/resume 可能重新加入正在运行的 listener,thread/fork 复制历史并生成新身份,thread/read 构造持久化视图,thread/archive 处理归档及其后代。

本文沿真实 handler 说明四个边界:哪些字段只在新线程生效,运行中 resume 为什么必须交给 listener,fork 如何区分历史来源与新 ThreadId,以及 archive 为什么不等于 task 已停止。

1. Start参数 ​

源码位置:codex-rs/app-server-protocol/src/protocol/v2/thread.rs :: ThreadStartParams、ThreadStartResponse

rust
pub struct ThreadStartParams {
    pub model: Option<String>,
    pub model_provider: Option<String>,
    pub allow_provider_model_fallback: bool,
    pub service_tier: Option<Option<String>>,
    pub cwd: Option<String>,
    pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>,
    pub approval_policy: Option<AskForApproval>,
    pub approvals_reviewer: Option<ApprovalsReviewer>,
    pub sandbox: Option<SandboxMode>,
    pub permissions: Option<String>,
    pub config: Option<HashMap<String, JsonValue>>,
    pub service_name: Option<String>,
    pub base_instructions: Option<String>,
    pub developer_instructions: Option<String>,
    pub personality: Option<Personality>,
    pub ephemeral: Option<bool>,
    pub history_mode: Option<ThreadHistoryMode>,
    pub thread_source: Option<ThreadSource>,
    pub project_id: Option<String>,
    pub environments: Option<Vec<TurnEnvironmentParams>>,
    pub dynamic_tools: Option<Vec<DynamicToolSpec>>,
    pub selected_capability_roots: Option<Vec<SelectedCapabilityRoot>>,
    pub experimental_raw_events: bool,
}

ThreadStartParams 把模型、权限、环境、历史模式、动态工具和 project identity 一次性提交给新线程。permissions 与 legacy sandbox 不能同时设置;部分字段带有实验能力标记,但仍由同一个 typed params 承载。

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: thread_start、thread_start_inner

rust
pub(crate) async fn thread_start(
    &self,
    request_id: ConnectionRequestId,
    params: ThreadStartParams,
    app_server_client_name: Option<String>,
    app_server_client_version: Option<String>,
    client_mcp_extensions: ClientMcpExtensions,
    request_context: RequestContext,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
    self.thread_start_inner(
        request_id,
        params,
        app_server_client_name,
        app_server_client_version,
        client_mcp_extensions,
        request_context,
    )
    .await
    .map(|()| None)
}

start handler 返回 None,因为 thread_start_inner 自己发送 response 和 started notification。这个设计让启动完成时可以在同一处装配 listener、状态和首个事件。

2. Start前置检查 ​

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: thread_start_inner

rust
if matches!(
    history_mode,
    Some(codex_app_server_protocol::ThreadHistoryMode::Paginated)
) && !self.thread_store.supports_paginated_history_lists()
{
    return Err(invalid_request(
        "paginated threads require thread/turns/list and thread/items/list support",
    ));
}
if sandbox.is_some() && permissions.is_some() {
    return Err(invalid_request(
        "`permissions` cannot be combined with `sandbox`",
    ));
}
if let Some(project_id) = project_id.as_ref() {
    if project_id.is_empty() {
        return Err(invalid_request("projectId must not be empty"));
    }
    let project = self
        .thread_store
        .read_project(project_id.clone())
        .await
        .map_err(|err| internal_error(format!("failed to read project: {err}")))?;
    if project.is_none() {
        return Err(invalid_request(format!("project not found: {project_id}")));
    }
}

这些检查发生在创建 live thread 前:分页模式要求 store 支持对应 list API,权限 profile 与 sandbox 互斥,project id 必须非空且存在。失败时不会产生可恢复的线程 owner。

3. Resume参数 ​

源码位置:codex-rs/app-server-protocol/src/protocol/v2/thread.rs :: ThreadResumeParams

rust
/// There are three ways to resume a thread:
/// 1. By thread_id: load the thread from disk by thread_id and resume it.
/// 2. By history: instantiate the thread from memory and resume it.
/// 3. By path: load the thread from disk by path and resume it.
pub struct ThreadResumeParams {
    pub thread_id: String,
    pub history: Option<Vec<ResponseItem>>,
    pub path: Option<PathBuf>,
    pub model: Option<String>,
    pub model_provider: Option<String>,
    pub service_tier: Option<Option<String>>,
    pub cwd: Option<String>,
    pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>,
    pub approval_policy: Option<AskForApproval>,
    pub approvals_reviewer: Option<ApprovalsReviewer>,
    pub sandbox: Option<SandboxMode>,
    pub permissions: Option<String>,
    pub config: Option<HashMap<String, serde_json::Value>>,
    pub base_instructions: Option<String>,
    pub developer_instructions: Option<String>,
    pub personality: Option<Personality>,
    pub exclude_turns: bool,
    pub initial_turns_page: Option<ThreadResumeInitialTurnsPageParams>,
}

当前注释规定:对未运行线程,history 优先于非空 path,path 再优先于 thread id;对运行线程,thread id 用于加入 live owner,path 只作为 active rollout path 的一致性检查。空 path 会被反序列化为 None。

4. Resume分流 ​

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: thread_resume_inner

rust
if let Ok(thread_id) = ThreadId::from_string(&params.thread_id)
    && self
        .pending_thread_unloads
        .lock()
        .await
        .contains(&thread_id)
{
    self.outgoing
        .send_error(
            request_id,
            invalid_request(format!(
                "thread {thread_id} is closing; retry thread/resume after the thread is closed"
            )),
        )
        .await;
    return Ok(());
}

let stored_thread_from_running_probe = match self
    .resume_running_thread(
        &request_id,
        &params,
        app_server_client_name.clone(),
        app_server_client_version.clone(),
        /*cold_resume_history*/ None,
    )
    .await
{
    Ok(RunningThreadResumeResult::Handled) => return Ok(()),
    Ok(RunningThreadResumeResult::NotRunning(stored_thread)) => stored_thread,
    Err(error) => {
        self.outgoing.send_error(request_id, error).await;
        return Ok(());
    }
};

resume 先检查 closing 状态,再探测 running thread。探测成功处理后直接返回;只有 NotRunning 才继续读取 history/store。这样运行中 resume 不会在“读取历史”和“重新订阅”之间暴露空窗。

5. Listener接管 ​

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: resume_running_thread;codex-rs/app-server/src/request_processors/thread_lifecycle.rs :: handle_pending_thread_resume_request

rust
if let Some(listener_tx) = self.thread_registry.listener_command_tx(thread_id) {
    listener_tx
        .send(ThreadListenerCommand::SendThreadResumeResponse(
            Box::new(pending_resume_request),
        ))
        .map_err(|_| invalid_request("thread listener is no longer available"))?;
    return Ok(RunningThreadResumeResult::Handled);
}

listener command 把 pending resume 交给拥有实时事件流的线程。listener 随后合并历史、active turn 和 initial page,再把连接加入 thread;调用方不需要自行拼接 live event。

源码位置:codex-rs/app-server/src/request_processors/thread_lifecycle.rs :: handle_pending_thread_resume_request

rust
let active_turn = {
    let state = thread_state.lock().await;
    state.active_turn_snapshot()
};
let has_live_in_progress_turn =
    matches!(conversation.agent_status().await, AgentStatus::Running)
        || active_turn
            .as_ref()
            .is_some_and(|turn| matches!(turn.status, TurnStatus::InProgress));

if pending.include_turns {
    if let Some(turns) = pending.paginated_turns.take() {
        thread.turns = turns;
    } else {
        populate_thread_turns_from_history(
            &mut thread,
            &pending.history_items,
            /*active_turn*/ None,
        );
    }
    if let Some(active_turn) = active_turn.as_ref() {
        merge_turn_history_with_active_turn(&mut thread.turns, active_turn.clone());
    }
}

运行中 resume 的返回视图可能包含 active turn。include_turns=false 则只返回 metadata/live state,客户端可以随后调用 turns/items list;这也是 exclude_turns 存在的原因。

6. Fork创建新身份 ​

源码位置:codex-rs/app-server-protocol/src/protocol/v2/thread.rs :: ThreadForkParams

rust
pub struct ThreadForkParams {
    pub thread_id: String,
    pub last_turn_id: Option<String>,
    pub before_turn_id: Option<String>,
    pub path: Option<PathBuf>,
    pub model: Option<String>,
    pub model_provider: Option<String>,
    pub service_tier: Option<Option<String>>,
    pub cwd: Option<String>,
    pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>,
    pub approval_policy: Option<AskForApproval>,
    pub approvals_reviewer: Option<ApprovalsReviewer>,
    pub sandbox: Option<SandboxMode>,
    pub permissions: Option<String>,
    pub config: Option<HashMap<String, JsonValue>>,
    pub developer_instructions: Option<String>,
    pub ephemeral: bool,
    pub thread_source: Option<ThreadSource>,
    pub exclude_turns: bool,
    pub defer_goal_continuation: bool,
}

last_turn_id 与 before_turn_id 是互斥的边界选择:前者包含指定 turn,后者排除指定 turn 及之后内容。fork 可以从 thread id 或 path 读取来源,并可延迟继承 goal 的自动 continuation。

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: thread_fork_inner

rust
async fn thread_fork_inner(
    &self,
    request_id: ConnectionRequestId,
    params: ThreadForkParams,
    app_server_client_name: Option<String>,
    app_server_client_version: Option<String>,
    client_mcp_extensions: ClientMcpExtensions,
) -> Result<(), JSONRPCErrorError> {
    let ThreadForkParams {
        thread_id,
        last_turn_id,
        before_turn_id,
        path,
        model,
        model_provider,
        service_tier,
        cwd,
        runtime_workspace_roots,
        approval_policy,
        approvals_reviewer,
        sandbox,
        permissions,
        config,
        developer_instructions,
        ephemeral,
        thread_source,
        exclude_turns,
        defer_goal_continuation,
    } = params;
}

fork handler 先拆出来源、边界和配置,再调用 ThreadManager。新线程的 ThreadId 与来源线程不同,forked_from_id 记录历史来源;如果传入的边界指向进行中的 turn,Core 会按当前 fork snapshot 规则冻结或拒绝。

7. Read构造视图 ​

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: thread_read_response_inner、apply_thread_read_store_fields

rust
async fn thread_read_response_inner(
    &self,
    params: ThreadReadParams,
) -> Result<ThreadReadResponse, JSONRPCErrorError> {
    let ThreadReadParams {
        thread_id,
        include_turns,
    } = params;

    let thread_uuid = ThreadId::from_string(&thread_id)
        .map_err(|err| invalid_request(format!("invalid thread id: {err}")))?;
    let thread = self
        .read_thread_view(thread_uuid, include_turns)
        .await
        .map_err(thread_read_view_error)?;
    Ok(ThreadReadResponse { thread })
}

read 是 view builder,不是把内存中的 CodexThread 原样序列化。include turns 时会从 store/history 填充;分页线程还可能只能通过 projected turns/items 读取,不能假设 rollout path 一定存在。

8. Archive与关闭 ​

源码位置:codex-rs/app-server/src/request_processors/thread_processor.rs :: thread_archive

rust
pub(crate) async fn thread_archive(
    &self,
    request_id: ConnectionRequestId,
    params: ThreadArchiveParams,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
    match self.thread_archive_inner(params).await {
        Ok((response, archived_thread_ids)) => {
            self.outgoing
                .send_response(request_id.clone(), response)
                .await;
            for thread_id in archived_thread_ids {
                self.outgoing
                    .send_server_notification(ServerNotification::ThreadArchived(
                        ThreadArchivedNotification { thread_id },
                    ))
                    .await;
            }
            Ok(None)
        }
        Err(error) => Err(error),
    }
}

archive 先返回 response,再对归档线程逐个发 notification。它改变的是持久化归档状态;live task、listener 和 pending requests 的关闭由独立 unload/close lifecycle 处理。归档成功不等于进程已经退出。

9. 失败与恢复 ​

Thread API 的失败阶段至少包括:参数互斥或 project 不存在、thread 正在 closing、running listener 不可用、store 不支持分页、rollout 未物化、历史边界指向进行中的 turn、恢复 payload 需要脱敏,以及按 client name 选择的 resume redaction。每一类错误都应回到对应 owner,而不是统一包装成“thread not found”。

源码位置:codex-rs/app-server/src/request_processors/thread_resume_redaction.rs :: should_redact_thread_resume_payloads、redact_thread_resume_payloads

rust
pub(super) fn should_redact_thread_resume_payloads(client_name: Option<&str>) -> bool {
    client_name.is_some_and(|name| name.eq_ignore_ascii_case("chatgpt"))
}

pub(super) fn redact_thread_resume_payloads(turns: &mut [Turn]) {
    for turn in turns {
        for item in &mut turn.items {
            redact_thread_item(item);
        }
    }
}

resume payload redaction 是公开视图处理,不是底层历史删除;同一线程可以因 client name 得到不同返回内容。

10. 测试与边界 ​

源码位置:codex-rs/app-server/tests/suite/v2/thread_start.rs :: thread_start_creates_thread_and_emits_started、thread_start_rejects_paginated_history_without_list_support

源码位置:codex-rs/app-server/tests/suite/v2/thread_resume.rs :: thread_resume_rejoins_running_paginated_thread_with_initial_page、thread_resume_rejects_mismatched_path_for_running_thread_id、thread_resume_returns_rollout_history

源码位置:codex-rs/app-server/tests/suite/v2/thread_fork.rs :: thread_fork_creates_new_thread_and_emits_started、thread_fork_at_last_turn_id_keeps_only_terminal_prefix、thread_fork_rejects_incompatible_boundaries_and_ephemeral_goal_deferral

源码位置:codex-rs/app-server/tests/suite/v2/thread_read.rs :: thread_read_returns_summary_without_turns、thread_read_can_include_turns、thread_read_returns_forked_from_id_for_forked_threads

源码位置:codex-rs/app-server/tests/suite/v2/thread_archive.rs :: thread_archive_requires_materialized_rollout、thread_archive_archives_spawned_descendants

text
cd codex-rs
cargo test -p codex-app-server thread_start_creates_thread_and_emits_started -- --nocapture --test-threads=1
cargo test -p codex-app-server thread_resume_rejoins_running_paginated_thread_with_initial_page -- --nocapture --test-threads=1
cargo test -p codex-app-server thread_resume_rejects_mismatched_path_for_running_thread_id -- --nocapture --test-threads=1
cargo test -p codex-app-server thread_fork_creates_new_thread_and_emits_started -- --nocapture --test-threads=1
cargo test -p codex-app-server thread_fork_at_last_turn_id_keeps_only_terminal_prefix -- --nocapture --test-threads=1
cargo test -p codex-app-server thread_read_returns_forked_from_id_for_forked_threads -- --nocapture --test-threads=1
cargo test -p codex-app-server thread_archive_requires_materialized_rollout -- --nocapture --test-threads=1

这些测试覆盖 start/resume/fork/read/archive 的代表性路径、listener rejoin、历史边界和物化要求;不能证明所有分页 cursor 组合、远程 thread store、跨平台路径或客户端脱敏策略都已覆盖。

11. 源码定位练习 ​

遇到 resume 漏事件,先确认线程是否 running,再沿 resume_running_thread 和 listener command 阅读;遇到 fork 内容过多或过少,检查 last_turn_id、before_turn_id、source path 和 fork snapshot;遇到 archive 后仍有执行,分别检查 archive store 状态和 unload/close lifecycle。

遇到 read 返回空 turns,检查 include_turns、paginated projection、rollout materialization 和 store history;遇到同一请求在不同客户端返回字段不同,检查 should_redact_thread_resume_payloads。Thread API 的核心不是五个方法名,而是 live owner、持久化历史、投影视图和订阅时机之间的契约。