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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
【用户按回车】──> emit("agent_start") ───【整个 Agent Loop 开始】

┌── Turn 1 ─────────▼──────────────────────────┐
│ emit("turn_start") │
│ 1. 调模型 -> 模型说“我要读 main.ts” │
│ 2. 执行 read 工具,拿到代码 │
│ emit("turn_end") ─── [Turn 1 结束] │
└───────────────────┬──────────────────────────┘
│ (判断:刚才调了工具,任务没完,继续转!)
┌── Turn 2 ─────────▼──────────────────────────┐
│ emit("turn_start") │
│ 1. 调模型 -> 模型说“我还要搜索 package.json” │
│ 2. 执行 grep 工具,拿到依赖 │
│ emit("turn_end") ─── [Turn 2 结束] │
└───────────────────┬──────────────────────────┘
│ (判断:刚才又调了工具,继续转!)
┌── Turn 3 ─────────▼──────────────────────────┐
│ emit("turn_start") │
│ 1. 调模型 -> 模型输出总结:“这个项目的入口在…” │
│ 2. 没有调用任何工具 │
│ emit("turn_end") ─── [Turn 3 结束] │
└───────────────────┬──────────────────────────┘
│ (判断:没调工具 + 队列无消息 -> 整个 Loop 可以停了)

emit("agent_end") ───────【整个 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
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
【用户在 UI 打字 / 扩展调用】

▼ session.steer("跳过测试文件,只改业务逻辑")
┌─────────────────────────────────────────────────────────────┐
│ 1. 同步入队:加入 steeringQueue(PendingMessageQueue) │
│ (完全不打断当前正在进行的 LLM 流式输出或 Tool 运行) │
└─────────────────────────────────────────────────────────────┘

▼ 当前 Turn 正在运行(调用模型 -> 执行完 Read 工具)
┌─────────────────────────────────────────────────────────────┐
│ 2. emit("turn_end") —— 当前 Turn 结束 │
└─────────────────────────────────────────────────────────────┘


┌─────────────────────────────────────────────────────────────┐
│ 3. 检查队列:pendingMessages = await getSteeringMessages() │
│ (从队列中取出刚才插队的消息) │
└─────────────────────────────────────────────────────────────┘


┌─────────────────────────────────────────────────────────────┐
│ 4. 循环判断:hasMoreToolCalls (false) || pendingMessages (>0)│
│ ★ 关键点:原本模型没要工具该停了,但被 Steering 强行“续命” │
└─────────────────────────────────────────────────────────────┘

▼ 进入新一轮 Turn
┌─────────────────────────────────────────────────────────────┐
│ 5. 注入消息: │
│ - emit("turn_start") │
│ - emit("message_start" / "message_end") │
│ - context.messages.append(steering_message) │
│ - new_messages.append(steering_message) │
└─────────────────────────────────────────────────────────────┘


┌─────────────────────────────────────────────────────────────┐
│ 6. 调用 LLM:模型在最新上下文中看到了 Steering 消息,调整决策 │
└─────────────────────────────────────────────────────────────┘

message_queue.py

消息类型枚举:MessageType

1
2
3
4
5
class MessageType(str, Enum):
"""排队干预消息类型。"""

STEERING = "steering" # 内层循环即时转向(安全点打断)
FOLLOWUP = "followup" # 外层循环排队追问(任务完成后驱动)

消息载体:QueuedMessage

1
2
3
4
5
6
7
@dataclass
class QueuedMessage:
"""队列中的一条干预消息。"""

content: str
type: MessageType
created_at: float = field(default_factory=time.time)

消息队列核心:MessageQueue

MessageQueue 内部维护一个纯列表 self.queue: list[QueuedMessage],并实现了四大类能力:

### 1. 构造与消费模式配置 (init)

1
2
3
4
5
6
7
8
def __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
5
def 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
12
def 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_messagesmessage_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