Pi 源码解析(三):Agent 的状态管理

Pi Agent Core 源码系列(3/8)|上一篇:Agent Loop:生命周期、事件与控制流|下一篇:Tool Runtime:流式响应与工具执行

agent.ts 在 Agent Loop 外包了一层内存状态。它保存模型、工具和消息,根据 Loop 事件更新 AgentState,并提供 prompt()continue()abort()、消息队列和订阅 API。

Agent 保存哪些状态

Agent Loop 每次调用都要求调用方提供上下文、配置、事件 Sink 和取消信号。直接使用它适合自定义 Runtime,但普通应用还需要维护流式消息、工具执行状态和运行互斥。

Agent 把这些工作收进一个对象:

flowchart LR API[prompt / continue / abort] --> Agent[Agent] Agent --> Snapshot[Context + Config 快照] Snapshot --> Loop[Agent Loop] Loop --> Events[AgentEvent] Events --> Reducer[processEvents] Reducer --> State[AgentState] Reducer --> Listeners[订阅者]

最小使用方式如下:

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
import { Agent } from "@earendil-works/pi-agent-core";
import { createModels } from "@earendil-works/pi-ai";
import { anthropicProvider } from "@earendil-works/pi-ai/providers/anthropic";

const models = createModels();
models.setProvider(anthropicProvider());

const model = models.getModel("anthropic", "claude-sonnet-4-6");
if (!model) throw new Error("Model not found");

const agent = new Agent({
initialState: {
model,
systemPrompt: "Answer as a software engineer.",
},
streamFunction: models.streamSimple.bind(models),
});

agent.subscribe((event) => {
if (
event.type === "message_update" &&
event.assistantMessageEvent.type === "text_delta"
) {
process.stdout.write(event.assistantMessageEvent.delta);
}
});

await agent.prompt("Explain the current module.");

构造函数只建立状态和配置,不会调用模型。真正的运行从 prompt()continue() 开始。

AgentState 分成稳定状态和瞬态状态

AgentState 同时保存跨 Run 配置和当前 Run 状态:

状态 生命周期 用途
systemPrompt 跨 Run 下一次模型请求的系统指令
modelthinkingLevel 跨 Run 当前模型配置
toolsmessages 跨 Run 工具集合与内存历史
isStreaming 当前 Run 是否正在处理请求
streamingMessage 当前消息 最新 assistant partial 或当前消息
pendingToolCalls 当前工具批次 正在执行的工具调用 ID
errorMessage 最近失败 Turn 提供给 UI 的错误文本

内部使用 MutableAgentState 更新只读运行字段。toolsmessages 的 setter 会复制顶层数组,但不会深拷贝数组内对象,因此工具定义和消息对象仍应被视为不可变值。

每次运行从状态创建浅快照

prompt() 最终调用 runAgentLoop()continue() 调用 runAgentLoopContinue()。两条路径都会先执行:

1
2
3
4
5
6
7
private createContextSnapshot(): AgentContext {
return {
systemPrompt: this._state.systemPrompt,
messages: this._state.messages.slice(),
tools: this._state.tools.slice(),
};
}

快照防止 Loop 直接持有状态数组,但它不是任意对象的深拷贝。一个 Run 启动后,初始模型、上下文和工具来自这一时刻;跨 Turn 更新则通过 prepareNextTurn 返回显式的新快照。

Agent 的快照职责比 AgentHarness 简单。它没有 Session、save point 或资源投影,只把当前内存状态适配成 AgentLoopConfig

processEvents 如何更新状态

Loop 产生的每个事件先进入 processEvents()。方法内部先更新状态,再按注册顺序通知订阅者:

1
2
3
4
5
AgentEvent
→ 更新 AgentState
→ Listener A
→ Listener B
→ 返回 Agent Loop

message_startmessage_update 更新 streamingMessagemessage_end 把完成消息追加到 messages;工具开始与结束事件分别修改 pendingToolCallsturn_end 提取 assistant 错误。

对应源码就是一个事件归约器。下面保留消息和工具状态的主要分支:

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
private async processEvents(event: AgentEvent): Promise<void> {
// 这一 switch 相当于 AgentEvent → AgentState 的 reducer。
switch (event.type) {
case "message_start":
// user、assistant 和 toolResult 开始时都会暂时成为 streamingMessage。
this._state.streamingMessage = event.message;
break;

case "message_update":
// 仅 assistant 有增量事件,用最新 partial 快照覆盖当前流式消息。
this._state.streamingMessage = event.message;
break;

case "message_end":
// 消息完整结束时移入正式历史,并清除流式占位。
this._state.streamingMessage = undefined;
this._state.messages.push(event.message);
break;

case "tool_execution_start": {
// 创建新 Set 再替换,便于依赖引用变化的 UI/状态系统观察更新。
const pendingToolCalls = new Set(this._state.pendingToolCalls);
pendingToolCalls.add(event.toolCallId);
this._state.pendingToolCalls = pendingToolCalls;
break;
}

case "tool_execution_end": {
const pendingToolCalls = new Set(this._state.pendingToolCalls);
pendingToolCalls.delete(event.toolCallId);
this._state.pendingToolCalls = pendingToolCalls;
break;
}
}

// 事件必须属于一个活动运行,否则说明生命周期调用顺序出现了内部错误。
const signal = this.activeRun?.abortController.signal;
if (!signal) {
throw new Error("Agent listener invoked outside active run");
}

// Set 保持插入顺序;逐个 await 使监听器之间形成确定的顺序屏障。
for (const listener of this.listeners) {
await listener(event, signal);
}
}

