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
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
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
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
/// 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
if let Ok(thread_id) = ThreadId::from_string(¶ms.thread_id)
&& self
.pending_thread_unloads
.lock()
.await
.contains(&thread_id)
{
self.outgoing
.send_error(
request_id,
invalid_request(format!(
"thread {thread_id} is closing; retry thread/resume after the thread is closed"
)),
)
.await;
return Ok(());
}
let stored_thread_from_running_probe = match self
.resume_running_thread(
&request_id,
¶ms,
app_server_client_name.clone(),
app_server_client_version.clone(),
/*cold_resume_history*/ None,
)
.await
{
Ok(RunningThreadResumeResult::Handled) => return Ok(()),
Ok(RunningThreadResumeResult::NotRunning(stored_thread)) => stored_thread,
Err(error) => {
self.outgoing.send_error(request_id, error).await;
return Ok(());
}
};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
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
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
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
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
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
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
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
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、持久化历史、投影视图和订阅时机之间的契约。
