my-pi-agent--loop微内核

前言

在现代 Agent 框架中,核心调度循环(ReAct Loop) 是整个系统的中枢神经。

早期大多数框架(包括我们重构前的版本)倾向于将循环硬编码在 Agent 类的方法内部(class Agent: def run(self): while True:)。这种“大泥球”写法导致循环逻辑与类实例的隐式状态(self.xxx)、磁盘文件(Session JSONL)、外围插件强行咬合在一起,既无法单独对状态机进行高并发压测,也无法向外暴露灵活的响应式流。

在对标 Tau (tau_agent/loop.py)Pi (@earendil-works/pi-agent-core) 的架构重塑中,我们做出了一个决定性的动作:将原先内联在 Agent 类中长达 260 行的循环彻底剥离,提炼为 packages/my-agent-core/src/my_agent_core/loop.py 中的独立无状态纯函数生成器 run_agent_loop


一、为什么要把 run_agent_loop 抽离为独立纯函数?

抽离为独立的 async def run_agent_loop(...) 后,系统获得了五大不可替代的工程红利:

1. 彻底消除隐式状态突变(零 self,绝对确定性)

在类方法中,循环随时可能修改 self.xxx,多轮对话后状态满天飞;纯函数微内核完全没有 self,零隐式状态!所有输入项(大模型门面 llm、消息列表 messages、工具注册表 tools、取消令牌 signal)全部通过参数显式传入。输入确定则输出事件流绝对确定,排查调试的确定性拉满。

2. 通用计算发动机:多场景极致复用(Write Once, Run Everywhere)

run_agent_loop 就像一台通用的汽车发动机,不绑定任何特定底盘,能够随处无缝挂载: - 交互式 CLI:配上 Agent 外壳,驱动 Session 树持久化与终端彩色打字机; - Web / API 流式服务:直接挂载到 FastAPI 的 StreamingResponse 或 WebSocket,中间事件流无需中间人直接推送到前端; - 大规模无头评测(Headless Benchmarks):跑 SWE-bench 评测 1000 道题时,无需任何 Session 落盘与复杂插件装配,直接拿微内核裸跑,内存开销降低 80%,极速并发; - Subagent 子代理:轻量级临时执行手脚任务,无需背负庞大的 Agent 实例。

3. 单元测试速度提升 100 倍(零装配、纯内存秒跑)

以前测试循环逻辑,必须先装配一整台 Agent 大机器(Session 树、MemoryStore、Skills、Plugins),单测准备代码冗长且慢;现在在 tests/test_agent_loop_pure.py 中,直接传入 FakeLLM 和内存列表,0.001 秒测完并发工具与动态转向(全套微内核单测在 0.05 秒内全绿)。

4. 动静分离,职责正交(SRP 原则)

  • agent.py(类)只负责静态的结构与持久化(持有 Session 树、长期记忆、扩展装配);
  • loop.py(函数)只负责动态的过程推理与调度(怎么推理、怎么并发调工具、怎么转向、怎么中断)。

5. 将“事件流”提升为系统级的一等公民(First-Class Stream)

大模型的运行本质上是一个基于时间轴的渐进流式过程:Token 在逐字产生、思考块在逐步形成、工具在按需触发并耗时运行。

在现代微内核架构下,事件流就是函数的正常返回值本身(AsyncIterator[Event])!没有任何侧路管子,没有控制权倒置! 调用方掌握 100% 绝对主动权,站在传送带旁边接住事件(async for event in run_agent_loop(...)),随时可以 break 中止,生成器负责自动触发优雅清理与断头自愈。


二、双层嵌套循环架构模型

