雨天小六

读懂 Codex(2.9):普通 User Turn 的端到端调用链

· 更新于 2026-08-02 · 专栏:读懂 Codex

#Codex#Agent Runtime#软件架构#Turn#Agent Loop

一条用户消息进入 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 依次完成:

  1. 标记 Turn 开始时间,保存起始 token usage;
  2. 清空本 Turn 的 Guardian rejection circuit breaker;
  3. 从 InputQueue 取可归属到该 Turn 的 pending items;
  4. 创建/取得 ActiveTurn 与 TurnState;
  5. 运行 Extension turn-start lifecycle;
  6. 取得 AgentControl execution guard;
  7. 建立 CancellationToken、done Notify 和后台 Tokio task;
  8. 把 RunningTask 写入 ActiveTurn。

execution guard 把一个 Task 的整个生命周期计入多 Agent 并发容量,不只是模型请求瞬间。CancellationToken 负责协作式停止,AbortOnDropHandle 则是超时后的强制兜底。

UserInput 从 Submission 分派、steer 判定、RegularTask 启动、TurnStarted、模型采样到 TurnComplete 的端到端时序
图 2.9-1:提交、Task 启动、模型 response 和 Turn 完成是四个不同阶段。

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 里有两层循环

RegularTask 外层 pending input 循环与 run_turn 内层模型工具自动压缩循环
图 2.9-2:内层循环处理一次 Turn 中的多次模型 response;外层循环收口临近完成时到达的 pending input。

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。后者:

  1. 取走 RunningTask,阻止重复完成;
  2. 处理仍未消费的 pending input/hook context;
  3. 计算本 Turn token delta、工具次数、memory citation、proxy/timing 指标;
  4. 执行 turn-stop 或 turn-abort lifecycle;
  5. 发送 TurnComplete/TurnAborted;
  6. 清 Guardian circuit breaker 与 ActiveTurn;
  7. 触发 thread-idle lifecycle;
  8. 再次 flush,把 terminal Event 也纳入 durability barrier;
  9. 若 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 errorTurnComplete,通常无最后消息
无效图片发 BadRequest ErrorTurnComplete
retryable stream errorStreamError + retry可继续
retry 次数耗尽ErrorTurnComplete
stop hook block 有 prompt写 continuation,继续尚未结束
用户 Interruptmarker + TurnAbortedAborted
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 外壳托管、内部可进行多次采样、每一步都能重新 捕获环境并最终以持久终态收口的执行单元。

评论


← 返回文章列表