如果只看快速教程,Agent Loop 往往被写成十几行:
while True:
response = model(history)
if response.has_tool_call:
history.append(run_tool(response.tool_call))
else:
return response.text
这段代码抓住了“采样—行动—观察—再采样”的骨架,但它无法解释一个可恢复 Coding Agent 的真实控制流:新用户输入什么时进历史?工具并发时何时才算一次采样完成?上下文达限是 终止还是压缩?模型说“完成”后 Hook 为什么还能要求它继续?取消与普通错误又为什么不是 同一个终态?
Codex 的当前实现并没有一个名为 AgentLoopState 的 Rust enum。它主要通过 run_turn 中的
异步函数边界、loop/continue/break、Future 收束和 CancellationToken 表达转移。本节要做的是
从这条真实控制流中抽出一台可实现、可测试的最小状态机,而不是声称源码里存在一组
同名状态。
先分清 Task 生命周期与 Agent Loop
用户看到的一个 Turn,外层是 Session Task,内层才是 Agent Loop。两层职责不同:
| 层 | 主要责任 | 关键终点 |
|---|---|---|
| Task wrapper | TurnStarted、CancellationToken、Tokio task、Rollout flush、TurnComplete/TurnAborted | 对客户端发出恰当的 Turn 终态 |
run_turn | 准备上下文、反复采样、执行工具、压缩、Stop Hook 与循环判定 | 返回 last message、普通收尾或 TurnAborted error |
Session 启动 Task 时创建一枚 Turn CancellationToken,发 TurnStarted,再把 Task body 放入异步任务。 Task body 结束后先 flush 已记录的普通项,然后由统一收尾路径计算 token 用量、耗时、差异等状态, 追加 TurnComplete 或 TurnAborted,再次 flush 这个终态事件。
这个分层保护了一个重要性质:Agent Loop 可以内含任意多次模型采样,但整个 Turn 仍然只有 一个对外开始和一个对外终态。如果每次 sampling 都自己发 TurnComplete,Tool Result 触发的第二次 采样就会在客户端看来变成一个新 Turn,从而破坏用户任务的语义边界。
最小状态集不只有“思考”与“行动”
从 Codex 的生产路径抽象,一个 Regular Turn 至少需要下列概念状态:
| 状态 | 入口条件 | 必须保护的事实 |
|---|---|---|
| Preparing | Turn 刚开始或准备下一次采样 | 本次请求的 context、tools 和后续 tool calls 共享一个 StepContext |
| Sampling | Prompt 已构建 | 流式 delta 可对 UI 可见,完整 Item 才入历史与分派 |
| JoiningTools | 本次 Response 启动了 client tools | 下一次采样前所有工具都有结果或明确终止 |
| Deciding | 采样和工具收束 | 用单一谓词整合模型 follow-up 与 pending input |
| Compacting | 需继续且 context 要 rollover | 压缩后仍能续上原 Turn,不把它变成新用户 Turn |
| StopGate | 没有 follow-up | Hook 可放行、强制停止或用明确 continuation prompt 改写为继续 |
| Completed | Loop 正常收尾 | 返回 last assistant message,交给 Task wrapper 发终态 |
| Aborted | Cancellation 穿透到当前阶段 | 不伪装成正常完成,先处理标记/清理/持久化顺序 |
| ErrorClosing | 非取消错误无法继续当前 Turn | 发错误事件,保留 Thread 可接下一个 Turn |
这个状态集仍然是“最小”的,因为它没有把 SSE 的每种 delta、每种 Tool Handler 或每种 Hook 拆成 主状态。拆分标准是:这个阶段是否有独立的等待边界、是否会改变后续控制流,是否需要独立的失败 或取消处理。
Preparing 不只是把字符串拼成 Prompt
首次采样前,run_turn 先处理一组会影响本 Turn 能否开始的工作:
- 判断旧上下文是否应在新输入进来前预压缩;
- 从用户输入中识别显式需要的 MCP server、插件与能力;
- 捕获第一个 StepContext;
- 记录模型可见的 World State/上下文基线;
- 构造 Skills 和 Plugins 注入 Item;
- 执行 session-start/input hooks;
- 记录用户输入、注入 Item 与当前 Turn 设置。
这些步骤中很多都会异步等待,也都可能被取消。例如初始 StepContext 捕获可能需要等环境与
MCP 能力准备。所以 Preparing 必须是真正状态,不能只当成一行 prompt = history 的无失败过渡。
尤其需要注意,“采样还没开始”并不意味着 Turn 还没开始。Task wrapper 早已发 TurnStarted, 也已经建立取消与持久化责任。因此预压缩或能力解析期间被中断,仍必须走 TurnAborted 收尾。
StepContext 是“一次采样的自洽视图”
TurnContext 与 StepContext 的时间尺度不同。TurnContext 跨越整个 Turn,保留模型、配置、
Turn ID、模式、扩展数据和终态错误等回合语义。StepContext 则会在后续采样前重新捕获。
为什么不在整个 Turn 里只建一个 StepContext?因为工具会改变环境:
- Apply Patch 可能改变 Git 状态与文件集;
- Shell 可能生成新文件或改变可发现工作区状态;
- pending user input 可能显式引用新 MCP server 或插件;
- 配置、帐号绑定或远程能力可能在 Turn 运行期间变脏。
但在某一次具体请求内,Context、模型看到的 Tool Spec 和真正处理 Tool Call 的 Router 又必须 共享同一视图。否则模型可能按快照 A 看到工具,Runtime 却按更新后快照 B 去解释调用,造成 无法复现的 TOCTOU 式不一致。
async def capture_request_view(turn, pending_input):
required_servers = await discover_required_servers(pending_input)
step = await session.capture_step_context(
turn_context=turn.context,
required_servers=required_servers,
cancellation=turn.cancellation,
)
# 同一 step 同时驱动 Prompt 与 ToolRouter。
prompt = await build_prompt(turn.history, step)
router = ToolRouter.from_step_context(step)
return RequestView(step=step, prompt=prompt, router=router)
pending input 不是随时插入正在发送的 Prompt
用户可能在模型或工具正在运行时追加指令。最粗暴的实现会直接改写已经打开的模型 stream, 但这会让“这次 Response 基于哪个 Prompt”无法回答。Codex 将运行中输入放进 Input Queue/Mailbox, 在采样边界接收,再于下一次循环写入历史。
主循环有一个 can_drain_pending_input 门。两种情况会故意延迟 drain:
- Turn 刚开始,确保函数入参中的新用户输入先被采样;
- 回合中压缩刚结束,当模型/工具本来就需要续行时,先恢复那条续行链。
对于普通后续循环,pending input 经 Hook 检查后记录。如果它包含新的用户输入,StepContext 捕获前 还会重新解析所需 MCP server。因此 Steer 不只是往 Prompt 末尾多加一句话,它还可能改变下一步的 能力准备集。
Sampling 的出口不是 Provider Completed 事件
1.4 已经解释过,Provider 发 response.completed 只表示模型输出完成。如果 Response 中含客户端
Tool Calls,工具 Future 可能仍在运行。try_run_sampling_request 会继续做三件事:
- 收束 FuturesOrdered,将每个 Tool Result 转成 ResponseItem 并记录;
- 在工具停顿结束后发 token count,避免 request-user-input 类工具等人时 UI 仍显示进展;
- 检查 cancellation,必要时使整次 sampling 返回 TurnAborted,然后才发 TurnDiff。
因此本节将 Sampling 与 JoiningTools 抽象为两个概念状态,但两者在源码中属于同一次
sampling request 的完整函数生命周期。外层 run_turn 不会在半批工具结果到达时进入 Deciding。
async def sampling_state(request_view, turn):
stream_outcome, tool_futures = await consume_provider_stream(
request_view,
cancellation=turn.cancellation.child_token(),
)
outputs = await join_tools_in_call_order(tool_futures)
for output in outputs:
await turn.history.record(output)
await emit_token_count_if_needed(turn)
turn.cancellation.raise_if_cancelled()
await emit_turn_diff_if_needed(turn)
return stream_outcome
Deciding 的核心只有一个继续谓词
工具、pending input、token 和 Hook 都会影响控制流,但 Codex 先把主循环的普通继续条件收敛为:
needs_follow_up = model_needs_follow_up or has_pending_input
model_needs_follow_up 表示当前模型交互还没闭环,可由下列原因产生:
- 出现了 client Tool Call,工具 Output 需要再给模型;
- Provider 明确给
end_turn=false; - 某些流式预占/mailbox 路径要求及时返回外层。
has_pending_input 则表示即使模型认为已可以停止,用户或协作方仍有新信息需要被处理。
用 OR 合并的价值是使循环回边可解释。不需要在主循环的每个角落各写一个 continue。之后的
Compact 与 Stop 判定都建立在这个布尔值上。
Compacting 是回边,不是终态
模型需要继续时,Context Window 可能已经达到 token 上限,也可能有显式新窗口请求。Codex 的 回合中 rollover 条件是:
should_roll_over = needs_follow_up and (
explicit_new_context_window_request
or token_limit_reached
)
前面的 needs_follow_up 门很重要。如果模型已经给出最终答案且没有 pending input,即使 token 很高,
当前 Turn 也不需要为一个不存在的下一次采样支付压缩延迟。压缩是为了使续行可能,不是回合
结束仪式。
压缩成功后,Loop 运行相关 Hook,调整 pending-input drain 门,然后回到 Preparing。下一次采样 使用重建后的历史/上下文,但仍处在同一 TurnContext、同一 Task 和同一 Turn 终态之前。
StopGate 让“模型没有 follow-up”不等于“系统必须完成”
当 needs_follow_up == false时,Agent 看起来已经可以完成。但 Codex 在 break 前还有 Stop Hook 门。
Hook 可以看到 Turn ID、cwd、transcript、model、permission mode、last assistant message 以及自己是否已经
阻止过一次。
最关键的分支是 should_block + continuation_fragments。例如一个 Stop Hook 检查发现 Agent 声称修复完成,
但尚未运行要求的验证,它可以返回一段续行原因。Runtime 会把这段内容编码为新 Prompt Item、
记录和发出 TurnItem lifecycle,将 stop_hook_active 置 true,然后重新采样。
async def stop_gate(turn, last_message):
outcome = await hooks.run_stop(
last_assistant_message=last_message,
stop_hook_active=turn.stop_hook_active,
)
if outcome.should_block:
prompt = build_continuation_message(outcome.continuation_fragments)
if prompt is not None:
await turn.history.record_and_emit(prompt)
turn.stop_hook_active = True
return CONTINUE
emit_warning("Stop hook blocked without a continuation prompt")
if outcome.should_stop:
return COMPLETE
return COMPLETE
“没有 continuation prompt 的 block 被忽略”是一个必要的防空转规则。如果 Hook 只说“不许停”, 却不给模型任何新信息,下一次采样会看到与上次几乎相同的状态,很容易无限重复。 续行必须有可观测的新 Prompt 原因。
终止不是一个布尔值
最小实现往往只有 done=True,但 Codex 需要区分至少三种对外含义:
正常收尾
run_turn 保存最后 assistant message 并返回。Task wrapper 记录指标、发 TurnComplete,清理 ActiveTurn,
必要时开始队列里的下一份工作。
明确取消
Sampling、Tool、Compaction 或 Preparing 中的取消传播为 TurnAborted。外部 interrupt 路径先取消 token,
给 Task 一段宽限时间优雅退出,再强制 abort 剩余 task,并调用 Task-specific abort 收尾。
若原因是 Interrupted,Runtime 还会把模型可见的 <turn_aborted> 上下文标记先写入历史,
flush 后才发 TurnAborted。这样客户端收到终态后立即重读 Rollout,不会看到旧状态。
错误收尾
无法继续的非取消错误会产生 Error lifecycle 和用户可见 Error。影响 Turn 状态的 Error 被写入
turn_context.terminal_error;统一 Task 终态信封可用 TurnComplete 携带它。这个协议细节意味着不能
用事件类型名字直接推断“任务业务成功”;要同时看 terminal error。
错误与取消的详细边界属于 1.6。在状态机层面只需保护:Aborted 是控制流终态,
ErrorClosing 是记录失败并使 Thread 可以继续下一 Turn 的收尾路径,两者不能被一个 except: break
无差别吞掉。
跨状态但不应放入状态 enum 的对象
状态机还需要一组跨转移对象。它们不是状态,而是拥有不同时间尺度的资源:
TurnContext:整 Turn 的语义与配置;CancellationToken:贯穿 Preparing、Sampling、Tools 和 Compacting 的终止信号;ModelClientSession:Turn-scoped,跨 sampling 和 retry 复用 WebSocket 与粘性路由;ContextHistory:保存消息、Call、Output 与压缩交接;TurnDiffTracker:跨多次工具执行累计用户视角的整 Turn 差异;WorldStatebaseline:用于发现下一 Step 相对上一 Step 的改变;last_agent_message:只在拟完成时交给 Stop Hook 与最终 TurnComplete。
如果把这些全部塞进一个巨大 State 并每次转移整体拷贝,资源所有权和生命周期会变得含糊。
更合理的实现是让 TurnRuntime 拥有长寿命资源,状态只描述当前控制位置与该状态局部数据。
Mini Codex 的显式状态机伪代码
虽然 Codex 当前使用异步控制流而非 enum,Mini Codex 可以用显式状态保护转移:
class Phase(Enum):
PREPARING = auto()
SAMPLING = auto()
JOINING_TOOLS = auto()
DECIDING = auto()
COMPACTING = auto()
STOP_GATE = auto()
COMPLETED = auto()
ABORTED = auto()
ERROR_CLOSING = auto()
@dataclass
class TurnRuntime:
turn_context: TurnContext
history: ContextHistory
client_session: ModelClientSession
cancellation: CancellationToken
diff_tracker: TurnDiffTracker
world_state: WorldState
phase: Phase = Phase.PREPARING
stop_hook_active: bool = False
last_agent_message: str | None = None
pending_tool_futures: FuturesOrdered = field(default_factory=FuturesOrdered)
async def run_agent_turn(runtime, initial_input):
await prepare_turn_once(runtime, initial_input)
while runtime.phase not in {
Phase.COMPLETED,
Phase.ABORTED,
Phase.ERROR_CLOSING,
}:
try:
runtime.cancellation.raise_if_cancelled()
if runtime.phase is Phase.PREPARING:
pending = await drain_pending_if_allowed(runtime)
runtime.request_view = await capture_request_view(runtime, pending)
runtime.phase = Phase.SAMPLING
elif runtime.phase is Phase.SAMPLING:
stream_result = await consume_model_stream(
runtime.request_view,
runtime.pending_tool_futures,
runtime.cancellation.child_token(),
)
runtime.stream_result = stream_result
runtime.phase = Phase.JOINING_TOOLS
elif runtime.phase is Phase.JOINING_TOOLS:
await drain_tools_and_record(runtime)
runtime.cancellation.raise_if_cancelled()
runtime.phase = Phase.DECIDING
elif runtime.phase is Phase.DECIDING:
has_pending = await runtime.history.has_pending_input()
needs_follow_up = (
runtime.stream_result.model_needs_follow_up
or has_pending
)
if needs_follow_up and await context_needs_rollover(runtime):
runtime.phase = Phase.COMPACTING
elif needs_follow_up:
runtime.phase = Phase.PREPARING
else:
runtime.last_agent_message = (
runtime.stream_result.last_agent_message
)
runtime.phase = Phase.STOP_GATE
elif runtime.phase is Phase.COMPACTING:
await compact_and_rebuild(runtime)
runtime.phase = Phase.PREPARING
elif runtime.phase is Phase.STOP_GATE:
outcome = await run_stop_hook(runtime)
if outcome.has_continuation_prompt:
await runtime.history.record(outcome.prompt)
runtime.stop_hook_active = True
runtime.phase = Phase.PREPARING
else:
runtime.phase = Phase.COMPLETED
except TurnCancelled:
runtime.phase = Phase.ABORTED
except NonRecoverableTurnError as error:
await record_turn_error(runtime, error)
runtime.phase = Phase.ERROR_CLOSING
return await finish_phase(runtime)
这个版本比生产代码更显式,但保留了实际设计的四个关键不变式:
- 一次采样的工具全部收束后才能 Deciding;
- StepContext 在 Preparing 捕获,而不是整 Turn 只捕获一次;
- Compact 和 Stop Hook continuation 都是回到 Preparing 的显式回边;
- Aborted 与 ErrorClosing 不伪装成 Completed。
转移表比流程图更适合写测试
状态机应该先用转移表列出可测的边:
| 当前状态 | 事件/条件 | 下一状态 | 必须发生的副作用 |
|---|---|---|---|
| Preparing | request view 成功 | Sampling | 保存单次 StepContext |
| Sampling | Provider Completed | JoiningTools | 不得发起新 sampling |
| JoiningTools | 所有 Future 收束 | Deciding | 所有 Output 已入历史 |
| Deciding | follow-up 且 token 未达限 | Preparing | 下次重新捕获 StepContext |
| Deciding | follow-up 且 token 达限 | Compacting | 不先采样 |
| Compacting | 成功 | Preparing | 压缩后历史已安装 |
| Deciding | 无 follow-up | StopGate | 保留 last assistant message |
| StopGate | block + prompt | Preparing | 记录 prompt,置 stop_hook_active |
| StopGate | allow/stop | Completed | 不再发 sampling |
| 任意活动状态 | cancellation | Aborted | 按本阶段的清理合同退出 |
| 任意活动状态 | 不可恢复错误 | ErrorClosing | 记录/发送 Error,Thread 可继续 |
对应测试不应只断言最后有一条文本,而应断言事件顺序与未发生的行为:
async def test_stop_hook_continuation_reenters_preparing():
model.enqueue(final("looks done"))
hooks.stop_returns(block=True, prompt="run tests first")
model.enqueue(final("tests passed"))
hooks.stop_returns(block=False)
await agent.run("fix bug")
assert model.request_count == 2
assert model.requests[1].contains("run tests first")
assert emitted_terminal_events() == ["TurnComplete"]
async def test_rollover_only_happens_when_follow_up_exists():
context.mark_token_limit_reached()
model.enqueue(final("done"))
await agent.run("answer")
assert compact.call_count == 0
async def test_cancel_during_tool_join_never_reaches_stop_gate():
tool.block_forever()
task = spawn(agent.run("act"))
await tool.started()
task.cancel()
await task
assert emitted_terminal_events() == ["TurnAborted"]
assert hooks.stop_call_count == 0
Codex 现有测试还对持久化顺序做了更强的断言:正常完成时,Task body 普通记录先 flush, TurnComplete 追加后再 flush;中断时,还要保证 interrupted marker 在 TurnAborted 前可持久重读。 这说明状态机的终态不只是内存枚举值,也是对外部观测顺序的承诺。
小结:Agent Loop 是一台有屏障、回边和终态的控制器
把 Agent Loop 简化为 while tool_call 会丢掉三类关键信息:
- 屏障:Provider Completed 不等于 sampling 完成,工具收束屏障决定何时可以进入决策;
- 回边:Tool Result、pending input、context rollover 和 Stop Hook continuation 都可以使同一 Turn 再采样;
- 终态:Completed、Aborted 和 ErrorClosing 需要不同的事件、清理与持久化顺序。
Codex 的主循环可以用一个很简洁的谓词表示普通续行,但这个谓词前有严格的采样/工具屏障, 后有 context 与 Stop Hook 门,外面还有 Task wrapper 统一所有终态。了解这些边界后,下一个问题就变得 可精确回答:哪些失败应该成为模型可见的观察,哪些应该在 Runtime 重试,哪些又必须终止当前 Turn。 这是 1.6 要建立的错误边界。
评论
登录后即可评论