run_agent_loop 并不是简单的一个 while True:,而是设计为精巧的 “双层嵌套状态机”

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
┌─────────────────────────────────────────────────────────────────────────────┐
│ 【外层循环: while True】(Follow-up 宏观任务接力) │
│ 负责串联宏观的“连续追问任务”。当当前任务完全结束后,检查是否有新的后续任务。 │
│ │
│ ┌─────────────────────────────────────────────────────────────────────┐ │
│ │ 【内层循环: while has_more_tools or pending_messages】 │ │
│ │ (ReAct 微观思考-行动循环 & Steer 即时转向) │ │
│ │ │ │
│ │ 步骤 1. 优先注入即时插话/转向 (Steering 消息) │ │
│ │ 步骤 2. 轮次上限熔断截断检查 (Max Turns Check) │ │
│ │ 步骤 3. 上下文拓扑清洗与 4 级廉价优先压缩 (Context Preparation) │ │
│ │ 步骤 4. 模型前置审查 Hook (BeforeModelCallHook) │ │
│ │ 步骤 5. 委托模型流式推理车间 (_assistant_turn) │ │
│ │ 步骤 6. 异常阻断与断头调用自愈 (Interruption Self-Healing) │ │
│ │ 步骤 7. 委托工具执行流水线 (阶段 1 截断防御 或 阶段 2~6 批执行) │ │
│ │ 步骤 8. 阶段 7 批量优雅熔断判定与轮次闭环 (TurnEnd 结算) │ │
│ │ 步骤 9. 收割内层即时转向消息 (get_steering_messages) │ │
│ └─────────────────────────────────────────────────────────────────────┘ │
│ │
│ 步骤 10. 收割外层宏观追问任务 (get_follow_up_messages) │
│ │
│ 【循环终点】发射终态事件: yield AgentEnd(stop_reason="end_turn") │
└─────────────────────────────────────────────────────────────────────────────┘

三、单轮微观迭代(Iteration)的 9 步标准时序流水线

在内层循环的每次运转中,微内核严格按照以下 9 个时序推进:

步骤 1:即时插话注入(Steering Ingestion)

在轮次开端先检查 pending_messages。若用户在中途发起了紧急插话(如“别删那个文件!”),优先追加进 messages 并发射 MessageStart/End,强行扭转大模型下一步意图。

步骤 2:最大轮次上限拦截(Max Turns Guard)

检查 iteration > effective_max,超限立即发射 AgentEnd(stop_reason="max_iterations") 安全退出,防止死循环无限消耗 Token。

步骤 3:上下文清洗与 4 级廉价优先压缩(Context Preparation)

调用 _provider_context 剔除无正文的残缺异常轮次,并由 ContextManager.prepare 生成满足当前模型窗口的“零污染只读视图”(view)。若触发压缩则广播 ContextCompacted

步骤 4:模型前置审查拦截(BeforeModelCallHook)

安全门禁审查即将发往模型的完整 view,可就地拦截(block=True)或动态改写送给大模型的消息列表。

