my-pi-agent--动态干预机制
followUp与steering
| 维度 | Steering(紧急插队) | FollowUp(后续追加) |
|---|---|---|
| 检查时机 | 内层循环每轮 Turn
结束后立刻拉取 |
内层循环全部结束(Agent 想停时)才拉取 |
| 所在层级 | 内层循环的核心驱动力之一 | 外层循环的续命重启机制 |
| 是否打断当前思路 | 会打断 Agent
既定的多轮工具链计划 |
不会打断,让 Agent
完整走完当前长链推理 |
| 对 Turn 数量的影响 | 在当前任务中途插入一个新
Turn |
在整个任务全部做完后,再追加若干个新
Turn |
| 生命周期关系 | 运行在同一个 Trace 内 |
仍然运行在同一个 Trace
内(共享 new_messages) |
| 消费模式 | 支持 steeringMode
(one-at-a-time / all) |
支持 followUpMode
(one-at-a-time / all) |
| 典型触发者 | 终端用户(发现 Agent
跑偏,中途纠正) |
系统/工作流(主任务做完自动追加质检/测试) |
一个turn什么时候算是结束呢
在 Agent Loop 的严谨架构中,一个 Turn(轮次) 的结束有一个非常精确的定义:
核心结论:当一次模型调用以及这次调用触发的所有工具全部执行完毕,并正式派发了 emit(“turn_end”) 事件时,这个 Turn 就算正式结束了。
一个 Turn 始终由 turn_start 和 turn_end 成对包裹,其内部经历了以下 3 个确定阶段:
1
2
3
4
5
6
7
8
9┌── [Turn 开始] ──> emit("turn_start")
│
│ 阶段 ①【调模型】:发请求给 LLM,接收流式 Token 直到生成完毕(拿到完整 AssistantMessage)
│
│ 阶段 ②【执行工具】:
│ ├── 如果模型输出了 ToolCalls ──> 等待这批工具全部执行完成,产出ToolResultMessage
│ └── 如果模型没有输出 ToolCalls ──> 跳过工具执行
│
└── [Turn 结束] ──> emit("turn_end", { message, toolResults })
视觉全景:一个 Loop 包含多个 Turn
1 | 【用户按回车】──> emit("agent_start") ───【整个 Agent Loop 开始】 |
agentloop双重循环
| 循环层级 | 负责的事情 | 本质角色 | 驱动信号 |
|---|---|---|---|
| 内层循环 ( Inner Loop) |
做完当前这一个具体任务 (例如:定位 Bug 并修复代码) |
微观执行引擎 (ReAct 循环) |
hasMoreToolCalls(还要调工具)或 steering(中途插队消息) |
| 外层循环 ( Outer Loop) |
当前任务彻底完成后,要不要接单下一个追加任务? (例如:“修完了?那顺便跑个测试吧”) |
宏观调度引擎 (任务续命) |
followUp(后置追加任务队列) |
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42[用户回车] ──> emit("agent_start")
│
▼
【外层第 1 圈】
│
├── ⚙️【内层第 1 圈 (Turn 1)】
│ 调模型 ──> 模型说“读 login.py” ──> 执行 read 工具 ──> turn_end
│ (判断:刚才调了工具,hasMoreToolCalls = True ──> 内层继续)
│
├── ⚙️【内层第 2 圈 (Turn 2)】
│ 调模型 ──> 模型说“修改 login.py” ──> 执行 edit 工具 ──> turn_end
│ (判断:刚才又调了工具,hasMoreToolCalls = True ──> 内层继续)
│
├── ⚙️【内层第 3 圈 (Turn 3)】
│ 调模型 ──> 模型说“Bug 已经修好了!”(stopReason="stop",无工具) ──> turn_end
│ (判断:没调工具 且 steering 队列为空 ──> 内层循环自然退出!)
│
▼
【内层结束,来到外层检查点】
│
├── 检查 followUp 队列 ──> 发现有一条:“修复完后自动追加运行 pytest 测试”
│
▼ (执行 continue,回到外层顶部,重新激活内层循环!)
│
【外层第 2 圈】
│
├── ⚙️【内层第 4 圈 (Turn 4)】
│ 注入 FollowUp 消息 ──> 调模型 ──> 模型说“执行 bash: pytest” ──> 跑测试 ──> turn_end
│ (判断:刚才调了工具 ──> 内层继续)
│
├── ⚙️【内层第 5 圈 (Turn 5)】
│ 调模型 ──> 模型输出“测试已全部通过,修改无误!” ──> turn_end
│ (判断:没调工具 ──> 内层循环退出!)
│
▼
【再次来到外层检查点】
│
├── 检查 followUp 队列 ──> 为空!
│
▼ (执行 break,退出外层循环!)
│
[整个任务圆满完成] ──> emit("agent_end")
steering 消息注入
端到端生命周期
1 | 【用户在 UI 打字 / 扩展调用】 |
message_queue.py
消息类型枚举:MessageType
1 | class MessageType(str, Enum): |
消息载体:QueuedMessage
1 |
|
消息队列核心:MessageQueue
MessageQueue 内部维护一个纯列表 self.queue: list[QueuedMessage],并实现了四大类能力:
### 1. 构造与消费模式配置 (init)
1
2
3
4
5
6
7
8def __init__(
self,
steering_mode: Literal["one-at-a-time", "all"] = "one-at-a-time",
followup_mode: Literal["one-at-a-time", "all"] = "one-at-a-time",
) -> None:
self.queue: list[QueuedMessage] = []
self.steering_mode = steering_mode
self.followup_mode = followup_mode
- one-at-a-time(默认推荐):单步推进模式。如果用户连发了两条转向指令,系统每次安全点只消费队首的第一条,等模型响应完后 再消费第二条,避免指令扎堆让大模型混淆;
- all(批量注入模式):一次性把队列里所有的同类消息全部提取出来,合并注入给模型。
### 2. 生产者入队方法 (add_steering / add_followup)
1
2
3
4
5def add_steering(self, message: str) -> None:
self.queue.append(QueuedMessage(content=message, type=MessageType.STEERING))
def add_followup(self, message: str) -> None:
self.queue.append(QueuedMessage(content=message, type=MessageType.FOLLOWUP))
- 外部调用 agent.steer(…) 或 agent.follow_up(…) 时,底层直接映射到这两个入队方法,将消息追加到列表尾部(保序 FIFO)
### 3. 消费者出队方法 (get_steering_messages / get_followup_messages)
这是队列中最关键的提取逻辑:
1
2
3
4
5
6
7
8
9
10
11
12def get_steering_messages(self) -> list[QueuedMessage]:
steering = [m for m in self.queue if m.type == MessageType.STEERING]
if not steering:
return []
if self.steering_mode == "one-at-a-time":
first = steering[0]
self.queue.remove(first)
return [first] # 只弹第一条,其余留在队列中
self.queue = [m for m in self.queue if m.type != MessageType.STEERING]
return steering # 弹出全部 steering 消息
- 分类过滤与移除:精准只提取目标类型(STEERING 或 FOLLOWUP),绝不误删另一种类型的排队消息;
- 提取后立即从 self.queue 物理移除,返回包含弹出消息的列表供循环注入
agent.py 与 loop.py 双层循环架构
双层循环控制流(Tau 对齐演进)
在经历阶段 18 深度重构后,这套双层循环控制流被正式抽离下沉为
packages/my-agent-core/src/my_agent_core/loop.py
中的无状态异步生成器 run_agent_loop。
Agent 作为轻量宿主外壳,通过将
message_queue.get_steering_messages 与
message_queue.get_follow_up_messages
作为无耦合回调注入微内核,在 三大安全点
驱动整个状态机的无缝换挡:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82# loop.py 中的纯函数微内核两层循环核心骨架:
async def run_agent_loop(
...,
get_steering_messages: Callable[[], Sequence[Message]] | None = None,
get_follow_up_messages: Callable[[], Sequence[Message]] | None = None,
) -> AsyncIterator[AgentEvent]:
...
# 初始化待处理消息列表
pending_messages: list[Message] = []
if get_steering_messages is not None:
init_steer = get_steering_messages()
if init_steer:
pending_messages.extend(init_steer)
iteration = 0
final_text: str | None = None
# ══════════════════════════════════════════════════════════════
# 【外层循环】:处理 Follow-up 宏观任务衔接
# ══════════════════════════════════════════════════════════════
while True:
has_more_tool_calls = True
# ──────────────────────────────────────────────────────────
# 【内层循环】:处理单任务的 ReAct 迭代与 Steer 即时转向
# ──────────────────────────────────────────────────────────
while has_more_tool_calls or len(pending_messages) > 0:
iteration += 1
if self._aborted:
await self._emit(AgentEnd(..., stop_reason="cancelled"))
return "(cancelled)"
await self._emit(TurnStart(iteration))
# ① 安全点 1:Turn 起始点注入 pending 消息并原子写盘
if pending_messages:
for text in pending_messages:
msg = Message(role="user", content=text)
self.messages.append(msg)
self.session.add_message("user", text) # 原子持久化
await self._emit(MessageStart(msg))
await self._emit(MessageEnd(msg))
pending_messages = []
# ── 准备上下文视图 + 调大模型(流式/非流式)
view = await self._ctx.prepare(self.messages)
...
# 记录 assistant 消息到 self.messages 与 session
...
# ── 检查是否有工具调用
if final_tool_calls:
# 执行工具并发/串行分流,写回 tool_results
tool_results = await self._execute_tool_batch(...)
has_more_tool_calls = True
else:
tool_results = []
has_more_tool_calls = False
final_text = content_acc
# 每个 Turn 结束时统一派发 TurnEnd(对齐 Pi 规范:成对闭环)
await self._emit(TurnEnd(message=assistant, tool_results=tool_results))
# ② & ③ 安全点:在 Turn 结束时统一检查 Steer 转向
if self.message_queue.has_steering():
steer_msgs = self.message_queue.get_steering_messages()
pending_messages = [m.content for m in steer_msgs]
# 关键:即使刚才 has_more_tool_calls=False(模型以为答完了),
# 因为 pending_messages > 0,内层循环绝不退出,立即带着转向指令开启新一轮!
# ──────────────────────────────────────────────────────────
# 内层循环自然结束 (无 tool_calls 且无 steering)
# ──────────────────────────────────────────────────────────
if self.message_queue.has_followup():
followup_msgs = self.message_queue.get_followup_messages()
pending_messages = [m.content for m in followup_msgs]
continue # 开启外层循环新一轮任务,重新激活内层循环!
break # 队列全清空,任务彻底完成
await self._emit(AgentEnd(..., stop_reason="end_turn"))
return final_text
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84 Agent.run(user_input)
│
[ 决策点 1: UserInput 拦截改写 ]
[ 决策点 2: AgentStart 提示词干预 ]
│
pending_messages = [ user_input ]
│
╔═══════════════════════════════════▼═══════════════════════════════════════════════════════╗
║ 【外层循环 (Outer Loop)】: 宏观任务生命周期与 Follow-up 队列管理 ║
║ while True: ║
║ ║
║ ┌───────────────────────────────────────────────────────────────────────────────────┐ ║
║ │ 【内层循环 (Inner Loop)】: 微观 ReAct 步骤迭代与 Steer 即时转向 │ ║
║ │ while has_more_tool_calls or len(pending_messages) > 0: │ ║
║ │ │ ║
║ │ ┌───────────────────────────────────────────────────────────────────────────┐ │ ║
║ │ │ ①【安全点 1: 消息注入与原子写盘】 │ │ ║
║ │ │ • 遍历 pending_messages ➔ session.add_message("user", ...) 原子落盘 │ │ ║
║ │ │ • 写入 self.messages ➔ 发射 MessageStart / MessageEnd │ │ ║
║ │ │ • pending_messages.clear() │ │ ║
║ │ └─────────────────────────────────────┬─────────────────────────────────────┘ │ ║
║ │ │ │ ║
║ │ ▼ │ ║
║ │ ┌───────────────────────────────────────────────────────────────────────────┐ │ ║
║ │ │ ②【Reason: 上下文准备与大模型推理】 │ │ ║
║ │ │ • view = await ctx.prepare(messages) ➔ 决策点 3: BeforeModelCall 拦截 │ │ ║
║ │ │ • assistant, tool_calls = await llm.achat_stream(view, tools) │ │ ║
║ │ │ • assistant 消息写入 session 与 self.messages ➔ 发射 MessageStart/End │ │ ║
║ │ └─────────────────────────────────────┬─────────────────────────────────────┘ │ ║
║ │ │ │ ║
║ │ 大模型本轮发起了工具调用吗? │ ║
║ │ / \ │ ║
║ │ [ 是 ] [ 否 ] │ ║
║ │ / \ │ ║
║ │ ▼ ▼ │ ║
║ │ ┌───────────────────────────────┐ ┌───────────────────────────────┐ │ ║
║ │ │ ③-A【Act+Observe: 工具批执行】│ │ ③-B【纯文本最终答复】 │ │ ║
║ │ │ • 决策点 4: 参数拦截与阻断 │ │ • 记录 final_text = content │ │ ║
║ │ │ • execute_batch 异步并发/串行│ │ • tool_results = [] │ │ ║
║ │ │ • 决策点 5: 篡改工具出参 │ │ • 标记: has_more_tools=False │ │ ║
║ │ │ • 保序写回 session 与消息 │ └───────────────┬───────────────┘ │ ║
║ │ │ • 标记: has_more_tools=True │ │ │ ║
║ │ └───────────────┬───────────────┘ │ │ ║
║ │ │ │ │ ║
║ │ └───────────────────┬───────────────────┘ │ ║
║ │ │ │ ║
║ │ ▼ │ ║
║ │ ┌───────────────────────────────────────────────────────────────────────────┐ │ ║
║ │ │ ④【Turn 闭环结算】 │ │ ║
║ │ │ • 派发配对的 TurnEnd(message=assistant, tool_results=tool_results) │ │ ║
║ │ └───────────────────────────────────┬───────────────────────────────────────┘ │ ║
║ │ │ │ ║
║ │ ▼ │ ║
║ │ 【安全点 2 & 3: 检查队列中有 Steer 转向消息吗?】 │ ║
║ │ / \ │ ║
║ │ [ 有 ] [ 无 ] │ ║
║ │ / \ │ ║
║ │ ┌────────────────────────┘ ▼ │ ║
║ │ │ 提取 Steer 消息存入 pending_messages 内层条件: 还有工具 或 pending非空? │ ║
║ │ │ (关键: 哪怕刚才走的是 ③-B 无工具分支, / \ │ ║
║ │ │ 只要 pending 非空, 内层循环绝不下班!) [ 是 ] [ 否 ] │ ║
║ │ │ / \ │ ║
║ │ └───────────────► 开启新一轮 Turn ◄───────────────┘ │ │ ║
║ └─────────────────────────────────────────────────────────────────────┼─────────────┘ ║
║ │ ║
║ (当前单项任务已彻底交付) ║
║ │ ║
║ ▼ ║
║ 【检查队列中有 Follow-up 追问吗?】 ║
║ / \ ║
║ [ 有 ] [ 无 ] ║
║ / \ ║
║ ┌───────────────────────────────────────────────────────┘ │ ║
║ │ 提取 Follow-up 追问消息存入 pending_messages │ ║
║ │ 执行 continue ➔ 回到外层顶部,重新唤醒内层循环! │ ║
║ └─────────────────────────────────────────────────────────────────────┤ ║
║ │ ║
║ ▼ ║
║ 执行 break 退出外层 ║
╚═══════════════════════════════════════════════════════════════════════════════╤═══════════╝
│
▼
[ 派发 AgentEnd 终结事件 ]
return final_text