03 — Agent 运行时核心 (pi-agent-core)

这是最核心的文档。理解了这一篇,你就理解了 Agent 的本质。

问题

一个 Agent 需要持续地与 LLM 交互:发送上下文 → 接收响应 → 执行工具 → 将工具结果加入上下文 → 再次发送。这个循环何时继续、何时停止?工具调用如何执行?如果 LLM 返回错误怎么办?如果用户想在 Agent 工作时插入新指令怎么办?

pi-agent-core 提供了两套实现来回答这些问题:

  1. 经典 Agentagent.ts + agent-loop.ts)— 简单的 while 循环,适合理解概念
  2. AgentHarnessharness/ 目录)— 生产级状态机,支持崩溃恢复、多 lane、compaction

第一层:经典 Agent

Agent 类

packages/agent/src/agent.ts

Agent 类持有 AgentState

1
2
3
4
5
6
7
8
9
10
11
interface AgentState {
systemPrompt: string;
model: Model<any>;
thinkingLevel: ThinkingLevel; // "off" | "low" | "medium" | "high"
tools: AgentTool<any>[];
messages: AgentMessage[];
isStreaming: boolean;
streamingMessage?: AgentMessage;
pendingToolCalls: Set<string>;
errorMessage?: string;
}

Agent 类本身不执行循环——它管理状态,通过 prompt() 方法启动 agentLoop()

AgentMessage 与 convertToLlm 边界

AgentMessage 是 Agent 内部的消息类型,比 LLM 的 Message 类型更丰富——可以包含 UI 通知、状态消息等 LLM 不需要看到的内容。

关键设计:在调用 LLM 之前,通过 convertToLlm()AgentMessage[] 转换为 LLM 能理解的 Message[]

1
2
AgentMessage[]  ──convertToLlm──▶  Message[]  ──▶  LLM
(内部格式) (过滤+转换) (标准格式) (provider API)

convertToLlm 的默认实现(agent.ts:33)只保留 userassistanttoolResult 三种角色的消息,过滤掉 UI-only 消息。

Agent Loop 主循环

packages/agent/src/agent-loop.ts:156-273

这是理解 Agent 的核心代码。循环结构是 双层 while 循环

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
外层循环 (while true)
│ 职责:处理 follow-up 消息(Agent 本该停止但有排队消息时继续)

└─ 内层循环 (while hasMoreToolCalls || pendingMessages.length > 0)
│ 职责:处理工具调用 + steering 消息

├─ 1. prepareNextTurn() — 准备下一轮(可能执行 compaction)
├─ 2. 注入 pendingMessages(steering 消息)
├─ 3. streamAssistantResponse() — 调用 LLM 获取响应
├─ 4. 检查 stopReason: error/aborted → 终止
├─ 5. 提取 toolCalls
├─ 6. executeToolCalls() — 执行工具
├─ 7. turn_end 事件
├─ 8. shouldStopAfterTurn() 检查
└─ 9. 获取新的 steering 消息

流式响应处理

streamAssistantResponse()agent-loop.ts:279)是实际调用 LLM 的地方:

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
// 1. 可选的上下文变换(AgentMessage 级别)
let messages = context.messages;
if (config.transformContext) {
messages = await config.transformContext(messages, signal);
}

// 2. 转换为 LLM 格式
const llmMessages = await config.convertToLlm(messages);

// 3. 构建 LLM 上下文
const llmContext = { systemPrompt, messages: llmMessages, tools };

// 4. 调用 streamFn(不抛异常,错误在流中)
const response = await streamFunction(config.model, llmContext, { ... });

// 5. 消费流式事件
for await (const event of response) {
switch (event.type) {
case "start": // assistant 消息开始
case "text_delta": // 文本增量
case "thinking_delta": // 思考增量
case "toolcall_delta": // 工具调用参数增量
case "done": // 完成
case "error": // 错误(不抛异常,在流中表达)
}
}

关键streamFunction 满足 StreamFn 契约——不会抛异常。如果 LLM 返回错误,stopReason 会是 "error",在 agent-loop.ts:215 被检测到后优雅地终止循环。

工具系统

packages/agent/src/types.ts:52-95

AgentTool 接口

每个工具实现 AgentTool 接口:

1
2
3
4
5
6
7
8
interface AgentTool<TResult = unknown> {
name: string;
description: string;
inputSchema: TSchema; // TypeBox schema,用于参数验证
execute(args, context): Promise<AgentToolResult<TResult>>;
prepareArguments?(toolCall): AgentToolCall; // 可选的参数预处理
executionMode?: "sequential" | "parallel"; // 执行模式
}

beforeToolCall / afterToolCall 钩子