完整实现位于 agent.ts。这里为便于阅读省略了 turn_endagent_end 和运行状态校验,但状态先更新、监听器后执行的顺序没有改变。

订阅者读取状态时已经能看到该事件对应的新值。监听器逐个 await,所以慢监听器会形成背压,抛错也会进入当前 Run 的错误处理路径。这种顺序适合需要一致 UI 的应用,但耗时且与运行无关的工作不应无界地阻塞监听器。

activeRun 防止重复运行

一个 Agent 同时只允许一个 prompt()continue()。运行期间再次调用会直接报错;新增意图应使用 steer()followUp()

stateDiagram-v2 [*] --> Idle Idle --> Running: prompt / continue Running --> Aborting: abort Aborting --> Running: 等待 Provider/工具响应 signal Running --> Idle: Loop + listeners + cleanup 完成

activeRun 保存 AbortController 和一个手动完成的 Promise。abort() 只发送取消信号,Provider、工具和 Hook 需要自行响应;waitForIdle() 等待 Loop、最后的 agent_end 监听器以及 finishRun() 清理全部完成。

这个区分解释了为什么看到 agent_end 不等于对象已经 idle。事件仍在被最后一批订阅者处理,isStreaming 也要到清理阶段才会变回 false

prompt 和 continue 表达不同起点

prompt() 接收字符串、单条 AgentMessage 或消息数组。字符串会被规范化为带时间戳的 user 消息,并可附带图片,然后作为本次 Run 的新消息进入上下文。

continue() 不添加新 prompt,而是从已有历史继续。最后一条消息必须是 user 或 toolResult;如果历史以 assistant 结束,Agent 会先尝试消费已排队的 steering,再尝试 follow-up,否则拒绝产生没有用户或工具输入的 assistant-to-assistant 请求。

这使失败重试和恢复场景不必伪造一条新的用户消息,同时保留 Provider 对消息角色顺序的要求。

两个队列何时消费

Agent 为 steering 和 follow-up 各维护一个 FIFO 队列。one-at-a-time 每个安全点只消费最早一条,all 一次消费当前全部消息。

1
2
3
4
5
steer(message)
→ 当前 Turn 结束后进入下一 Turn

followUp(message)
→ Agent 原本准备停止时进入下一 Turn

队列只负责保存和取出消息,实际消费时机仍由 Agent Loop 的双层循环决定。createLoopConfig() 把两个 drain() 方法作为回调交给 Loop,因此 Agent 不需要复制控制流。

封装层异常也会补齐生命周期

如果异常没有被模型流或工具协议归一化,runWithLifecycle() 会构造一条 stopReason: "error""aborted" 的 assistant 消息,并依次发送:

1
2
3
4
message_start
→ message_end
→ turn_end
→ agent_end

UI 和监听器不需要为“Loop 正常失败”和“封装层抛错”维护两套状态机。finally 中的 finishRun() 仍会清理流式消息、工具集合和 idle barrier。

Agent 和 AgentHarness 怎么选

能力 Agent AgentHarness
消息保存 内存数组 Session Tree 与 Storage
状态更新 Event → AgentState Event → Session、快照与应用事件
消息队列 steering、follow-up steering、follow-up、next-turn
资源 直接传入工具和 prompt Skills、Templates、System Prompt 资源
长会话 调用方实现 transformContext 内置 Compaction 与 Branch Summary 编排
适用场景 嵌入式、临时或自定义状态层 需要恢复、分支和扩展的完整应用

两者都直接复用 Agent Loop。选择哪一个取决于是否需要 Session 级语义,而不是是否需要工具调用。

源码阅读路径

createMutableAgentState() 看状态初始化,再阅读 prompt()continue()runWithLifecycle() 理解运行入口。最后跟踪 createLoopConfig()processEvents(),就能看到 Agent 如何在 Loop 两端完成“状态转配置”和“事件转状态”。

总结

Agent 运行前创建快照,运行中根据事件更新状态,结束后等待监听器和清理工作完成。它不管 Session 和资源加载,适合临时任务或自定义状态层的应用。