Skip to content

V2 Turn与Item协议

沿着 turn/start、steer、interrupt、item 生命周期和 completed 通知,理解 V2 Turn 的状态与投影边界。

基于rust-v0.150.0
CodexRustAppServerTurn

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

rust
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

rust
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

rust
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(&params.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 的输入提交路径

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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

text
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、增量通知和历史视图分开阅读。