在工具执行前后,Agent loop 调用可选的钩子:

beforeToolCalltypes.ts:278):

1
beforeToolCall?: (context: BeforeToolCallContext, signal?) => Promise<BeforeToolCallResult | undefined>;

返回 { block: true } → 阻止工具执行,返回错误结果给 LLM。
返回 { terminate: true } → 暗示这批工具执行后应该停止。

afterToolCalltypes.ts:84):

1
afterToolCall?: (context: AfterToolCallContext) => Promise<AfterToolCallResult | undefined>;

可以覆盖工具结果:contentisErrorusageterminate

顺序 vs 并行执行

agent-loop.ts:409-424

1
2
3
4
5
6
7
8
9
async function executeToolCalls(...) {
const hasSequentialToolCall = toolCalls.some(
(tc) => tools?.find((t) => t.name === tc.name)?.executionMode === "sequential",
);
if (config.toolExecution === "sequential" || hasSequentialToolCall) {
return executeToolCallsSequential(...);
}
return executeToolCallsParallel(...);
}
  • 顺序模式:逐个执行工具,每个工具完成后才执行下一个。适合有副作用的工具(如文件写入)。
  • 并行模式:先顺序做 prepare(验证参数、beforeToolCall),然后并行执行所有工具。tool_execution_end 按完成顺序发出,但 toolResultMessage 按 assistant 消息中的源顺序发出。

截断消息的处理

agent-loop.ts:379:如果 stopReason === "length"(输出被 token 限制截断),所有工具调用都不执行,直接返回错误结果。因为截断的 JSON 参数可能恰好能解析但内容不完整,执行会有风险。

Steering 与 Follow-up

这是 Pi 的一个独特设计——用户可以在 Agent 工作时插入消息

  • Steering 消息getSteeringMessagestypes.ts:245):在 Agent 执行完工具调用后、下一轮 LLM 调用前注入。用户可以在 Agent 工作时”转向”。
  • Follow-up 消息getFollowUpMessagestypes.ts:258):在 Agent 本该停止时检查,如果有排队的后续消息则继续运行。
1
2
3
4
5
6
7
8
9
10
11
12
13
Agent 正在工作    用户输入 "也看看 src/ 目录"
│ │
│ ┌─────▼─────┐
│ │ steering │ 被注入到下一轮
│ │ queue │
│ └─────┬─────┘
│ │
▼ │
工具执行完成 ◄──────────┘

▼ getSteeringMessages() → ["也看看 src/ 目录"]

▼ 下一轮 LLM 调用时,steering 消息已加入上下文

事件序列

完整的 Agent 事件流:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
agent_start
turn_start
message_start (user)
message_end
message_start (assistant) ← streamAssistantResponse
message_update (text_delta...) ← 流式增量
message_update (toolcall_delta...)
message_end
tool_execution_start ← 工具开始执行
tool_execution_end ← 工具执行完成
tool_result ← 工具结果消息
turn_end
turn_start ← 如果有更多工具调用
message_start (assistant)
...
message_end
turn_end
agent_end

第二层:AgentHarness(生产级架构)

经典 Agent 简单但不支持:崩溃恢复、上下文压缩、多会话分支、持久化。AgentHarness 解决这些问题。

核心概念

概念 文件 说明
Harness harness/runtime/harness.ts 管理多个 Lane、配置、钩子、事件
Lane harness/runtime/lane.ts 一条执行通道(一个会话分支),持有状态和操作
Drive harness/runtime/drive.ts 驱动一次操作通过状态机直到完成或等待
Operation lane 内部状态 一次 run/compact/navigation 操作的完整生命周期
Session harness/session/ 会话持久化(Entry 树、分支)

Lane 状态机

packages/agent/src/harness/runtime/drive.ts:29-80driveOperation() 是状态机的核心调度:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
async function driveOperation(lane, drive): Promise<DriveOutcome> {
for (;;) {
operation = currentOperation(lane, drive);
const state = operation.state;
let result: ProcedureResult;

if (state.control.status === "cancel_requested") {
result = await reconcileOperation(lane, drive); // 取消协调
} else switch (state.at) {
case "starting": → startRun() // 读取 prompt,提交分支
case "checkpoint": → runCheckpoint() // 检查 compaction、注入消息
case "assistant.ready": → runGeneration() // 调用 LLM
case "assistant.retry_wait": → runGeneration() // 重试等待
case "assistant.effect_pending": → recoverAssistantGeneration() // 恢复中断的生成
case "tools": → runTools() // 执行工具
case "deferred.suspended": → runDeferred() // 异步响应轮询
case "deferred.effect_pending": → runDeferred() // 恢复中断的轮询
case "summary.deciding": → runStructuralDecision() // compaction 决策
case "summary.ready": → runStructuralGeneration() // 执行 compaction
}
// result 决定:继续循环、等待、完成
}
}

