1. 模块概述
AgentLoop 和 AgentRunner 是 nanobot 最核心的两个类,共同构成了 Agent 的处理引擎:
- **AgentLoop**(
nanobot/agent/loop.py):产品层面的协调器,管理会话、状态机、消息路由、记忆巩固、MCP 连接等 - **AgentRunner**(
nanobot/agent/runner.py):纯执行引擎,负责”消息→LLM→工具→LLM→…”的迭代循环,不包含产品层逻辑
📍 源码:
nanobot/nanobot/agent/loop.py#L140、nanobot/nanobot/agent/runner.py#L132
1.1 职责边界
| 关注点 | AgentLoop | AgentRunner |
|---|---|---|
| 会话管理 | ✅ | ❌ |
| 状态机驱动 | ✅ | ❌ |
| LLM 调用循环 | ❌ | ✅ |
| 工具执行编排 | ❌ | ✅ |
| 上下文窗口治理 | ❌ | ✅ |
| 记忆巩固触发 | ✅ | ❌ |
| MCP 连接管理 | ✅ | ❌ |
| 流式输出回调 | ✅ | ✅ (透传) |
| 错误恢复/重试 | ❌ | ✅ (通过 Provider) |
1.2 在系统中的位置
graph TB
MB[MessageBus] -->|consume_inbound| AL[AgentLoop]
AL -->|state machine| AL
AL -->|_run_agent_loop| AR[AgentRunner]
AR -->|chat_stream| PR[LLMProvider]
AR -->|execute| TR[ToolRegistry]
AL --> SM[SessionManager]
AL --> CS[Consolidator]
AL --> MS[MemoryStore]
2. 命名体系与易混淆函数对比
2.1 命名规律拆解
| 前缀 | 模块/对象 | 操作 | 含义 |
|---|---|---|---|
_state_ |
restore |
— | 状态机处理器:恢复会话和检查点 |
_state_ |
compact |
— | 状态机处理器:自动压缩检查 |
_state_ |
command |
— | 状态机处理器:命令分发 |
_state_ |
build |
— | 状态机处理器:构建上下文 |
_state_ |
run |
— | 状态机处理器:执行 LLM 循环 |
_state_ |
save |
— | 状态机处理器:持久化会话 |
_state_ |
respond |
— | 状态机处理器:组装响应 |
_run_ |
agent_loop |
— | 调用 AgentRunner 执行核心迭代 |
_run_ |
core |
— | AgentRunner 的核心迭代循环 |
_process_ |
message |
— | 处理单条消息的完整生命周期 |
_process_ |
system_message |
— | 处理系统级消息(子 Agent 结果等) |
_dispatch |
— | — | 消息分发入口,管理会话锁和并发 |
_save_ |
turn |
— | 将本轮消息持久化到会话 |
_build_ |
initial_messages |
— | 为 LLM 调用构建初始消息列表 |
_build_ |
bus_progress_callback |
— | 构建进度回调发布到消息总线 |
_persist_ |
user_message_early |
— | 在 turn 开始前提前持久化用户消息 |
_try_ |
drain_injections |
— | AgentRunner:尝试排空中途注入消息 |
_try_ |
finalize_after_max_iterations |
— | AgentRunner:迭代上限后的最终化尝试 |
命名规律总结:
_state_xxx方法对应状态机的 7 个处理阶段,每个返回事件字符串驱动状态转换_run_xxx是实际执行 LLM 调用的方法,_run_agent_loop(AgentLoop)委托给_run_core(AgentRunner)_try_xxx方法表示可能失败/为空的尝试性操作,失败时降级而非崩溃
2.2 易混淆函数对比表
| 对比维度 | _run_agent_loop (AgentLoop) |
_run_core (AgentRunner) |
|---|---|---|
| 操作对象 | TurnContext(含会话、历史等产品上下文) | 纯消息列表 + AgentRunSpec |
| 调用者 | _state_run 和 _process_system_message |
AgentRunner.run() |
| 调用阶段 | 状态机 RUN 阶段 | AgentRunner 入口 |
| 附加逻辑 | 注入回调、检查点、工作区作用域绑定、流式回调 | 上下文治理(orphan/backfill/snip/microcompact) |
| 一句话区分 | _run_agent_loop 是产品层包装,绑定上下文和回调后委托给 AgentRunner |
_run_core 是纯执行核心,只管消息→LLM→工具的迭代 |
| 对比维度 | _dispatch (AgentLoop) |
_process_message (AgentLoop) |
|---|---|---|
| 操作对象 | 原始 InboundMessage | 会话锁内的消息 |
| 调用者 | run() 主循环 |
_dispatch |
| 调用阶段 | 消息从总线消费后立即 | 获取会话锁后 |
| 附加逻辑 | 会话锁管理、并发门控、待处理队列路由、流式回调构建 | 状态机驱动、错误恢复 |
| 一句话区分 | _dispatch 是并发控制层,负责锁和排队 |
_process_message 是业务处理层,负责状态机流转 |
3. API Signatures
graph TB
subgraph "AgentLoop 公共接口"
run["run()"]
stop["stop()"]
process_direct["process_direct()"]
submit_cron_turn["submit_cron_turn()"]
set_model_preset["set_model_preset()"]
from_config["from_config() [classmethod]"]
end
subgraph "AgentLoop 内部调用链"
dispatch["_dispatch()"]
process_message["_process_message()"]
process_system["_process_system_message()"]
run_agent_loop["_run_agent_loop()"]
state_handlers["_state_restore/compact/command/build/run/save/respond()"]
end
subgraph "AgentRunner 公共接口"
runner_run["run(spec) -> AgentRunResult"]
end
subgraph "AgentRunner 内部调用链"
run_core["_run_core()"]
request_model["_request_model()"]
execute_tools["_execute_tools()"]
snip_history["_snip_history()"]
microcompact["_microcompact()"]
drop_orphan["_drop_orphan_tool_results()"]
backfill["_backfill_missing_tool_results()"]
end
run --> dispatch
dispatch --> process_message
process_message --> state_handlers
state_handlers --> run_agent_loop
run_agent_loop --> runner_run
runner_run --> run_core
run_core --> request_model
run_core --> execute_tools
run_core --> snip_history
AgentLoop
1 | class AgentLoop: |
AgentRunner
1 | class AgentRunner: |
4. 数据结构深度解析
4.1 TurnContext
4.a 结构体存在的理由
TurnContext 是 AgentLoop 状态机各阶段之间共享状态的载体。如果没有它,每个状态处理器需要独立的参数传递和返回值约定,状态之间的数据流转将变得脆弱且难以扩展。它将一个消息处理的完整生命周期中的所有临时数据集中管理,使状态机各阶段只需读取/写入 context 字段即可协作。
4.b 结构定义
📍 源码:
nanobot/nanobot/agent/loop.py#L99-L137
1 |
|
4.c 字段三层分析表
| 字段 | 设计动机 | 反事实 | 替代方案 |
|---|---|---|---|
msg |
状态机各阶段都需要访问原始消息内容 | 各阶段需要独立传参,重复且易出错 | 可以用全局变量,但破坏了无状态设计 |
session_key |
唯一标识会话,用于锁、队列、持久化 | 每次需要时从 msg 重新计算,浪费且有 bug 风险 | 可以存在 session 对象上,但 session 可能为 None |
state |
跟踪当前状态机位置 | 无法知道当前处于哪个阶段 | 可以用方法调用栈隐式表示,但失去了显式追踪 |
turn_id |
日志追踪和调试的唯一标识 | 日志中无法关联同一 turn 的各阶段 | 可以用 (session_key, timestamp),但不保证唯一 |
trace |
记录每个状态的耗时,用于性能分析 | 无法定位性能瓶颈 | 可以用外部 profiler,但侵入性更强 |
ephemeral |
标记临时执行(如 Dream),不持久化会话 | 临时任务污染正常会话历史 | 可以用单独的 AgentLoop 实例,但资源开销大 |
pending_queue |
支持中途消息注入(子Agent完成通知等) | 子Agent完成后无法通知主Agent继续 | 可以用回调,但队列模式与消息总线一致 |
suppress_response |
某些工具(如 message 工具)需要抑制自动回复 | 用户看到重复或多余的回复 | 可以在上层判断,但需要传递额外状态 |
4.d 生命周期状态图
stateDiagram-v2
[*] --> Created: _process_message() 创建
Created --> RESTORE: 状态机启动
RESTORE --> COMPACT: event="ok"
COMPACT --> COMMAND: event="ok"
COMMAND --> BUILD: event="dispatch"
COMMAND --> Destroyed: event="shortcut"
BUILD --> RUN: event="ok"
RUN --> SAVE: event="ok"
SAVE --> RESPOND: event="ok"
RESPOND --> Destroyed: event="ok"
Destroyed --> [*]: ctx.outbound 已设置
4.2 AgentRunSpec
4.a 结构体存在的理由
AgentRunSpec 将一次 Agent 运行的所有配置参数集中在一个 dataclass 中,避免了 AgentRunner.run() 拥有 20+ 个参数。它也使得添加新配置项不需要修改方法签名,只需在 dataclass 中增加字段即可。
4.b 字段三层分析表
| 字段 | 设计动机 | 反事实 | 替代方案 |
|---|---|---|---|
injection_callback |
支持 LLM 循环中途注入新消息(子Agent结果、目标续推) | 子Agent完成后无法通知主循环 | 可以用 polling,但延迟高 |
checkpoint_callback |
在工具执行间隙持久化进度,支持 /stop 后恢复 | /stop 后丢失所有进度 | 可以每次工具完成后全量保存,但 I/O 开销大 |
goal_active_predicate |
判断持续目标是否活跃,影响超时策略 | 持续目标被普通超时杀死 | 可以用全局标志,但无法按 session 区分 |
finalize_on_max_iterations |
控制达到迭代上限时是否尝试无工具最终化 | 迭代上限后只能返回固定错误消息 | 可以始终最终化,但某些场景不需要 |
context_block_limit |
替代基于 token 的上下文窗口限制 | 无法适应不同模型的 block 限制 | 可以只用 context_window_tokens,但某些模型有 block 数限制 |
5. 函数逐行精讲
5.1 AgentLoop.run() — 主循环
5.a 场景卡片
函数:
AgentLoop.run()
- 调用时机:网关启动后,作为长期运行的 asyncio 任务调用
- 典型调用者:
nanobot/cli/commands.py的 gateway 命令- 前置条件:MessageBus 已创建、Provider 已初始化、MCP 已连接
- 目的:持续从 MessageBus 消费消息,为每条消息创建独立的 asyncio.Task
5.b 逐行注释式精讲
1 | async def run(self) -> None: |
5.2 AgentLoop._dispatch() — 消息分发
5.a 场景卡片
函数:
AgentLoop._dispatch()
- 调用时机:主循环为每条非注入消息创建 asyncio.Task 时
- 典型调用者:
run()主循环- 前置条件:消息已从总线消费,已通过优先命令和路由检查
- 目的:管理会话锁和并发门控,然后调用
_process_message执行实际处理
关键逻辑:
- 获取会话级别的
asyncio.Lock,确保同一会话的消息串行处理 - 获取全局并发信号量(默认 3),限制同时处理的会话数
- 为需要流式输出的消息构建
on_stream/on_stream_end回调 - 处理完成后清理待处理队列、发布
turn_completed事件
5.3 AgentRunner._run_core() — 核心迭代循环
5.a 场景卡片
函数:
AgentRunner._run_core()
- 调用时机:AgentRunner.run() 在 before_run hook 之后
- 典型调用者:
AgentRunner.run()- 前置条件:AgentRunSpec 已完整填充,hook 已就绪
- 目的:执行”LLM调用→工具执行→LLM调用→…”的迭代循环,直到获得最终回复或达到迭代上限
5.b 逐行注释式精讲
1 | async def _run_core(self, spec, hook, messages) -> AgentRunResult: |
6. 关键算法剖析
6.1 上下文窗口治理(Context Governance)
每次 LLM 调用前,AgentRunner 对消息列表执行 5 步治理流水线:
1 | 原始消息列表 |
时间复杂度:O(n) 每步,总体 O(5n),n 为消息数。
边界处理:
- 裁剪后确保以 user 消息开头(Provider 要求)
- 如果找不到 user 消息,插入合成
"(conversation continued)"消息 - 保留 system 消息不受裁剪影响
6.2 工具分区并发执行
_partition_tool_batches() 将工具调用列表分为可并发组和必须串行组:
1 | def _partition_tool_batches(self, spec, tool_calls): |
7. 设计决策分析
7.1 为什么 AgentLoop 和 AgentRunner 是两个类?
决策:将产品层关注点(会话、记忆、状态机、MCP)与纯执行逻辑(LLM 调用、工具执行、上下文治理)分离。
权衡:
- ✅ AgentRunner 可独立测试,不需要完整的 MessageBus/Channel 环境
- ✅ 产品层逻辑变更(如增加新的状态机阶段)不影响执行核心
- ✅ 执行核心的优化(如上下文治理策略)不影响产品层
- ❌ 增加了一层间接调用,调试时需要跨类追踪
替代方案:合并为一个类。这会使类超过 2000 行,难以维护和测试。
7.2 为什么使用状态机而不是简单的顺序调用?
决策:使用 8 状态状态机 (TurnState 枚举 + _TRANSITIONS 表) 驱动消息处理。
权衡:
- ✅ 每个状态是独立方法,职责单一
- ✅ 状态转换显式声明,易于追踪和调试
- ✅ 可以独立测试每个状态处理器
- ✅
StateTraceEntry记录每个状态的耗时,支持性能分析 - ❌ 简单场景(如纯文本回复)也要走完整状态链
替代方案:线性调用链。但这会使错误处理和条件分支难以管理。
7.3 为什么使用 checkpoint 机制而不是每次工具调用后全量保存?
决策:在工具执行间隙通过 checkpoint_callback 保存最小状态(assistant_message + tool_results),而不是每次全量保存会话。
权衡:
- ✅ 减少 I/O 开销(只保存增量)
- ✅ /stop 取消后可以恢复部分进度
- ✅ checkpoint 保存在 session metadata 中,不污染消息历史
- ❌ 恢复时需要合并逻辑(
_restore_runtime_checkpoint) - ❌ 如果 checkpoint 与实际消息历史不同步,需要 overlap 检测
8. 学习检查点
📝 本章小结
- AgentLoop 是产品协调器,AgentRunner 是纯执行引擎,两者职责分明
- 8 状态状态机覆盖了从消息恢复到最终响应的完整生命周期
- 上下文治理流水线(5 步)确保每次 LLM 调用前的消息列表干净合法
- 工具分区并发执行根据
concurrency_safe属性自动分组 - checkpoint + injection 机制支持中途恢复和中途消息注入
🤔 思考题
如果
_snip_history裁剪后第一条非 system 消息是 assistant(不含 tool_calls),会发生什么?为什么需要_enforce_role_alternation中的 safety net?参考答案
某些 Provider(如 GLM/Zhipu)会拒绝
system → assistant的消息序列(错误 1214)。_enforce_role_alternation在nanobot/nanobot/providers/base.py#L548-L552检测到这种情况时,会插入一条合成 user 消息"(conversation continued)",确保消息序列合法。这是 Provider 兼容性的最后防线。为什么
_try_drain_injections需要在工具执行后 和 最终回复后两个位置调用?只在一个位置调用会有什么问题?参考答案
工具执行后的注入(
nanobot/nanobot/agent/runner.py#L505-L508)处理子Agent完成等异步事件——这些事件可能在 LLM 思考期间到达,需要在下一轮 LLM 调用前注入。最终回复后的注入(nanobot/nanobot/agent/runner.py#L586-L591)处理持续目标续推——即使 LLM 已经生成了回复,如果有活跃的 sustained goal,应该继续工作而不是返回给用户。如果只在工具执行后检查,最终回复后会丢失目标续推;如果只在最终回复后检查,工具执行期间到达的子Agent结果会被延迟一轮。_microcompact为什么只保留最近 10 条 compactable 工具结果?这个数字(_MICROCOMPACT_KEEP_RECENT=10)的选择有什么考量?参考答案
10 条是平衡上下文质量和 LLM 记忆的折中选择(
nanobot/nanobot/agent/runner.py#L70)。太少的保留(如 3 条)会导致 LLM 快速”忘记”最近的操作,可能在后续迭代中重复读取同一文件。太多的保留(如 50 条)则压缩效果有限,上下文窗口仍会被大量旧工具结果占据。10 条确保 LLM 能看到最近约 2-3 轮工具调用的完整结果,同时将更早的结果压缩为一行摘要,释放上下文空间给新的推理。