一条用户消息进入 Codex 后,不是直接拼成 HTTP 请求。它先经过 Thread 级设置排序、活动 Turn 判定、 Task 所有权建立、上下文快照、Hook 与技能注入,才生成第一次模型采样。模型返回 final text 也不一定 结束 Turn:工具调用、steering、自动压缩和 stop hook 都可能把它带回下一轮采样。
本节跟踪没有省略工具但先以普通文本结束为主线,说明一条 UserInput 怎样变成可持久、可中止、可继续 的 Regular Turn。
第一步不是创建 Turn,而是判断能否 steer
submission loop 把 Op::UserInput 交给 user_input_or_turn_inner。载荷里的 ThreadSettingsOverrides 必须先
应用,final output JSON Schema 也进入新的 TurnContext。随后 Handler 调 Session steering:
- 存在 Regular 活动 Turn:输入进入当前 Turn 的 pending queue;
- 没有活动 Turn:构造 TurnInput,启动 RegularTask;
- 活动 Task 是 Review/Compact:返回 ActiveTurnNotSteerable;
- expected Turn ID 不符:拒绝,防止把输入送错 Turn。
async def accept_user_input(session, submission):
op = submission.op
await session.apply_thread_settings(op.thread_settings)
context = await session.new_turn_context(
turn_id=submission.id,
output_schema=op.final_output_json_schema,
)
steer = await session.steer_input(op.items)
if steer.ok:
return
if steer.error != NoActiveTurn:
await session.emit_error(submission.id, steer.error)
return
turn_input = TurnInput.user(op.items, op.client_user_message_id)
await session.spawn_task(context, [turn_input], RegularTask())
因此 App Server 在建立更强的 turn/start/turn/steer 公开语义时,需要在桥接层明确选择入口,不能只把
两者都翻成 UserInput 后假设 Core 会替它区分。
start_task 建立 Turn 的运行所有权
spawn_task 会先以 Replaced 原因中止已有 Task、清 connector selection,再进入 start_task。正常
UserInput 已经通过 steer-or-start 排除了活动 Turn,但特殊 Task 复用这套入口时仍需要替换保护。
start_task 依次完成:
- 标记 Turn 开始时间,保存起始 token usage;
- 清空本 Turn 的 Guardian rejection circuit breaker;
- 从 InputQueue 取可归属到该 Turn 的 pending items;
- 创建/取得 ActiveTurn 与 TurnState;
- 运行 Extension turn-start lifecycle;
- 取得 AgentControl execution guard;
- 建立 CancellationToken、done Notify 和后台 Tokio task;
- 把 RunningTask 写入 ActiveTurn。
execution guard 把一个 Task 的整个生命周期计入多 Agent 并发容量,不只是模型请求瞬间。CancellationToken 负责协作式停止,AbortOnDropHandle 则是超时后的强制兜底。
RegularTask 负责 TurnStarted 与外层 pending-input 循环
RegularTask 进入 run 时先发送 TurnStarted,重置 server-reasoning-included 标记,再消费 Session startup
prewarm。预热成功就取得预建 ModelClientSession;不可用则新建;若等待预热时被取消,仍运行输入 Hook
并记录输入,然后无模型采样返回。
RegularTask 外层还有一个循环:一次 run_turn 返回后,如果 InputQueue 仍有 pending input,就用空的
initial input 再调用 run_turn。这处理 Task 正准备结束时到达的 steering,避免输入落在“Task 尚未清除、
内部循环已退出”的窄窗口。
run_turn 先处理压缩和工具依赖
进入第一次采样前,run_turn 做的事情比构造 Prompt 多:
- 若模型切换导致 comp hash 变化或新模型 context window 更小,运行 pre-sampling compact;
- 从用户输入解析显式 MCP server/plugin 需求;
- 等待所需 MCP binding,并 capture 第一个 StepContext;
- 记录 context/world-state 更新并计算 diff 展示根;
- 解析 skill/plugin mentions,注入相关 instructions 和 connector selection;
- 运行 pending SessionStart hooks 与 UserPromptSubmit hooks;
- 把用户输入和 hook additional context 写入历史;
- 记录本 Turn 的已解析配置 analytics。
如果必要 MCP 初始化被取消,源码仍记录已接受输入,再以 TurnAborted 结束。这样中止不会让用户消息从 历史中凭空消失。
StepContext 是一次请求的一致性快照
StepContext 把当前 TurnContext、环境快照、MCP binding、ToolRouter 与 Step ExtensionData 固定在一起。 同一个模型请求看到的工具 spec 和真正执行 Tool Call 的 Router 来自同一 StepContext,避免“请求时工具 A 存在,回包时却用刷新后的 Router B 执行”。
下一次 response 前可以重新 capture StepContext,以吸收 MCP refresh、环境或工具面变化。TurnContext 则跨整个 Turn 保持 model、permission profile、mode 等 Turn 级合同。
Prompt 来自历史快照,不是只含最新消息
采样前 Session clone 当前 History,再按模型支持的 input modalities 生成 Prompt input。Prompt 还包含:
- ToolRouter 当前 model-visible specs;
- 模型是否支持 parallel tool calls;
- base instructions;
- final output JSON Schema;
- Guardian review source 下不同的 strict schema 策略。
def build_prompt(session, step, turn):
history = session.history.clone().for_prompt(turn.model.input_modalities)
return Prompt(
input=history,
tools=step.tool_router.model_visible_specs(),
parallel_tool_calls=turn.model.supports_parallel_tool_calls,
base_instructions=session.base_instructions,
output_schema=turn.final_output_json_schema,
output_schema_strict=not turn.is_guardian_reviewer,
)
同一个 ModelClientSession 在 Turn 内复用,保存 WebSocket 与 sticky routing state;流重试也在这一个
Session 上进行。若 retryable stream error 尚未超出 Provider max retries,Runtime 发 StreamError、退避,
再重新用历史生成 Prompt。ContextWindowExceeded 和 UsageLimitReached 则直接进入专门错误路径。
Response stream 怎样成为历史与界面事件
try_run_sampling_request 消费类型化 ResponseEvent:
- OutputItemAdded 建立 TurnItem Started;
- text/reasoning delta 低延迟发给客户端;
- OutputItemDone 完成 Item、写 ResponseItem 历史,或启动工具 Future;
- RateLimits/usage 更新 Session 状态;
- Completed 发 RawResponseCompleted、记录 token、标记是否需要下一次 response;
- stream 在 Completed 前关闭,视为可重试的断流错误。
文本 final response 会更新 last_agent_message。但这里返回的是一次 sampling request 的结果,不是整个
Turn 终态。
为什么普通 Turn 里有两层循环
run_turn 的内层循环在以下任一条件成立时继续:
- 模型发出 Tool Call,工具 output 已写回历史;
- Provider Completed 明确
end_turn=false; - 活动 Turn 收到 steering/mailbox input;
- stop hook 请求继续,并提供 continuation prompt;
- token/context 状态要求 mid-turn auto compact 后继续。
只有没有 follow-up、stop hook 不阻止、legacy after-agent hook 也未终止时,run_turn 才返回最后消息。
Steering 何时进入模型历史
第一次采样前,初始 input 先记录;因此 can_drain_pending_input 初值取决于 initial input 是否为空。
完成一次采样后才允许从 InputQueue 取新 steering,运行同样的输入 Hook、解析新增 MCP 需求,并重新
capture StepContext。
工具调用或 commentary/reasoning 还会开放当前 Turn 的 mailbox delivery;而一个 final assistant message 可能把 mailbox 延后到下一 Turn。这避免 Agent 已经给出最终答案后又把迟到 mail 悄悄塞进同一个结论。
自动压缩不是独立用户 Turn
后续工作存在且 token limit 达到,run_turn 在当前 Turn 内执行 auto compact。TokenBudget、remote v2、 remote 或 local compact 由 Feature/Provider 决定。压缩后会刷新历史窗口,再继续未完成的工具/steer 逻辑; 不会向客户端伪造第二个 User Turn。
stop hook 可以阻止“看似完成”
当模型不再要求 follow-up,Runtime 才运行 stop hooks。Hook 可以:
- 正常放行;
- 要求 stop;
- 要求 block 并给 continuation fragments。
第三种情况下,Runtime 将 fragments 合成为模型可见消息、记录并发 Item,再继续采样。如果 Hook 说要 continue 却没有 prompt,Runtime 发 Warning 并忽略 block,避免空循环。
Task 结束由统一外壳收口
Task body 返回后,start_task 创建的后台闭包先 flush rollout,再调用 on_task_finished。后者:
- 取走 RunningTask,阻止重复完成;
- 处理仍未消费的 pending input/hook context;
- 计算本 Turn token delta、工具次数、memory citation、proxy/timing 指标;
- 执行 turn-stop 或 turn-abort lifecycle;
- 发送 TurnComplete/TurnAborted;
- 清 Guardian circuit breaker 与 ActiveTurn;
- 触发 thread-idle lifecycle;
- 再次 flush,把 terminal Event 也纳入 durability barrier;
- 若 mailbox 仍有 trigger work,自动启动下一 Turn。
async def finish_task(session, turn, result):
await session.flush_rollout() # 普通 Item
terminal = build_terminal_event(result)
await session.events.send(terminal)
if session.clear_active_turn_if_same(turn):
await session.emit_thread_idle()
await session.flush_rollout() # terminal Event
await session.maybe_start_pending_mailbox_turn()
TurnComplete 的 error 字段来自 TurnContext terminal error;一次采样中发过 Error 后,Task 外壳仍可能正常 形成 TurnComplete,以便客户端收口状态,而不是让 UI 永远停在 running。
中止路径有协作与强制两层
Interrupt 先 cancel Token,等待最多 100ms 让 Task 自行退出,再 abort Tokio handle。随后调用 Task.abort, 必要时写入 interrupted history marker,并在 TurnAborted 前 flush。这样收到终态后立刻重读 rollout 的 客户端能看到一致的中止上下文。
后台 terminal 默认不因普通 Interrupt 被清理;它们由 CleanBackgroundTerminals 或 Session shutdown 管理。
失败矩阵
| 失败点 | 对历史/事件的影响 | Turn 终态 |
|---|---|---|
| Thread setting 约束失败 | 不启动 Task,发 Error | 无新 Turn |
| 必要 MCP 等待被 cancel | 输入仍记录 | TurnAborted |
| pre-sampling compact 非中止错误 | 发 lifecycle error | TurnComplete,通常无最后消息 |
| 无效图片 | 发 BadRequest Error | TurnComplete |
| retryable stream error | StreamError + retry | 可继续 |
| retry 次数耗尽 | Error | TurnComplete |
| stop hook block 有 prompt | 写 continuation,继续 | 尚未结束 |
| 用户 Interrupt | marker + TurnAborted | Aborted |
| terminal rollout flush 失败 | Warning,后台 writer 继续重试 | Event 已发 |
测试应验证跨层顺序
async def test_regular_turn_event_order(thread):
sid = await thread.submit(user_input("explain module"))
events = await collect_until_terminal(thread, sid)
assert events[0].type == "turn_started"
assert events[-1].type == "turn_complete"
async def test_follow_up_response_is_same_turn(fake_model, thread):
fake_model.enqueue(tool_call_response())
fake_model.enqueue(final_text_response("done"))
events = await run(thread)
assert count(events, "turn_started") == 1
assert fake_model.request_count == 2
async def test_terminal_event_is_flushed(thread, store):
await run_one_turn(thread)
persisted = await store.load(thread.id)
assert persisted.last_event.type == "turn_complete"
普通 Turn 的核心不是一次 HTTP 往返,而是一个被 Task 外壳托管、内部可进行多次采样、每一步都能重新 捕获环境并最终以持久终态收口的执行单元。
评论
登录后即可评论