Submission与Op总表
本文承接Codex协议分层总览。默认读者了解 Rust enum、async/await 和 channel,但还没有把 CodexThread、SessionIo、Submission、Op 与 Event 串起来。本文回答一个内部操作从哪里生成 ID、如何进入 submission queue、由哪个 session owner 分派,以及 Started、Steered、Recover、Suspend 和 Shutdown 分别在什么时候结束;不展开具体审批动作、MCP 工具或模型请求字段。
0.150.0 的重要变化是:普通 Turn 输入不再走旧的 Op::UserInput 分支,而是 Op::TurnInput 携带 TurnInputRequest、路由模式和 oneshot reply;恢复使用 Op::RecoverTurn,根 Turn 交接使用 Op::SuspendTurnAndShutdown。因此“提交成功”至少有三个时刻:入队成功、Core 作出路由决定、实际 Turn/task 完成。
建议先读Session与Turn状态理解 active turn,再读Session输入队列理解 pending mailbox。读完本篇后,你应能定位一个 Op 的 ID、队列 owner、状态回复和清理动作,并解释 channel 关闭时为什么没有业务事件。
1. Submission结构
Submission 是 Core SQ 的单项。它包含生成于入队前的 id、待分派的 Op、可选 W3C trace,以及 agent handoff 的 parent/root lineage。它没有旧版 client_user_message_id 字段;用户消息的 client id 现在位于 TurnInput 内部 payload。
源码位置:codex-rs/protocol/src/protocol.rs :: Submission
/// Submission Queue Entry - requests from user
#[derive(Debug)]
pub struct Submission {
/// Unique id for this Submission to correlate with Events
pub id: String,
/// Payload
pub op: Op,
/// Optional W3C trace carrier propagated across async submission handoffs.
pub trace: Option<W3cTraceContext>,
/// Core-provided ID of the parent turn that directly initiated this submission.
///
/// This is only used for inter-agent communication.
pub parent_turn_id: Option<String>,
/// Core-provided ID of the top-level turn that causally initiated this submission.
pub root_turn_id: Option<String>,
}id 是 submission 与 event 的关联键,但不是所有结果都通过 Event 返回:TurnInput、RecoverTurn 和 SuspendTurnAndShutdown 还通过 reply channel 返回结构化状态。parent_turn_id 和 root_turn_id 只描述跨 agent 因果关系,不能拿来代替当前 Turn 的 sub_id。
2. 入队入口
SessionIo::submit 为普通无 reply 操作生成 ID;submit_with_trace 在构造 Submission 后调用统一的 submit_with_id。真正发送前,如果调用者没有带 trace,Core 从当前 tracing span 补一个 W3C carrier;发送端关闭时统一映射为 InternalAgentDied。
源码位置:codex-rs/core/src/session/mod.rs :: SessionIo::submit、submit_with_trace、submit_with_id
pub(crate) async fn submit(&self, op: Op) -> CodexResult<String> {
self.submit_with_trace(
op,
/*trace*/ None,
/*parent_turn_id*/ None,
/*root_turn_id*/ None,
)
.await
}
pub(crate) async fn submit_with_trace(
&self,
op: Op,
trace: Option<W3cTraceContext>,
parent_turn_id: Option<String>,
root_turn_id: Option<String>,
) -> CodexResult<String> {
let id = new_submission_id();
let sub = Submission {
id: id.clone(),
op,
trace,
parent_turn_id,
root_turn_id,
};
self.submit_with_id(sub).await?;
Ok(id)
}
pub(crate) async fn submit_with_id(&self, mut sub: Submission) -> CodexResult<()> {
if sub.trace.is_none() {
sub.trace = current_span_w3c_trace_context();
}
self.tx_sub
.send(sub)
.await
.map_err(|_| CodexErr::InternalAgentDied)?;
Ok(())
}这里的 Ok(id) 只说明 tx_sub.send 接受了消息,不说明 handler 已运行。对于 reply-bearing 操作,调用者还必须等待另一个 channel;这正是普通 submit 与 submit_turn_input 的语义差异。
源码位置:codex-rs/core/src/session/mod.rs :: new_submission_id
pub(crate) fn new_submission_id() -> String {
Uuid::now_v7().to_string()
}UUID v7 让 submission ID 同时具备唯一性和时间排序特征;它仍只是应用层字符串,不能当作数据库 row id 或远端 model response id。
3. Op分类
Op 是内部操作协议,使用 #[non_exhaustive] 保留扩展空间。它的分组由资源 owner 决定:实时会话操作交给 realtime handler,Turn 输入交给 turn_input,审批/响应交给挂起请求,配置与清理直接修改 session services。
源码位置:codex-rs/protocol/src/protocol.rs :: Op
#[derive(Debug)]
#[allow(clippy::large_enum_variant)]
#[non_exhaustive]
pub enum Op {
Interrupt,
CleanBackgroundTerminals,
RealtimeConversationStart(ConversationStartParams),
RealtimeConversationAudio(ConversationAudioParams),
RealtimeConversationText(ConversationTextParams),
RealtimeConversationSpeech(ConversationSpeechParams),
RealtimeConversationClose,
RealtimeConversationListVoices,
TurnInput {
request: Box<TurnInputRequest>,
mode: TurnInputMode,
reply: oneshot::Sender<CodexResult<TurnInputSubmission>>,
},
RecoverTurn {
thread_settings: ThreadSettingsOverrides,
reply: oneshot::Sender<CodexResult<TurnInputSubmission>>,
},
SuspendTurnAndShutdown {
reply: oneshot::Sender<CodexResult<SuspendTurnOutcome>>,
},
ThreadSettings {
thread_settings: ThreadSettingsOverrides,
},
InterAgentCommunication {
communication: InterAgentCommunication,
},
ExecApproval {
id: String,
turn_id: Option<String>,
decision: ReviewDecision,
},
PatchApproval {
id: String,
decision: ReviewDecision,
},
UserInputAnswer {
id: String,
response: RequestUserInputResponse,
},
RequestPermissionsResponse {
id: String,
response: RequestPermissionsResponse,
},
DynamicToolResponse {
id: String,
response: DynamicToolResponse,
},
RefreshMcpServers,
ReloadUserConfig,
Compact,
SetThreadMemoryMode { mode: ThreadMemoryMode },
ThreadRollback { num_turns: u32 },
Review { review_request: ReviewRequest },
ApproveGuardianDeniedAction { event: GuardianAssessmentEvent },
Shutdown,
RunUserShellCommand { command: String },
}TurnInput、RecoverTurn 和 SuspendTurnAndShutdown 的 reply sender 使“路由决定”可以返回给调用者,而无需伪造一个 EQ event。相反,Interrupt、Compact 或 ReloadUserConfig 的结果由 handler 通过事件、状态或后续行为呈现。
4. Session分派
session loop 是唯一顺序消费点。它为每个 Submission 创建 dispatch span,再按 Op variant 调用 handler。下面截取实时操作、Turn 输入、恢复和暂停四个相邻分支;它们展示了异步 handler 与队列顺序的关系。
源码位置:codex-rs/core/src/session/handlers.rs :: submission_loop
pub(super) async fn submission_loop(
sess: Arc<Session>,
config: Arc<Config>,
rx_sub: Receiver<Submission>,
) {
// To break out of this loop, send Op::Shutdown.
let mut shutdown_received = false;
while let Ok(sub) = rx_sub.recv().await {
debug!(?sub, "Submission");
let dispatch_span = submission_dispatch_span(&sub);
let should_exit = async {
match sub.op {
Op::Interrupt => {
interrupt(&sess).await;
false
}
Op::CleanBackgroundTerminals => {
clean_background_terminals(&sess).await;
false
}
Op::RealtimeConversationStart(params) => {
if let Err(err) =
handle_realtime_conversation_start(&sess, sub.id.clone(), params).await
{
sess.send_event_raw(Event {
id: sub.id.clone(),
msg: EventMsg::Error(ErrorEvent {
message: err.to_string(),
codex_error_info: Some(CodexErrorInfo::Other),
}),
})
.await;
}
false
}
Op::TurnInput {
request,
mode,
reply,
} => {
let result = turn_input::handle(&sess, *request, mode, sub.id.clone()).await;
let _ = reply.send(result);
false
}
Op::RecoverTurn {
thread_settings,
reply,
} => {
let result =
turn_input::handle_recovery(&sess, thread_settings, sub.id.clone()).await;
let _ = reply.send(result);
false
}
Op::SuspendTurnAndShutdown { reply } => {
let result =
super::turn_suspension::suspend_turn_and_shutdown(&sess, sub.id.clone())
.await;
let should_exit = matches!(
&result,
Ok(codex_protocol::turn_input::SuspendTurnOutcome::Suspended { .. })
);
let _ = reply.send(result);
should_exit
}
_ => false,
}
}
.instrument(dispatch_span)
.await;
if should_exit {
shutdown_received = true;
break;
}
}Op::TurnInput 的 handler 返回后立刻 send reply,reply 的完成点早于 RegularTask 的采样完成。暂停则相反:只有持久化 flush、writer close 和 stop lifecycle 完成后才返回 Suspended 并退出 loop。
5. Turn输入回复
SessionIo::submit_turn_input 为 TurnInput 创建 oneshot channel,把 request、mode 和 reply 封装进 Submission,然后等待 Core 的 routing decision。调用者取消等待不会撤回已入队的操作;如果 session loop 先退出,reply receiver 会得到 InternalAgentDied。
源码位置:codex-rs/core/src/session/mod.rs :: SessionIo::submit_turn_input、submit_recover_turn
pub(crate) async fn submit_turn_input(
&self,
mut request: TurnInputRequest,
mode: TurnInputMode,
) -> CodexResult<TurnInputSubmission> {
let id = new_submission_id();
let (reply_tx, reply_rx) = oneshot::channel();
let trace = request.trace.take();
self.submit_with_id(Submission {
id,
op: Op::TurnInput {
request: Box::new(request),
mode,
reply: reply_tx,
},
trace,
parent_turn_id: None,
root_turn_id: None,
})
.await?;
reply_rx.await.unwrap_or(Err(CodexErr::InternalAgentDied))
}
pub(crate) async fn submit_recover_turn(
&self,
thread_settings: ThreadSettingsOverrides,
trace: Option<W3cTraceContext>,
turn_id: String,
) -> CodexResult<TurnInputSubmission> {
let (reply_tx, reply_rx) = oneshot::channel();
self.submit_with_id(Submission {
id: turn_id,
op: Op::RecoverTurn {
thread_settings,
reply: reply_tx,
},
trace,
parent_turn_id: None,
root_turn_id: None,
})
.await?;
reply_rx.await.unwrap_or(Err(CodexErr::InternalAgentDied))
}恢复操作故意用已有 turn_id 作为 Submission id,而不是生成新的 ID;这样恢复后的 sampling 继续使用原 Turn 标识。普通 TurnInput 则生成新的 UUID v7,成为新 Turn 的 submission/turn id。
6. 三种输入模式
turn_input::handle 只做模式分派:StartOrSteer 优先尝试 steering,没有 active turn 才创建新 Turn;StartIfIdle 在非 idle 时返回 NotIdle;Steer 要求调用者提供仍然 active 的 expected turn id。
源码位置:codex-rs/core/src/session/turn_input.rs :: handle
pub(super) async fn handle(
session: &Arc<Session>,
request: TurnInputRequest,
mode: TurnInputMode,
submission_id: String,
) -> CodexResult<TurnInputSubmission> {
match mode {
TurnInputMode::StartOrSteer => start_or_steer(session, request, submission_id).await,
TurnInputMode::StartIfIdle => {
start_if_idle(session, request, submission_id, /*is_recovery*/ false).await
}
TurnInputMode::Steer { expected_turn_id } => {
steer(session, request, expected_turn_id, submission_id).await
}
}
}
pub(super) async fn handle_recovery(
session: &Arc<Session>,
thread_settings: ThreadSettingsOverrides,
submission_id: String,
) -> CodexResult<TurnInputSubmission> {
let request = TurnInputRequest::user_input(Vec::new()).with_thread_settings(thread_settings);
start_if_idle(session, request, submission_id, /*is_recovery*/ true).await
}6.1 StartOrSteer
在 start_or_steer 中,settings 先 preview,避免拒绝请求修改 session;然后调用 session.steer_input。成功就 apply steered settings 并返回 Steered;只有 NoActiveTurn 才转入新 Turn 分支,其他 NotSubmittedReason 原样返回。
源码位置:codex-rs/core/src/session/turn_input.rs :: start_or_steer
let settings = PreparedTurnInputSettings::prepare(session, thread_settings, start).await?;
match session
.steer_input(
&mut items,
additional_context.clone(),
/*expected_turn_id*/ None,
settings.required_active_final_output_json_schema(),
client_id.clone(),
responsesapi_client_metadata.clone(),
incoming_root_turn_id,
)
.await
{
Ok(turn_id) => {
settings.apply_steered(session, submission_id).await?;
Ok(TurnInputSubmission::Steered { turn_id })
}
Err(NotSubmittedReason::NoActiveTurn) => {
let turn_context = settings
.apply_started(session, submission_id.clone())
.await?;
if can_start_root_turn
&& !items.is_empty()
&& turn_context
.turn_metadata_state
.can_start_root_turn(&turn_context.session_source)
{
turn_context
.turn_metadata_state
.set_root_turn_id(submission_id.clone());
}
session
.spawn_task(turn_context, task_input, RegularTask::new())
.await;
Ok(TurnInputSubmission::Started {
turn_id: submission_id,
})
}
Err(reason) => Ok(TurnInputSubmission::NotSubmitted { reason }),
}这条路径说明 StartOrSteer 不是“先看 active_turn 再决定”的无锁检查,而是把 steering 尝试和新 Turn 创建串在同一个 session owner 内,减少检查与提交之间的竞态。
6.2 StartIfIdle
start_if_idle 先检查 trigger-turn mailbox,再判断 Plan mode 下的空自动唤醒,然后在 active_turn 锁内预留 idle Turn。settings preview 或 apply_started 失败都会清除 reservation,避免留下看似 active、实际没有 task 的空状态。
源码位置:codex-rs/core/src/session/turn_input.rs :: start_if_idle
let has_user_input = has_nonempty_user_input(&input);
let is_automatic_idle_work = !has_user_input && !is_recovery;
if session.input_queue.has_trigger_turn_mailbox_items().await {
return Ok(TurnInputSubmission::NotSubmitted {
reason: NotSubmittedReason::PendingTriggerTurn,
});
}
if is_automatic_idle_work && session.collaboration_mode().await.mode == ModeKind::Plan {
return Ok(TurnInputSubmission::NotSubmitted {
reason: NotSubmittedReason::PlanMode,
});
}
let turn_state = {
let mut active_turn = session.active_turn.lock().await;
if active_turn.is_some() {
return Ok(TurnInputSubmission::NotSubmitted {
reason: NotSubmittedReason::NotIdle,
});
}
let active_turn = active_turn.get_or_insert_with(ActiveTurn::default);
Arc::clone(&active_turn.turn_state)
};真正创建 context 后,带 user input 的请求会清除 connector selection、记录 user prompt 并把 pending_turn_input 放入 task;空的 automatic work 则只把 response-item input 放入 queue,不添加一条空 user message。
7. 设置与生效时机
PreparedTurnInputSettings 将 thread settings 和 start-only options 分开。prepare 只 preview,不 mutates session;apply_started 在创建新 Turn 前应用 settings 和 schema,并写入 parent/root lineage;apply_steered 只更新后续 Turn 可见的持久设置,当前 active Turn 的 context 不变。
源码位置:codex-rs/core/src/session/turn_input.rs :: PreparedTurnInputSettings
async fn prepare(
session: &Session,
thread_settings: ThreadSettingsOverrides,
start_options: TurnStartOptions,
) -> CodexResult<Self> {
let thread_settings_update = if thread_settings == ThreadSettingsOverrides::default() {
None
} else {
let updates = thread_settings::prepare_update(session, thread_settings).await;
session
.preview_settings(&updates)
.await
.map_err(|error| CodexErr::InvalidRequest(error.to_string()))?;
Some(updates)
};
Ok(Self {
thread_settings_update,
start_options,
})
}
async fn apply_steered(self, session: &Session, submission_id: String) -> CodexResult<()> {
let Some(thread_settings_update) = self.thread_settings_update else {
return Ok(());
};
thread_settings::apply_update(session, submission_id, thread_settings_update)
.await
.map_err(|error| CodexErr::InvalidRequest(error.to_string()))
}因此,一个 Steered 回复并不表示新 settings 已进入当前 sampling;它只表示 steering 成功,持久 thread settings 会影响后续 Turn。NotSubmitted 则保证 settings 和 start options 都未应用。
8. 恢复与暂停
恢复和暂停解决的是不同问题。recover_turn_if_idle 在 idle 时重新启动原 Turn ID,不添加新的 user input;suspend_turn_and_shutdown 只允许 root thread,将活动 regular Turn 停止而不记录 terminal Turn event,以便另一个 worker 用原 ID 恢复。
源码位置:codex-rs/core/src/codex_thread.rs :: recover_turn_if_idle、suspend_turn_and_shutdown
pub async fn recover_turn_if_idle(
&self,
request: RecoverTurnRequest,
) -> CodexResult<StartIfIdleSubmission> {
self.session
.services
.agent_control
.ensure_execution_capacity_for_turn_start(self)
.await?;
let RecoverTurnRequest {
turn_id,
thread_settings,
trace,
} = request;
match self
.io
.submit_recover_turn(thread_settings, trace, turn_id)
.await?
{
TurnInputSubmission::Started { turn_id } => {
Ok(StartIfIdleSubmission::Started { turn_id })
}
TurnInputSubmission::NotSubmitted { reason } => {
Ok(StartIfIdleSubmission::NotSubmitted { reason })
}
TurnInputSubmission::Steered { .. } => {
unreachable!("recovered turn submission cannot steer")
}
}
}暂停在发送 Suspended reply 前要依次 flush、取消 task、清理 pending input、停止 producers、再次 flush、关闭 writer 并发出 ShutdownComplete。如果任一持久化步骤失败,当前 worker 保留所有权,不报告成功。
源码位置:codex-rs/core/src/session/turn_suspension.rs :: suspend_turn_and_shutdown
// Flush before canceling execution so a persistence failure leaves the original turn running.
live_thread.flush().await.map_err(|error| {
CodexErr::Fatal(format!("flush before root turn suspension failed: {error}"))
})?;
task.cancellation_token.cancel();
task.turn_context
.turn_metadata_state
.cancel_git_enrichment_task();
session.input_queue.clear_pending(&turn).await;
handlers::shutdown_session_runtime(session).await;
live_thread.flush().await.map_err(|error| {
CodexErr::Fatal(format!("flush after root turn suspension failed: {error}"))
})?;
live_thread.shutdown().await.map_err(|error| {
CodexErr::Fatal(format!("close suspended root turn writer failed: {error}"))
})?;
handlers::emit_thread_stop_lifecycle(session.as_ref()).await;
session
.deliver_event_raw(Event {
id: submission_id,
msg: EventMsg::ShutdownComplete,
})
.await;
Ok(SuspendTurnOutcome::Suspended { turn_id })9. 中断与关闭
Interrupt 只取消当前 task,不终止后台 terminal;CleanBackgroundTerminals 才关闭长期运行的 unified exec processes。Shutdown 是 session loop 的显式终止操作;submission channel 意外关闭时,loop 仍执行 runtime shutdown、thread-stop lifecycle 和 live thread writer shutdown。
源码位置:codex-rs/core/src/session/handlers.rs :: submission_loop 尾部
if should_exit {
shutdown_received = true;
break;
}
}
// If the submission loop exits because the channel closed without an
// explicit shutdown op, still run session teardown.
if !shutdown_received {
shutdown_session_runtime(&sess).await;
emit_thread_stop_lifecycle(sess.as_ref()).await;
if let Some(live_thread) = sess.live_thread()
&& let Err(err) = live_thread.shutdown().await
{
warn!("failed to shutdown thread persistence after submission channel closed: {err}");
}
}
debug!("Agent loop exited");关闭顺序是资源生命周期的一部分:发送端 drop 只意味着 loop 收不到下一项,不等于当前 Turn、MCP producer 或 rollout writer 已经停止。
10. 测试路径
Turn input 测试直接调用 handle,避免依赖完整模型。活动 Turn 上调用 StartIfIdle,断言 NotSubmitted::NotIdle 且 pending input 为空;Plan mode 的空自动输入返回 PlanMode;trigger mailbox 非空时返回 PendingTriggerTurn 且 mailbox 保留。
源码位置:codex-rs/core/src/session/turn_input_tests.rs :: start_only_rejects_active_turn_without_injecting、start_only_rejects_empty_user_input_in_plan_mode、start_only_rejects_pending_trigger_turn_without_injecting
assert_eq!(
TurnInputSubmission::NotSubmitted {
reason: NotSubmittedReason::NotIdle,
},
submission
);
assert_eq!(
(Vec::<TurnInput>::new(), None, None),
session
.input_queue
.get_pending_input(&session.active_turn)
.await
);Steer-only 测试输入不存在的 expected turn id,断言 NoActiveTurn;输入不同于当前 active turn 的 id,断言 ExpectedTurnMismatch { expected, actual }。非 regular task(Review、Compact)则返回 ActiveTurnNotSteerable,并保留当前 root metadata。
源码位置:codex-rs/core/src/session/turn_input_tests.rs :: steer_only_requires_active_turn、steer_only_enforces_expected_turn_id、rejects_non_regular_turns
assert_eq!(
TurnInputSubmission::NotSubmitted {
reason: NotSubmittedReason::ExpectedTurnMismatch {
expected: "different-turn-id".to_string(),
actual: turn_context.sub_id.clone(),
},
},
submission
);session loop 测试关闭 submission sender,不发送 Op::Shutdown,随后断言 thread lifecycle contributor 只执行一次、thread store 的 shutdown_thread 调用发生;另一测试在活动 Turn 上关闭 channel,断言 abort lifecycle 先于 thread-stop lifecycle。它验证的是兜底清理,不是每个 Op 的业务语义。
源码位置:codex-rs/core/src/session/tests.rs :: submission_loop_channel_close_runs_full_thread_teardown、submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle
cd codex-rs
cargo test -p codex-core --lib start_only_rejects_active_turn_without_injecting -- --nocapture --test-threads=1
cargo test -p codex-core --lib start_only_rejects_empty_user_input_in_plan_mode -- --nocapture --test-threads=1
cargo test -p codex-core --lib steer_only_enforces_expected_turn_id -- --nocapture --test-threads=1
cargo test -p codex-core --lib rejects_non_regular_turns -- --nocapture --test-threads=1
cargo test -p codex-core --lib submission_loop_channel_close_runs_full_thread_teardown -- --nocapture --test-threads=1
cargo test -p codex-core --lib submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle -- --nocapture --test-threads=1这些测试证明路由决定、settings 不变式、Turn 类型限制和 channel-close teardown;不证明模型采样已经完成、远端 provider 收到取消,或后台 OS 进程在所有平台都按同样时序退出。
11. 阅读闭环
遇到“调用返回成功但还没有输出”时,先区分普通 submit 的入队结果和 TurnInputSubmission::Started/Steered 的路由结果,再去看 session task 和 EQ。遇到“恢复创建了新 Turn”时,检查 submit_recover_turn 是否使用原 turn id;遇到“关闭后仍有资源”时,检查 channel-close teardown、live thread writer 和 background terminal owner。
最后任选 Interrupt、TurnInput、RecoverTurn 或 SuspendTurnAndShutdown,写出它的 Submission.id 来源、reply/event 形式、状态 owner、成功时机和失败清理。能把这五项写清楚,就已经真正理解了 Submission 与 Op 在 0.150.0 中的职责边界。