状态转换图:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
starting


checkpoint ────────▶ summary.deciding ──▶ summary.ready ──▶ checkpoint
│ (需要 assistant) (compaction 完成后)

assistant.ready


assistant.effect_pending ──(崩溃恢复)──▶ checkpoint


tools ──▶ checkpoint


deferred.suspended ──▶ deferred.effect_pending ──▶ checkpoint


完成 (terminal)

Drive 各阶段详解

1. starting → checkpoint

drive/checkpoint.ts:25startRun()

  • 读取 prompt entries
  • 运行 before_run 钩子(可注入消息)
  • 提交 prompt entries + 分支 tip
  • 转换到 checkpoint 状态

2. checkpoint — 边界决策点

drive/checkpoint.ts:95runCheckpoint()

  1. 检查 compaction 阈值(上下文是否接近满)
  2. 规划边界 inbox(注入 steering/follow-up/write 消息)
  3. 如果有触发 entry → 转到 assistant.ready
  4. 如果阈值超标 → 转到 summary.deciding(compaction)
  5. 如果可以结束 → finish_pending

drive/boundary.ts:76planBoundaryInbox():选择 steer/followUp/write 项,加载待处理 payload,链式连接 entries。

3. assistant.ready — LLM 生成

drive/generation.ts

  • prepareGeneration()(line 68):解析模型、验证工具、读取有界上下文、解析系统提示、运行 before_request 钩子
  • publishGenerationIntent()(line 132):提交 assistant.effect_pending 状态(预留 response/usage ID)
  • performGeneration()(line 174):打开 assistant 响应,调用 streamHarnessAssistant()
  • runGeneration()(line 285):编排:prepare → publish intent → perform → publish response

关键publishGenerationIntent() 在实际调用 LLM 之前就将 effect_pending 状态持久化。如果进程在此之后崩溃,恢复时知道”有一次生成未完成”,可以尝试从已提交的帧中恢复部分响应。

4. assistant.effect_pending — 崩溃恢复

drive/recovery.ts:44recoverAssistantGeneration()

  • 读取已提交的流式帧(partial response)
  • 将帧归约为 partial message
  • 标记为 interrupted error
  • 发布恢复事件
  • 通过 publishResponse() 结算

这就是 Harness 的核心价值:即使 LLM 调用中途崩溃,也能从已持久化的帧中恢复部分响应

5. tools — 工具执行

drive/tools.ts

  • runTools()(line 656):检查恢复(有 effect_pending 的调用需要恢复),读取工具批次,分发到 runSequentialrunParallel
  • startToolInvocation()(line 475):准备工具 → publishToolIntent()(提交 effect_pending)→ performToolInvocation()publishToolOutcome()
  • recoverToolInvocation()(line 515):对于 effect_pending 的工具调用,如果工具支持 replay === "safe" 则重新执行,否则从检查点发布中断结果

与经典 Agent 的区别:每个工具调用在执行前先持久化 intent(effect_pending),执行后持久化 outcome(outcome_ready)。崩溃后可以恢复——知道哪些工具已执行、哪些未执行。

drive/tool-placement.ts:工具结果按 assistant 源顺序(而非完成顺序)提交,保证消息顺序与 LLM 输出一致。

6. response 分类

drive/response.ts:181publishResponse() 是响应分类器,根据 stopReason 决定下一步:

stopReason 处理
cancel_requested 标记为 aborted,转到 checkpoint
上下文溢出 首次 → 准备 overflow compaction;已试过 → 失败
deferred 验证 handle,转到deferred.suspended
error 可重试 →assistant.retry_wait;否则失败
有 tool calls 规划工具调用,转到tools
无 tool calls 但toolUse 失败(格式错误)
正常结束 转到checkpoint,可以结束

7. summary — Compaction

drive/structural.ts

  • runStructuralDecision():运行 before_compaction 钩子,决定是否执行 compaction
  • runStructuralGeneration():准备摘要请求,调用 LLM 生成摘要
  • publishStructuralOutcome():处理三种边界类型:
    • resume_checkpoint — compaction 后继续 run
    • finish — 独立 compaction 完成
    • commit_navigation — 分支摘要完成

Compaction 机制

packages/agent/src/harness/compaction/

当对话接近上下文窗口限制时:

  1. 触发checkpoint 阶段检查 compaction 阈值
  2. 决策before_compaction 钩子决定是否压缩
  3. 生成:向 LLM 发送摘要请求,总结旧消息
  4. 提交:将摘要作为 CompactionEntry 提交,替代旧消息
  5. 恢复:后续轮次只读取从最近 compaction 边界开始的消息

