V2 Turn与Item协议
本文承接V2 Thread协议和ConversationItem类型体系。Thread 是长期 owner,Turn 是一次模型工作周期,Item 是周期内可观察的消息、命令、MCP、动态工具、计划或推理单元。V2 协议把三者拆开:turn/start 的 response 只确认输入被接受,item/started/item/completed 描述局部生命周期,turn/completed 才描述整轮终态。
如果把它们当作一个状态字段,就会误判:item 已完成不代表 Turn 完成,Turn 被中断也不保证每个 item 都有 completed,Token usage 和 plan notification 还可能在 Turn 终态之前到达。
1. Turn状态与请求
源码位置:codex-rs/app-server-protocol/src/protocol/v2/turn.rs :: TurnStatus、TurnStartParams、TurnSteerParams、TurnInterruptParams、Turn
pub enum TurnStatus {
Completed,
Interrupted,
Failed,
InProgress,
}
pub struct TurnStartParams {
pub thread_id: String,
pub client_user_message_id: Option<String>,
pub input: Vec<UserInput>,
pub responsesapi_client_metadata: Option<HashMap<String, String>>,
pub additional_context: Option<HashMap<String, AdditionalContextEntry>>,
pub environments: Option<Vec<TurnEnvironmentParams>>,
pub cwd: Option<PathBuf>,
pub runtime_workspace_roots: Option<Vec<AbsolutePathBuf>>,
pub approval_policy: Option<AskForApproval>,
pub approvals_reviewer: Option<ApprovalsReviewer>,
pub sandbox_policy: Option<SandboxPolicy>,
pub permissions: Option<String>,
pub model: Option<String>,
pub effort: Option<ReasoningEffort>,
pub summary: Option<ReasoningSummary>,
pub personality: Option<Personality>,
pub output_schema: Option<JsonValue>,
pub collaboration_mode: Option<CollaborationMode>,
}
pub struct TurnSteerParams {
pub thread_id: String,
pub client_user_message_id: Option<String>,
pub input: Vec<UserInput>,
pub responsesapi_client_metadata: Option<HashMap<String, String>>,
pub additional_context: Option<HashMap<String, AdditionalContextEntry>>,
pub expected_turn_id: String,
}
pub struct TurnInterruptParams {
pub thread_id: String,
pub turn_id: String,
}turn/start 与 turn/steer 都带 UserInput,但 steer 多了 expected_turn_id 前置条件;interrupt 只带 thread/turn 身份。start 的 settings 可影响后续 turn,steer 不会替换活动 turn 的上下文。
源码位置:codex-rs/app-server-protocol/src/protocol/v2/thread_data.rs :: Turn、TurnItemsView、TurnError
pub struct Turn {
pub id: String,
pub items: Vec<ThreadItem>,
pub items_view: TurnItemsView,
pub status: TurnStatus,
pub error: Option<TurnError>,
pub started_at: Option<i64>,
pub completed_at: Option<i64>,
pub duration_ms: Option<i64>,
}
pub enum TurnItemsView {
NotLoaded,
Summary,
Full,
}items_view 说明 items 的加载程度,不是 Turn 状态。NotLoaded 可以出现在 in-progress 或 metadata-only response,Summary 只包含展示摘要,Full 才表示持久化历史中可用的完整 item 集合。
2. turn/start
源码位置:codex-rs/app-server/src/request_processors/turn_processor.rs :: turn_start、turn_start_inner
pub(crate) async fn turn_start(
&self,
request_id: ConnectionRequestId,
params: TurnStartParams,
app_server_client_name: Option<String>,
app_server_client_version: Option<String>,
) -> Result<Option<ClientResponsePayload>, JSONRPCErrorError> {
validate_user_input_image_urls(¶ms.input)?;
self.turn_start_inner(
request_id,
params,
app_server_client_name,
app_server_client_version,
)
.await
.map(|response| Some(response.into()))
}具体的 turn_start_inner 会解析 thread、环境、权限、模型和输入,再提交 Core 的 reply-bearing turn input。方法返回的 TurnStartResponse 只带一个初始 Turn 视图,后续 item/turn 通知通过独立事件路径发送。
源码位置:codex-rs/app-server/src/request_processors/turn_processor.rs :: turn_start_inner 的输入提交路径
let submission = thread
.start_or_steer_turn(
TurnInputRequest::new(TurnInput::UserInput {
content: mapped_items,
client_id: client_user_message_id,
})
.with_thread_settings(thread_settings)
.on_start(TurnStartOptions {
final_output_json_schema: params.output_schema,
..Default::default()
})
.with_additional_context(additional_context)
.with_responses_metadata(params.responsesapi_client_metadata)
.with_trace(self.request_trace_context(&request_id).await),
)
.await
.map_err(|err| internal_error(format!("failed to submit turn input: {err}")))?;提交结果还会区分 Started、Steered 和 NotSubmitted:
let (turn_id, started) = match submission {
TurnInputSubmission::Started { turn_id } => (turn_id, true),
TurnInputSubmission::Steered { turn_id } => (turn_id, false),
TurnInputSubmission::NotSubmitted { reason } => {
return Err(internal_error(format!("failed to submit turn input: {reason:?}")));
}
};提交成功表示 Core 接受了输入,不表示 hook、历史写入或模型采样完成。若线程不存在、权限/环境参数无法解析或输入超过限制,错误会在初始 response 阶段返回。
3. steer与interrupt
源码位置:codex-rs/app-server/src/request_processors/turn_processor.rs :: turn_steer_inner
let submission = thread
.steer_turn(
TurnInputRequest::new(TurnInput::UserInput {
content: mapped_items,
client_id: params.client_user_message_id,
})
.with_additional_context(additional_context)
.with_responses_metadata(params.responsesapi_client_metadata),
params.expected_turn_id,
)
.await
.map_err(|err| {
let error = internal_error(format!("failed to steer turn: {err}"));
self.track_error_response(request_id, &error, /*error_type*/ None);
error
})?;steer 必须匹配当前活动 Turn ID;没有活动 Turn、ID 不匹配或当前任务不可 steering 时,Core 返回明确失败,而不是创建一个隐式新 Turn。
源码位置:codex-rs/app-server/src/request_processors/turn_processor.rs :: turn_interrupt_inner
let is_startup_interrupt = thread.agent_status().await == AgentStatus::PendingInit;
if !is_startup_interrupt {
let thread_state = self.thread_state_manager.thread_state(thread_uuid).await;
let mut thread_state = thread_state.lock().await;
thread_state.pending_interrupts.push(request_id.clone());
}
match self
.submit_core_op(request_id, thread.as_ref(), Op::Interrupt)
.await
{
Ok(_) => {}
Err(err) => {
return Err(internal_error(format!("failed to interrupt turn: {err}")));
}
}非 startup interrupt 会先登记 pending request,直到 Core 发出 TurnAborted 才回复客户端。interrupt response 因而是生命周期事件驱动的,不是“收到请求就立即成功”。
4. Item生命周期通知
源码位置:codex-rs/app-server-protocol/src/protocol/v2/item.rs :: ItemStartedNotification、ItemCompletedNotification
pub struct ItemStartedNotification {
pub item: ThreadItem,
pub thread_id: String,
pub turn_id: String,
pub started_at_ms: i64,
}
pub struct ItemCompletedNotification {
pub item: ThreadItem,
pub thread_id: String,
pub turn_id: String,
pub completed_at_ms: i64,
}item 的配对键来自 ThreadItem.id,外层 thread_id/turn_id 决定它属于哪一条 Thread/Turn。不同 item 类型的终态字段不同:命令有 exit code,MCP 有 result/error,动态工具有 success/content,计划只有最终文本。
源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: EventMsg::ItemStarted、EventMsg::ItemCompleted
EventMsg::ItemStarted(event) => {
let should_emit = match &event.item {
CoreTurnItem::CommandExecution(item) => thread_state
.lock()
.await
.turn_summary
.command_execution_started
.insert(item.id.clone()),
_ => true,
};
if should_emit {
let notification = item_event_to_server_notification(
EventMsg::ItemStarted(event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}
}
EventMsg::ItemCompleted(event) => {
apply_canonical_item_completed_side_effects(
&thread_manager,
&thread_watch_manager,
&thread_state,
&event.item,
)
.await;
let notification = item_event_to_server_notification(
EventMsg::ItemCompleted(event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}command execution 有额外的 started 去重集合,因为 approval/guardian 流程可能先发一次命令开始通知。ItemCompleted 还会执行线程 watch 等副作用,再发客户端通知,所以它不是单纯的结构转换。
5. Turn开始与完成
源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: EventMsg::TurnStarted
EventMsg::TurnStarted(payload) => {
outgoing.abort_pending_server_requests().await;
thread_watch_manager
.note_turn_started(&conversation_id.to_string())
.await;
let turn = {
let state = thread_state.lock().await;
let mut turn = state.active_turn_snapshot().unwrap_or_else(|| Turn {
id: payload.turn_id.clone(),
items: Vec::new(),
items_view: TurnItemsView::NotLoaded,
error: None,
status: TurnStatus::InProgress,
started_at: payload.started_at,
completed_at: None,
duration_ms: None,
});
turn.items.clear();
turn.items_view = TurnItemsView::NotLoaded;
turn
};
outgoing
.send_server_notification(ServerNotification::TurnStarted(
TurnStartedNotification {
thread_id: conversation_id.to_string(),
turn,
},
))
.await;
}TurnStarted 会清理上一轮遗留的 server requests,并以 NotLoaded 视图发布活动 Turn。它不是 item started 的别名。
源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: emit_turn_completed_with_status
async fn emit_turn_completed_with_status(
conversation_id: ThreadId,
event_turn_id: String,
turn_completion_metadata: TurnCompletionMetadata,
outgoing: &ThreadScopedOutgoingMessageSender,
) {
let (items, items_view) = match turn_completion_metadata.last_agent_message {
Some(item) => (vec![item], TurnItemsView::Summary),
None => (Vec::new(), TurnItemsView::NotLoaded),
};
let notification = TurnCompletedNotification {
thread_id: conversation_id.to_string(),
turn: Turn {
id: event_turn_id,
items,
items_view,
error: turn_completion_metadata.error,
status: turn_completion_metadata.status,
started_at: turn_completion_metadata.started_at,
completed_at: turn_completion_metadata.completed_at,
duration_ms: turn_completion_metadata.duration_ms,
},
};
outgoing
.send_server_notification(ServerNotification::TurnCompleted(notification))
.await;
}完成通知默认只携带最后一条 agent message 摘要,因此 items_view=Summary;没有可用摘要时是 NotLoaded。需要完整 item 的客户端必须再读 Thread history,而不能从 completed notification 反推全量历史。
6. 正常完成与中断
源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: handle_turn_complete、handle_turn_interrupted
async fn handle_turn_complete(
conversation_id: ThreadId,
event_turn_id: String,
turn_complete_event: TurnCompleteEvent,
outgoing: &ThreadScopedOutgoingMessageSender,
thread_state: &Arc<Mutex<ThreadState>>,
) {
let turn_summary = find_and_remove_turn_summary(conversation_id, thread_state).await;
let (status, error, last_agent_message) = match turn_summary.last_error {
Some(error) => (TurnStatus::Failed, Some(error), None),
None => (TurnStatus::Completed, None, turn_summary.last_agent_message),
};
emit_turn_completed_with_status(
conversation_id,
event_turn_id,
TurnCompletionMetadata {
status,
error,
last_agent_message,
started_at: turn_summary.started_at,
completed_at: turn_complete_event.completed_at,
duration_ms: turn_complete_event.duration_ms,
},
outgoing,
)
.await;
}
async fn handle_turn_interrupted(
conversation_id: ThreadId,
event_turn_id: String,
turn_aborted_event: TurnAbortedEvent,
outgoing: &ThreadScopedOutgoingMessageSender,
thread_state: &Arc<Mutex<ThreadState>>,
) {
let turn_summary = find_and_remove_turn_summary(conversation_id, thread_state).await;
emit_turn_completed_with_status(
conversation_id,
event_turn_id,
TurnCompletionMetadata {
status: TurnStatus::Interrupted,
error: None,
last_agent_message: None,
started_at: turn_summary.started_at,
completed_at: turn_aborted_event.completed_at,
duration_ms: turn_aborted_event.duration_ms,
},
outgoing,
)
.await;
}正常路径根据 turn summary 的 last error 选择 Completed 或 Failed;中断路径固定为 Interrupted 且不携带 error。两条路径都会清空 turn summary,再通过同一个 completed notification builder 发布。
7. 计划与Token通知
源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: handle_turn_plan_update、handle_token_count_event
async fn handle_turn_plan_update(
conversation_id: ThreadId,
event_turn_id: &str,
plan_update_event: UpdatePlanArgs,
outgoing: &ThreadScopedOutgoingMessageSender,
) {
let notification = TurnPlanUpdatedNotification {
thread_id: conversation_id.to_string(),
turn_id: event_turn_id.to_string(),
explanation: plan_update_event.explanation,
plan: plan_update_event
.plan
.into_iter()
.map(TurnPlanStep::from)
.collect(),
};
outgoing
.send_server_notification(ServerNotification::TurnPlanUpdated(notification))
.await;
}
async fn handle_token_count_event(
conversation_id: ThreadId,
turn_id: String,
token_count_event: TokenCountEvent,
outgoing: &ThreadScopedOutgoingMessageSender,
) {
let TokenCountEvent { info, rate_limits } = token_count_event;
if let Some(token_usage) = info.map(ThreadTokenUsage::from) {
outgoing
.send_server_notification(ServerNotification::ThreadTokenUsageUpdated(
ThreadTokenUsageUpdatedNotification {
thread_id: conversation_id.to_string(),
turn_id,
token_usage,
},
))
.await;
}
if let Some(rate_limits) = rate_limits {
outgoing
.send_server_notification(ServerNotification::AccountRateLimitsUpdated(
AccountRateLimitsUpdatedNotification {
rate_limits: rate_limits.into(),
},
))
.await;
}
}TurnPlanUpdated 是 checklist 全量更新;模型 proposed-plan 的增量是 PlanDelta,属于 item progress。Token usage 和 account rate limits 也分别发布,可能只到达其中一个。
8. 历史与实时投影
源码位置:codex-rs/app-server-protocol/src/protocol/event_mapping.rs :: item_event_to_server_notification
pub fn item_event_to_server_notification(
msg: EventMsg,
thread_id: &str,
turn_id: &str,
) -> ServerNotification {
let thread_id = thread_id.to_string();
let turn_id = turn_id.to_string();
match msg {
EventMsg::AgentMessageContentDelta(event) => {
ServerNotification::AgentMessageDelta(AgentMessageDeltaNotification {
thread_id,
turn_id,
item_id: event.item_id,
delta: event.delta,
})
}
EventMsg::PlanDelta(event) => ServerNotification::PlanDelta(PlanDeltaNotification {
thread_id,
turn_id,
item_id: event.item_id,
delta: event.delta,
}),
EventMsg::ItemStarted(event) => {
ServerNotification::ItemStarted(ItemStartedNotification {
thread_id,
turn_id: event.turn_id,
item: event.item.into(),
started_at_ms: event.started_at_ms,
})
}
EventMsg::ItemCompleted(event) => {
ServerNotification::ItemCompleted(ItemCompletedNotification {
thread_id,
turn_id: event.turn_id,
item: event.item.into(),
completed_at_ms: event.completed_at_ms,
})
}
_ => unreachable!("only item events are accepted by this mapper"),
}
}这个 helper 只覆盖 stateless 的一对一 item 事件;TurnStarted/Completed、approval、usage 和需要跨事件聚合的历史状态由其他 handler 负责。实时 mapper 不能替代 ThreadHistoryBuilder。
9. 失败与取消
Turn 输入可能因 oversized input、inactive thread、steer ID mismatch、不可 steering task 或环境/权限解析失败而拒绝。执行中可能收到 stream error、tool failure 或 interrupt;App Server 会清理与 Turn 绑定的 pending server requests,但不会回滚已经产生的文件、网络或进程副作用。客户端断开则可能只意味着通知丢失,不能据此推断 Core Turn 已停止。
源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: EventMsg::TurnAborted、EventMsg::StreamError
EventMsg::TurnAborted(turn_aborted_event) => {
outgoing.abort_pending_server_requests().await;
respond_to_pending_interrupts(&thread_state, &outgoing).await;
thread_watch_manager
.note_turn_interrupted(&conversation_id.to_string())
.await;
handle_turn_interrupted(
conversation_id,
event_turn_id,
turn_aborted_event,
&outgoing,
&thread_state,
)
.await;
}
EventMsg::StreamError(ev) => {
let turn_error = TurnError {
message: ev.message,
codex_error_info: ev.codex_error_info.map(V2CodexErrorInfo::from),
additional_details: ev.additional_details,
};
outgoing
.send_server_notification(ServerNotification::Error(ErrorNotification {
error: turn_error,
will_retry: true,
thread_id: conversation_id.to_string(),
turn_id: event_turn_id.clone(),
}))
.await;
}StreamError 是可重试的中间错误通知,TurnAborted 才触发 Interrupted completed notification。两者都可能出现在最终 Turn 之前。
10. 测试与边界
源码位置:codex-rs/app-server/tests/suite/v2/turn_start.rs :: turn_start_steers_active_turn_and_returns_active_turn_id、turn_start_rejects_combined_oversized_text_input、turn_start_emits_user_message_item_with_text_elements
源码位置:codex-rs/app-server/tests/suite/v2/turn_steer.rs :: turn_steer_requires_active_turn、turn_steer_returns_active_turn_id、turn_steer_rejects_context_only_input_without_merging_context
源码位置:codex-rs/app-server/tests/suite/v2/turn_interrupt.rs :: turn_interrupt_aborts_running_turn、turn_interrupt_rejects_completed_turn、turn_interrupt_resolves_pending_command_approval_request
源码位置:codex-rs/app-server/tests/suite/v2/plan_item.rs :: plan_mode_uses_proposed_plan_block_for_plan_item、plan_mode_without_proposed_plan_does_not_emit_plan_item
cd codex-rs
cargo test -p codex-app-server turn_start_steers_active_turn_and_returns_active_turn_id -- --nocapture --test-threads=1
cargo test -p codex-app-server turn_start_rejects_combined_oversized_text_input -- --nocapture --test-threads=1
cargo test -p codex-app-server turn_start_emits_user_message_item_with_text_elements -- --nocapture --test-threads=1
cargo test -p codex-app-server turn_steer_requires_active_turn -- --nocapture --test-threads=1
cargo test -p codex-app-server turn_interrupt_rejects_completed_turn -- --nocapture --test-threads=1
cargo test -p codex-app-server plan_mode_uses_proposed_plan_block_for_plan_item -- --nocapture --test-threads=1
cargo test -p codex-app-server plan_mode_without_proposed_plan_does_not_emit_plan_item -- --nocapture --test-threads=1这些测试覆盖输入接受/拒绝、steer 前置条件、interrupt 终态和 plan item 生成;不能证明每个 item 都有完整生命周期、断线后通知仍可达,或所有远程/分页场景都被覆盖。
11. 源码定位练习
遇到“turn/start 成功但没有完成通知”,先保存 response 中的 Turn ID,再检查 TurnStarted、pending server request、item 事件和 TurnComplete/TurnAborted 是否出现。遇到“item 完成但 Turn 仍运行”,这是正常的层级差异,应继续查看是否还有工具或模型采样。
遇到“steer 修改了当前上下文”,检查 expected_turn_id 和 Core steer 分支;遇到“计划被当作 assistant 文本”,沿 parser → PlanDelta → Plan item 完成路径定位。V2 Turn 协议的关键是把整轮状态、局部 item、增量通知和历史视图分开阅读。