步骤 5:委托模型流式推理车间(_assistant_turn

专职消费 LLM 底层产出的高阶 StreamEvent 事件流,转译发射 MessageStart、逐字流式打字的 MessageUpdate 与定型的 MessageEnd。若模型产生空响应,安全合成防守性错误消息。

步骤 6:异常阻断与断头自愈(Interruption Self-Healing)

若大模型报错或被协作取消,调用 _synthesize_interrupted_tool_calls 为未完成的 tool_calls 自动补齐中断消息,从根源消除下一次请求 API 400 校验死锁。

步骤 7:工具流水线批处理(_execute_tools_turn

  • 分支 A(阶段 1 截断防御):若命中 Token 超限截断(stop_reason == "length"),断然拒绝执行任何工具,自动生成重试提示;
  • 分支 B(阶段 2~6 批执行):交给专职工具车间执行 Preflight 广播 before_tool_call 改参 并发批执行(ToolExecutionUpdate 实时流式进度) after_tool_call 改写 消息保序发射。

步骤 8:阶段 7 批量优雅熔断与轮次闭环(TurnEnd & Termination)

检查是否有工具声明了 terminate=Trueany() 语义)。若存在,将 has_more_tools 关停,保全 final_text,优雅结束 ReAct 循环。随后发射 TurnEnd 封闭本轮。

步骤 9:收割内层即时转向消息(Steering Harvesting)

调用 get_steering_messages()。若在模型思考或工具执行期间用户插入了新的转向指令,立即收割注入 pending_messages,驱动内层循环立即开启下一轮应对。


四、两级动态干预队列:Steering(插队转向) vs. Follow-up(排队接力)

run_agent_loop 内置了对两级消息队列的精密调度:

1
2
3
4
5
6
7
8
9
10
11
12
# 1. 内层收割即时转向
if get_steering_messages is not None:
steer_msgs = get_steering_messages()
if steer_msgs:
pending_messages = _as_messages(steer_msgs)

# 2. 外层收割宏观追问
if get_follow_up_messages is not None:
followups = get_follow_up_messages()
if followups:
pending_messages = _as_messages(followups)
continue # 重启外层循环!
队列类别 触发方法 收割时序 核心语义 典型应用场景
Steering(即时转向) agent.steer("不要改了") 内层工具执行完毕后立刻收割 “抢占式插话”:直接插在模型刚拿到的工具结果后面,强行扭转模型下一次推理决策 用户在中途纠偏、阻止破坏性操作、注入急迫上下文
Follow-up(宏观追问) agent.follow_up("写完后发通知") 内层 ReAct 循环完全自然闭环后收割 “任务队列接力”:当前任务完全大功告成后,开启下一个全新的宏观业务轮次 链式长任务排队、后台轮询通知、自动化流水线接力

五、两大专职子生成器车间的分工与协作

为了让主循环保持极其轻盈的百余行代码,loop.py 将两大核心脏活分别下沉给两个专职子生成器:

1
2
3
4
5
6
7
8
9
10
11
12
13
               run_agent_loop (总调度指挥台)

┌───────────┴───────────┐
▼ ▼
【车间 ①: _assistant_turn】 【车间 ②: _execute_tools_turn】
(大模型流式推理车间) (工业级七阶段工具执行流水线)
• 消费 astream_events 事件流 • 阶段 1: 截断防御挂起 (_fail_tool_calls_from_truncated_message)
• 逐字 yield MessageUpdate • 阶段 2: 畸形参数归一化防崩 (_coerce_tool_call)
• 发射 MessageEnd 终态定型 • 阶段 3: Preflight 广播 (ToolExecutionStart)
• 异常/取消时优雅闭环 • 阶段 4: 前置门禁与改参 (before_tool_call)
• 阶段 5: 跨线程异步队列实时流式进度 (ToolExecutionUpdate)
• 阶段 6: 后置改写与单工具终态 (ToolExecutionEnd)
• 阶段 7: 消息保序归档与 any 熔断退出 (MessageEnd & terminate)

1. 模型推理车间:_assistant_turn 的流式桥接

  • 协议解耦:优先调用 llm.astream_events(...) 消费高阶流式事件;若 Provider 不支持事件流,自动退化为基础的 astream(...) 并接入 StreamAccumulator 累加器进行状态机还原;
  • 打字机驱动:遇到 StreamStartEvent 时发射 MessageStart(assistant),遇到 TextDeltaEvent / ThinkingDeltaEvent 时实时发射 MessageUpdate 驱动终端打字机;
  • 协作取消:在消费每个 Chunk 时校验 signal.is_cancelled(),若检测到取消则标记 stop_reason="cancelled" 优雅闭环;
  • 防守性兜底:若上游发生网络崩溃或产出空响应,自动生成防守型合成消息,绝不把未初始化的 assistant 变量遗留给下游。

2. 工具执行车间:_execute_tools_turn 的七阶段流水线

  • 输入容错(_coerce_tool_call:对大模型吐出的任何畸形参数字典进行强制包裹,反序列化报错时自动生成 ToolCall(error="..."),由阶段 2 转化为结构化错误供大模型自愈,绝不让生成器崩溃;
  • 截断保护:当 stop_reason == "length" 时由 _fail_tool_calls_from_truncated_message 严防死守,安全挂起截断的危险工具;
  • 跨线程单向传送带(asyncio.Queue:针对 asyncio.to_thread 执行的同步工具,利用 loop.call_soon_threadsafe 跨线程安全推送 ToolExecutionUpdate,并在 finally 中发射 _SENTINEL 独一无二哨兵,彻底破解并发工具与流式生成器之间的死锁困境;
  • 锁存器丢弃迟到更新(accepting_updates:工具执行完毕(settle)后立即关门,静默丢弃任何残留子线程的延时更新;
  • 单批提前退出(any() 语义):当 Human-in-the-loop 或终止型工具返回 terminate=True 时,整批工具安全执行完毕后立即终结 ReAct 循环,并完整保留 final_text 作为最终产出。

六、现代微内核架构设计精髓总结

  1. 分层分治,彻底消灭上帝类:状态存储交给 agent.py,模型流转交给 _assistant_turn,工具并发交给 _execute_tools_turn,主状态机只用约 110 行代码统筹大局;
  2. 纯粹无状态与一等公民事件流run_agent_loopself 隐式状态,不读写磁盘,所有生命周期状态通过 yield Event 抛出,调用者拥有 100% 消费控制权;
  3. Never-Throw 与自愈闭环:全链路各层异常皆被捕获转化为可自愈的模型反馈或中断结果,绝不发生主进程崩溃;
  4. 统一扁平消息模型:采用单一 Message + metadata,彻底消除多态继承的序列化摩擦,对 JSONL 持久化与 Provider 协议 1:1 零成本适配。