compaction/branch-summarization.ts:处理分支摘要——当在分支间导航时,为未摘要的分支生成摘要。

Session 持久化模型

packages/agent/src/harness/session/

Entry 类型

类型 说明
MessageEntry 用户/assistant/toolResult 消息
CompactionEntry compaction 摘要(替代旧消息)
BranchSummaryEntry 分支导航摘要
CustomEntry 自定义条目(扩展使用)

分支 (Branch)

Entry 通过 parentId 形成树结构。分支是从根到某个叶子的路径。当用户 fork 一个会话或导航到历史中的某一点再继续时,会创建新分支。

Fork

session/ 中的 fork 逻辑:从一个会话的快照创建新会话,包含所有 entries 和 scalar/list values。可以是完整树 fork 或仅当前分支 fork。

Hook 与事件总线

HookRegistry

packages/agent/src/harness/hooks.ts

生命周期钩子,在每个关键节点被调用:

钩子 时机 可做
before_run 操作开始前 注入消息
before_request LLM 调用前 修改请求
before_tool 工具执行前 阻止执行
after_tool 工具执行后 修改结果
after_response LLM 响应后 修改响应
before_compaction compaction 前 决定是否执行
before_run_end 操作结束前 注入 follow-up
before_drive 每次 drive 开始 取消控制

HarnessEventBus

packages/agent/src/harness/events.ts:事件分发总线,将 harness 内部事件(message_starttool_startentry_addedrun_end 等)分发给订阅者。

Reducer — 状态归约

harness/runtime/reducer.tsreduceLaneSnapshot() 将事件序列归约为 lane 的当前状态。每个事件类型更新状态的不同部分:

  • message_start/end → 管理流式消息
  • tool_start/end → 更新 runningTools
  • entry_added → 推入 transcript,更新 tipId
  • run_end/compaction_end → 设置 lastResult,清除 operation
  • fault → 设置 faulted: true

Restore — 恢复

harness/runtime/restore.ts

  • restoreSession()(line 92):扫描所有分支 tip、lane 配置、lane 状态,恢复每个完整的 lane
  • restoreLaneState()(line 131):恢复一个 lane,包括其活动操作(读取操作 meta + state,验证 intent 与 state 匹配)

这就是崩溃恢复的核心:进程重启后,从持久化存储中恢复所有 lane 状态,包括正在进行中的操作。effect_pending 状态的操作会进入恢复路径。

Effect Gate — 副作用门控

harness/execution/effect-gate.ts

1
2
// Gate(过程侧):signal, admit()
// GateControl(所有者侧):beginAbort(), signalAbort(), close()

状态机:openaborting(抛出 AbortRequested)→ closed

在工具执行或 LLM 调用时,通过 gate.admit() 获取执行许可。如果取消请求到达,admit() 抛出 AbortRequested,携带一个 cancellation promise。这让取消操作可以安全地等待正在进行的副作用完成。

两层架构对比

维度 经典 Agent AgentHarness
状态管理 内存中的AgentState 持久化的LaneState+ 操作状态机
崩溃恢复 不支持 支持(从持久化状态恢复)
Compaction 需要外部实现 内置(summary 状态)
多会话 不支持 多 Lane 管理
工具执行恢复 不支持 支持(effect_pendingoutcome_ready
复杂度 ~800 行 ~5000+ 行
适用场景 简单用例、学习理解 生产环境

学习要点

  • Agent loop 的本质是双层 while 循环:内层处理工具调用,外层处理 follow-up。理解了这个结构,就理解了 Agent 的运行模型。
  • StreamFn 不抛异常是状态机一致性的基础:错误在流中表达,让 loop 以统一路径处理正常和异常情况。
  • Harness 的核心价值是可恢复性:通过 effect_pendingoutcome_ready 的两阶段提交,保证每个副作用都有明确的”已提交”或”未执行”状态。崩溃后可以精确恢复。
  • 工具结果按源顺序而非完成顺序提交:保证消息顺序与 LLM 输出一致,这对 LLM 理解上下文很重要。
  • Compaction 不是删除消息,而是创建摘要 entry:原始 entries 仍在存储中,只是后续读取从 compaction 边界开始。这让 compaction 是可逆的。
  • convertToLlm 是内部格式与 LLM 格式的边界:Agent 内部可以持有比 LLM 需要更丰富的消息类型,在调用 LLM 前过滤和转换。

上一篇:02-LLM 抽象层-pi-ai
下一篇:04-会话持久化-pi-protocol 与 session-backends