Pi Agent 源码解析(一):从调试环境到双层主循环
最近一段时间我一直在使用 Pi。它的交互界面看起来很简单,但代码里把模型适配、会话管理、上下文处理、工具执行和终端 UI 分成了比较清楚的几层,很适合拿来理解一个 Coding Agent 到底是怎么运转的。
前面已经写过一篇 Pi Agent 工具系统:Read、Write、Edit 与 Bash 的实现原理,那篇关注 Agent 的“手脚”,本文继续向下看驱动这些工具反复工作的主循环。
这篇先从源码调试环境开始,然后沿着一次真实请求的调用链,重点看下面几个问题:
- Pi 的源码怎样安装、补齐模型数据并启动调试;
- 模型的
baseUrl和apiKey从哪里读取; Agent在哪里初始化,streamFn又是在哪里注入的;agent-loop.ts中的两层循环分别解决什么问题;- LLM 的流式事件如何变成 Agent 事件;
- 模型发起工具调用以后,工具结果如何回到上下文并触发下一轮请求。
本文对应我本地阅读的源码版本:1
2
3
4repository: https://github.com/earendil-works/pi
commit: 71dca871bc80b6bc97be37f0ca3189399d651fff
describe: v0.85.1-67-g71dca871b
date: 2026-09-11
Pi 的代码还在快速变化,后续如果发现文件名或行号对不上,应当先确认源码版本。本文代码片段均摘自上述提交;为了方便阅读,我只在少数位置加入了以 // 解读: 开头的注释,没有把真实逻辑改写成伪代码。
一、先看项目分层
仓库根目录的开发文档给出的主要结构如下:1
2
3
4
5packages/
ai/ # LLM provider 抽象、模型目录以及各家 API 适配
agent/ # Agent 状态、事件、主循环和工具执行
tui/ # 终端 UI
coding-agent/ # CLI、会话、配置、扩展、内置工具和交互模式
阅读时最容易混淆的是 agent 与 coding-agent:
packages/agent是较底层、与具体模型供应商无关的 Agent Core;packages/coding-agent是完整产品层,负责加载用户配置、模型运行时、会话、工具、扩展和 TUI;packages/ai才真正知道 OpenAI、Anthropic、Google 等供应商的请求协议。
所以 agent-loop.ts 不会直接调用 OpenAI SDK。它只认识一个抽象的 StreamFn,具体实现由上层注入。这是理解整套代码最关键的一点。
二、搭建源码调试环境
1. Node.js 版本
根目录 package.json 的要求不是笼统的 Node 22,而是:1
2
3"engines": {
"node": ">=22.19.0"
}
因此建议先确认版本:1
2node --version
npm --version
我使用的是 Node 22。如果是较早的 22.x,也要确认小版本不低于 22.19.0。
2. Clone 与安装依赖
1 | git clone https://github.com/earendil-works/pi.git |
这里补充一下 npm install 和 npm ci 的区别。这两个通常是二选一,并不需要固定地先后执行:
npm ci严格按照已有的package-lock.json做干净安装,适合复现仓库环境;npm install更适合要调整依赖、更新 lockfile 的开发场景。
如果只是阅读和调试源码,我更推荐 npm ci。仓库自己的开发规则还要求使用 --ignore-scripts,因此更谨慎的做法是:1
npm ci --ignore-scripts
如果确实需要仓库生命周期脚本,再在理解脚本用途后单独执行。
3. 下载模型目录
安装依赖以后还要补齐模型数据:1
npm run hydrate:model-data
根目录只是把命令转发给 packages/ai:1
2"hydrate:model-data": "npm --prefix packages/ai run hydrate-model-data",
"check:model-data": "npm --prefix packages/ai run check:model-data"
而 packages/ai/package.json 里的实际命令是:1
2"hydrate-model-data": "node scripts/generate-models.ts --strict --data-only",
"check:model-data": "node scripts/check-model-data.ts"
hydrate:model-data 不是下载某个模型权重,它下载的是 Pi 内置模型目录所需的元数据,例如模型 ID、API 类型、上下文窗口、最大输出长度、价格和是否支持 reasoning 等。
generate-models.ts 会访问多个在线目录,代码里可以直接看到这些请求: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
30async function fetchNvidiaNimModelIds(): Promise<Map<string, string>> {
try {
console.log("Fetching models from NVIDIA NIM API...");
const response = await fetch(`${NVIDIA_BASE_URL}/models`);
if (!response.ok) throw new Error(`NVIDIA NIM API returned ${response.status}`);
// ...
} catch (error) {
console.error("Failed to fetch NVIDIA NIM models:", error);
if (generatorOptions.strict) throw error;
return new Map();
}
}
async function fetchOpenRouterModels(): Promise<Model<any>[]> {
try {
console.log("Fetching models from OpenRouter API...");
const response = await fetch("https://openrouter.ai/api/v1/models");
if (!response.ok) throw new Error(`OpenRouter API returned ${response.status}`);
// ...
}
}
async function loadModelsDevData(): Promise<Model<any>[]> {
try {
console.log("Fetching models from models.dev API...");
const response = await fetch("https://models.dev/api.json");
if (!response.ok) throw new Error(`models.dev API returned ${response.status}`);
// ...
}
}
命令带有 --strict,某个关键数据源访问失败时会直接报错,而不是悄悄生成残缺目录。因此在部分网络环境下可能需要配置代理,也就是大家常说的“科学上网”。
数据不会直接覆盖旧目录。生成器先在临时目录写入并校验,成功后才替换 src/providers/data:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21const providersDir = join(packageRoot, "src/providers");
const dataDir = join(providersDir, "data");
const stagingRoot = mkdtempSync(join(providersDir, ".model-generation-"));
const stagedDataDir = join(stagingRoot, "data");
const previousDataDir = join(stagingRoot, "previous-data");
// 先向 stagedDataDir 写入各 provider 的 JSON 和 manifest
validateModelDataDirectory(modelDataStructure, stagedDataDir);
const hadPreviousData = existsSync(dataDir);
if (hadPreviousData) renameSync(dataDir, previousDataDir);
try {
renameSync(stagedDataDir, dataDir);
validateGeneratedModelData(packageRoot);
} catch (error) {
rmSync(dataDir, { recursive: true, force: true });
if (hadPreviousData && existsSync(previousDataDir)) {
renameSync(previousDataDir, dataDir);
}
throw error;
}
这实际上是一个带回滚的目录级原子替换:新数据先通过校验,替换后再做一次整体校验;失败时恢复旧目录,避免半生成状态污染工作区。
4. 验证模型数据
1 | npm run check:model-data |
对应的检查入口很短:1
2
3
4
5
6
7
8
9
10
11try {
validateGeneratedModelData(packageRoot);
console.log("Generated model data is valid.");
} catch (error) {
console.error(error instanceof Error ? error.message : String(error));
console.error(
"\nModel data is missing or stale. " +
"Run `npm run hydrate:model-data` from the repository root.",
);
process.exitCode = 1;
}
实际校验不只是看目录存不存在,还会检查:
- provider JSON 文件是否齐全;
.manifest.json的 schema 版本;- 模型目录结构哈希;
- 每个文件的 SHA-256;
- 模型的
id、provider、api是否匹配; baseUrl、reasoning、输入模态、上下文窗口、最大输出和价格字段是否合法。
看到下面的输出才说明这一步完成:1
Generated model data is valid.
三、配置模型和密钥
Pi 默认读取的全局目录是:1
~/.pi/agent/
经常使用 Pi 的机器上一般已经有 auth.json、models.json 和相关设置,所以直接从源码启动时也能读取已有配置。sdk.ts 中创建模型运行时的代码如下:1
2
3
4const authPath = options.agentDir ? join(agentDir, "auth.json") : undefined;
const modelsPath = options.agentDir ? join(agentDir, "models.json") : undefined;
const modelRuntime =
options.modelRuntime ?? (await ModelRuntime.create({ authPath, modelsPath }));
如果 options.agentDir 没有显式传入,ModelRuntime.create() 内部仍会回退到默认路径:1
2
3
4
5const modelsPath =
options.modelsPath === null
? undefined
: (options.modelsPath ?? join(getAgentDir(), "models.json"));
const config = await ModelConfig.load(modelsPath);
一个 OpenAI 兼容服务的配置例子
在 ~/.pi/agent/models.json 中可以增加: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{
"providers": {
"my-openai-compatible": {
"baseUrl": "https://example.com/v1",
"api": "openai-completions",
"apiKey": "$PI_DEBUG_API_KEY",
"models": [
{
"id": "my-model-id",
"name": "My Debug Model",
"reasoning": true,
"input": ["text"],
"contextWindow": 128000,
"maxTokens": 16384,
"cost": {
"input": 0,
"output": 0,
"cacheRead": 0,
"cacheWrite": 0
}
}
]
}
}
}
apiKey 支持直接写值、引用环境变量或执行命令。为了避免密钥进入 Git,调试时建议使用 $PI_DEBUG_API_KEY 这种环境变量引用。
四、配置 VS Code 调试入口
源码入口是 packages/coding-agent/src/cli.ts。项目使用 TypeScript ESM,Node 通过 tsx 直接执行源码,因此一个单次请求的调试配置可以写成: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{
"version": "0.2.0",
"configurations": [
{
"name": "Pi:单次 Prompt,只读工具",
"type": "node",
"request": "launch",
"runtimeExecutable": "node",
"runtimeArgs": ["--import", "tsx"],
"program": "${workspaceFolder}/packages/coding-agent/src/cli.ts",
"cwd": "${workspaceFolder}",
"console": "integratedTerminal",
"internalConsoleOptions": "neverOpen",
"sourceMaps": true,
"smartStep": true,
"autoAttachChildProcesses": true,
"skipFiles": [
"<node_internals>/**",
"${workspaceFolder}/node_modules/**"
],
"args": [
"--print",
"读取 package.json,并告诉我项目名称",
"--tools",
"read",
"--no-session",
"--no-extensions",
"--no-skills",
"--no-prompt-templates",
"--no-context-files"
],
"env": {
"PI_DEBUG_API_KEY": "${env:PI_DEBUG_API_KEY}"
}
}
]
}
如果 VS Code 本身能继承终端中的环境变量,env 可以省略。macOS 上从图形界面启动 VS Code 时,经常读不到 shell 中设置的变量,这时可以显式增加 env,或者使用只保存在本机且被 .gitignore 排除的 envFile。
这里选择单次 Prompt 模式而不是一上来调 TUI,原因很实际:
- 请求固定,方便反复复现;
- 只开放
read,不会因为模型输出而修改工作区; - 禁用 session、extension、skill 和上下文文件后,调用链更短;
- 请求结束进程就退出,断点更容易观察。
环境跑通以后,再切换到仓库自带的“Pi:交互模式”配置分析 steering、follow-up 和 TUI。
建议先在这些位置打断点:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19packages/coding-agent/src/core/sdk.ts
createAgentSession()
new Agent({...})
streamFn: async (...) => {...}
packages/coding-agent/src/core/agent-session.ts
AgentSession.prompt()
_runAgentPrompt()
packages/agent/src/agent.ts
Agent.prompt()
runPromptMessages()
processEvents()
packages/agent/src/agent-loop.ts
runAgentLoop()
runLoop()
streamAssistantResponse()
executeToolCalls()
五、一次请求的完整调用链
先给出总链路,后面再逐层展开:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17packages/coding-agent/src/cli.ts
└─ main(process.argv.slice(2))
├─ createAgentSessionRuntime(...)
│ └─ createAgentSession(...)
│ ├─ ModelRuntime.create(...)
│ └─ new Agent({ streamFn, convertToLlm, ... })
└─ runPrintMode(...) / InteractiveMode.run()
└─ AgentSession.prompt(...)
└─ Agent.prompt(...)
└─ runAgentLoop(...)
└─ runLoop(...)
├─ streamAssistantResponse(...)
│ └─ streamFunction(...)
│ └─ ModelRuntime.streamSimple(...)
│ └─ provider.streamSimple(...)
└─ executeToolCalls(...)
└─ tool.execute(...)
CLI 入口本身很薄:1
2
3
4
5
6
import { setupCli } from "./cli/setup.ts";
import { main } from "./main.ts";
setupCli();
main(process.argv.slice(2));
单次 Prompt 模式最后调用 runPrintMode(),里面再进入 AgentSession:1
2
3
4
5
6
7if (initialMessage) {
await session.prompt(initialMessage, { images: initialImages });
}
for (const message of messages) {
await session.prompt(message);
}
六、Agent 初始化:streamFn 在 sdk.ts 中注入
1. 先安装一个兼容性默认值
packages/coding-agent/src/core/sdk.ts 文件顶部有:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17import {
Agent,
type AgentMessage,
setDefaultStreamFn,
type ThinkingLevel,
} from "@earendil-works/pi-agent-core";
import {
clampThinkingLevel,
type Message,
type Model,
streamSimple,
} from "@earendil-works/pi-ai/compat";
// Preserve the pre-0.81 fallback for extensions that construct Agent instances
// or invoke low-level agent loops without supplying streamFn. Agent core remains
// provider-agnostic and does not import pi-ai/compat itself.
setDefaultStreamFn(streamSimple);
这一步是兼容旧扩展的全局 fallback。Agent Core 自己不依赖 pi-ai/compat,由 Coding Agent 在启动时把默认流函数注册进去。
packages/agent/src/stream-fn.ts 的全部核心逻辑如下:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17import type { StreamFn } from "./types.ts";
let defaultStreamFn: StreamFn | undefined;
export function setDefaultStreamFn(streamFn: StreamFn | undefined): void {
defaultStreamFn = streamFn;
}
export function getDefaultStreamFn(): StreamFn {
if (!defaultStreamFn) {
throw new Error(
"No default stream function configured. " +
"Pass streamFn explicitly or call setDefaultStreamFn().",
);
}
return defaultStreamFn;
}
2. 实际 Agent 使用的是显式注入的 streamFn
创建会话时,sdk.ts 并不是只依赖上面的 fallback,而是给 new Agent() 显式传入一个包装函数: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
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114agent = new Agent({
// initialState 只是 Agent Core 的初始快照。
// systemPrompt 和 tools 暂时为空,AgentSession 随后会根据资源加载结果填充。
initialState: {
systemPrompt: "",
model,
thinkingLevel,
tools: [],
},
// Coding Agent 中除了 user/assistant/toolResult,还有扩展产生的 custom message。
// 这里负责把产品层消息转成 pi-ai 接受的 Message,并按设置屏蔽图片。
convertToLlm: convertToLlmWithBlockImages,
// 这是本文最关键的注入点。
// Agent Core 只保存此函数的引用,不知道认证文件、供应商 SDK 和重试配置。
streamFn: async (model, context, options) => {
// 每次请求时动态读取设置,所以运行中修改配置也能影响后续 turn。
const providerRetrySettings = settingsManager.getProviderRetrySettings();
const httpIdleTimeoutMs = settingsManager.getHttpIdleTimeoutMs();
// SDKs treat timeout=0 as 0ms (immediate timeout), not "no timeout".
// Use max int32 to effectively disable the timeout.
const effectiveTimeoutMs =
httpIdleTimeoutMs === 0 ? 2147483647 : httpIdleTimeoutMs;
// 优先级:单次调用参数 > provider 重试配置 > 全局 HTTP 空闲超时。
const timeoutMs =
options?.timeoutMs ??
providerRetrySettings.timeoutMs ??
effectiveTimeoutMs;
// WebSocket 连接超时独立于普通 HTTP 空闲超时。
const websocketConnectTimeoutMs =
options?.websocketConnectTimeoutMs ??
settingsManager.getWebSocketConnectTimeoutMs();
// extensionRunnerRef 在 Agent 创建后由 AgentSession 绑定。
// 用 ref 是为了让闭包始终拿到当前 runner,而不是初始化时的 undefined。
const headerRunner = extensionRunnerRef.current;
// 真正进入模型运行时。ModelRuntime 会解析 provider、认证与协议实现。
return modelRuntime.streamSimple(model, context, {
// 先保留 Agent Core 传来的 signal、apiKey、thinking 等请求选项。
...options,
timeoutMs,
websocketConnectTimeoutMs,
// 单次请求可覆盖默认重试策略;未传入才读取 SettingsManager。
maxRetries:
options?.maxRetries ?? providerRetrySettings.maxRetries,
maxRetryDelayMs:
options?.maxRetryDelayMs ??
providerRetrySettings.maxRetryDelayMs,
// provider 即将发送请求前执行。先补 Pi 的来源/session header,
// 再让扩展做最后修改,扩展因此拥有最终 header 的控制权。
transformHeaders: async (requestHeaders) => {
const headers = mergeProviderAttributionHeaders(
model,
settingsManager,
options?.sessionId,
requestHeaders,
);
return headerRunner?.hasHandlers("before_provider_headers")
? headerRunner.emitBeforeProviderHeaders(headers ?? {})
: (headers ?? {});
},
});
},
// onPayload 处理已经转换成供应商协议、但尚未真正发出的请求体。
// 没有扩展监听时直接原样返回,避免无意义复制。
onPayload: async (payload, _model) => {
const runner = extensionRunnerRef.current;
if (!runner?.hasHandlers("before_provider_request")) {
return payload;
}
return runner.emitBeforeProviderRequest(payload);
},
// onResponse 暴露 HTTP 状态与响应头,而正文仍由流事件处理。
onResponse: async (response, _model) => {
const runner = extensionRunnerRef.current;
if (!runner?.hasHandlers("after_provider_response")) {
return;
}
await runner.emit({
type: "after_provider_response",
status: response.status,
headers: response.headers,
});
},
// sessionId 会一路传到 provider,可供支持会话亲和或缓存的后端使用。
sessionId: sessionManager.getSessionId(),
// 在 convertToLlm 之前运行扩展的上下文 hook。
// 它操作的仍然是 AgentMessage[],适合裁剪、插入或改写上下文。
transformContext: async (messages) => {
const runner = extensionRunnerRef.current;
if (!runner) return messages;
return runner.emitContext(messages);
},
// steering 选择何时向进行中的任务插话;followUp 决定后续任务如何出队。
steeringMode: settingsManager.getSteeringMode(),
followUpMode: settingsManager.getFollowUpMode(),
// 这些值不会在 sdk.ts 中消费,而是继续透传给 Agent Loop 和 provider。
transport: settingsManager.getTransport(),
thinkingBudgets: settingsManager.getThinkingBudgets(),
maxRetryDelayMs: settingsManager.getProviderRetrySettings().maxRetryDelayMs,
});
这个包装函数不只是把请求转发出去,它还在真正请求前补上:
- HTTP 空闲超时和 WebSocket 连接超时;
- 最大重试次数、最大重试等待时间;
- provider attribution headers;
- 扩展提供的
before_provider_headershook; - 请求 payload 与响应 headers 的扩展 hook。
因此 streamFn 是 Agent Core 与真实模型运行时之间的注入点,也是调试“最终请求为什么长这样”的最佳断点之一。
3. StreamFn 的契约
它不是一个普通的 Promise<AssistantMessage>,而是返回事件流:1
2
3
4
5export type StreamFn = (
model: Model<Api>,
context: Context,
options?: SimpleStreamOptions,
) => AssistantMessageEventStream | Promise<AssistantMessageEventStream>;
源码注释还规定,请求失败不应该直接 throw,而要在流中产生错误事件,并以 stopReason: "error" 或 "aborted" 的最终 AssistantMessage 收尾。这样上层可以用统一的事件协议处理正常完成、失败和取消。
4. ModelRuntime 再选择具体 Provider
sdk.ts 注入的函数最终来到 ModelRuntime.streamSimple():1
2
3
4
5
6
7
8
9
10
11
12
13
14streamSimple(
model: Model<Api>,
context: Context,
options?: ModelsSimpleStreamOptions,
): AssistantMessageEventStream {
return lazyStream(model, async () => {
const prepared = await this.prepareRequest(model, options);
return prepared.provider.streamSimple(
prepared.model,
context,
prepared.options as SimpleStreamOptions,
);
});
}
到这里才根据 model 找到对应 provider,准备认证和请求参数,然后调用各 provider 的 streamSimple。agent-loop.ts 完全不用知道下层究竟是 OpenAI Responses、OpenAI Completions、Anthropic Messages 还是 Google Generative AI。
七、从 AgentSession 进入 Agent Core
AgentSession.prompt() 在进入 Core 前会处理很多产品层能力,包括扩展命令、Prompt 模板、Skill、模型与认证检查、压缩检查以及扩展 hook。最后构造用户消息并调用:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16private async _runAgentPrompt(
messages: AgentMessage | AgentMessage[],
): Promise<void> {
this._isAgentRunActive = true;
try {
await this.agent.prompt(messages);
while (await this._handlePostAgentRun()) {
await this.agent.continue();
}
} finally {
this._systemPromptOverride = undefined;
this._flushPendingBashMessages();
this._flushPendingCustomMessages();
await this._emitAgentSettled();
}
}
这里还有一层“运行结束后的继续处理”,主要服务于 Coding Agent 会话层的压缩等逻辑。不要把它和 agent-loop.ts 内部的双层循环混在一起。
进入 Agent 后,字符串会先被规范化为 AgentMessage[],然后调用 runAgentLoop():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
42async prompt(message: AgentMessage | AgentMessage[]): Promise<void>;
async prompt(input: string, images?: ImageContent[]): Promise<void>;
async prompt(
input: string | AgentMessage | AgentMessage[],
images?: ImageContent[],
): Promise<void> {
// activeRun 同时充当并发保护:一个 Agent 实例不能并行跑两个 prompt。
if (this.activeRun) {
throw new Error(
"Agent is already processing a prompt. Use steer() or followUp() " +
"to queue messages, or wait for completion.",
);
}
// string、单个 AgentMessage 和 AgentMessage[] 最终统一成数组。
// 字符串会被包装为 role=user,并附加 timestamp 与可选图片。
const messages = this.normalizePromptInput(input, images);
await this.runPromptMessages(messages);
}
private async runPromptMessages(
messages: AgentMessage[],
options: { skipInitialSteeringPoll?: boolean } = {},
): Promise<void> {
// runWithLifecycle 创建 AbortController、设置 isStreaming,
// 并保证成功、异常或取消后都能清理 activeRun。
await this.runWithLifecycle(async (signal) => {
await runAgentLoop(
// 本次新输入;runAgentLoop 会把它拼到历史上下文末尾。
messages,
// 使用数组副本,避免低层循环直接修改 Agent.state 的顶层数组。
this.createContextSnapshot(),
// 把模型、hook、队列读取器和执行策略冻结为本次运行配置。
this.createLoopConfig(options),
// 主循环只 emit;Agent 通过 processEvents 归并自己的状态。
(event) => this.processEvents(event),
// 同一个 AbortSignal 贯穿模型请求、hook 与工具调用。
signal,
// 即 sdk.ts 构造 Agent 时显式注入的 streamFn。
this.streamFunction,
);
});
}
注意最后一个参数 this.streamFunction,它就是 sdk.ts 创建 Agent 时传进来的函数。到这里,初始化阶段的依赖注入和运行阶段的实际调用闭环了。
八、runAgentLoop:准备上下文和事件
runAgentLoop() 没有立刻发请求,而是先复制上下文、发出生命周期事件: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
38export async function runAgentLoop(
prompts: AgentMessage[],
context: AgentContext,
config: AgentLoopConfig,
emit: AgentEventSink,
signal: AbortSignal | undefined,
streamFn: StreamFn,
): Promise<AgentMessage[]> {
// newMessages 是本次运行的增量结果,初始就包含用户新提交的 prompts。
const newMessages: AgentMessage[] = [...prompts];
// currentContext 是本次循环使用的工作副本:保留旧历史,再拼接新输入。
// 展开 context 只浅复制对象,因此 messages 必须单独创建新数组。
const currentContext: AgentContext = {
...context,
messages: [...context.messages, ...prompts],
};
// 事件顺序是协议的一部分:先开始 Agent 和 turn,再完整发布用户消息。
await emit({ type: "agent_start" });
await emit({ type: "turn_start" });
for (const prompt of prompts) {
await emit({ type: "message_start", message: prompt });
await emit({ type: "message_end", message: prompt });
}
// 正常的 Coding Agent 总会显式传 streamFn;右侧是旧调用方的兼容 fallback。
await runLoop(
currentContext,
newMessages,
config,
signal,
emit,
streamFn ?? getDefaultStreamFn(),
);
// runLoop 内部通过 push 修改 newMessages,这里返回的就是本次运行完整增量。
return newMessages;
}
这里有两个消息集合:
currentContext.messages是发给模型的完整上下文;newMessages只记录本次 Agent Run 新产生的消息,最后随agent_end返回。
这两个集合在循环中同时追加,但语义不同。调试时如果只盯一个数组,很容易误以为消息被重复添加。
九、核心:两层 Agent Loop 到底在循环什么
下面是 runLoop() 的核心代码。片段较长,但只有结合真实代码才能看清两个循环的退出条件: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
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178async function runLoop(
initialContext: AgentContext,
newMessages: AgentMessage[],
initialConfig: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
streamFunction: StreamFn,
): Promise<void> {
// 这两个变量都可能被 prepareNextTurn 替换,所以不能使用 const。
let currentContext = initialContext;
let config = initialConfig;
// 第一轮之前没有已完成的 turn;从第二轮起它携带上一轮快照。
let lastCompletedTurn: PrepareNextTurnContext | undefined;
// 进入循环前先拉取一次 steering。例如 Agent 真正开始执行前,
// 用户已经排入了一条纠偏消息,它也应进入第一次请求。
let pendingMessages: AgentMessage[] =
(await config.getSteeringMessages?.()) || [];
// 外层:Agent 本来要停下时,如果还有 follow-up,就重新进入内层
while (true) {
// 初始为 true,保证没有排队消息时也至少请求一次 LLM。
// 后续它表示上一条 assistant message 是否留下了待反馈的工具结果。
let hasMoreToolCalls = true;
// 内层:只要还有工具结果需要交给模型,或有 steering 消息,就继续
while (hasMoreToolCalls || pendingMessages.length > 0) {
// 第一轮会跳过。第二轮起,在每次 LLM 请求前给上层一次机会,
// 用上一轮完整结果更新上下文、模型或 thinking level。
if (lastCompletedTurn) {
const nextTurnSnapshot =
await config.prepareNextTurn?.(lastCompletedTurn);
if (nextTurnSnapshot) {
// hook 可以只返回需要替换的字段,其余部分继续沿用。
currentContext =
nextTurnSnapshot.context ?? currentContext;
config = {
...config,
model: nextTurnSnapshot.model ?? config.model,
reasoning:
nextTurnSnapshot.thinkingLevel === undefined
? config.reasoning
: nextTurnSnapshot.thinkingLevel === "off"
? undefined
: nextTurnSnapshot.thinkingLevel,
};
}
// prepareNextTurn 可能在做耗时的上下文压缩。压缩期间到达的
// steering 也要进入下一轮;已有 pending 时不再拉取,避免
// one-at-a-time 模式在一个 turn 中取出两条消息。
if (pendingMessages.length === 0) {
pendingMessages =
(await config.getSteeringMessages?.()) || [];
}
// turn_start 放在准备结束后,表示下一轮现在才正式开始。
await emit({ type: "turn_start" });
}
// 在下一次 assistant 响应前注入 steering 消息
if (pendingMessages.length > 0) {
for (const message of pendingMessages) {
// 先发标准消息事件供 Session/UI 消费,再同步写入两个数组。
await emit({ type: "message_start", message });
await emit({ type: "message_end", message });
currentContext.messages.push(message);
newMessages.push(message);
}
// 已转移到上下文,必须清空,防止下个 turn 重复注入。
pendingMessages = [];
}
// 一次真正的 LLM 请求与流式响应消费。函数返回时,
// currentContext 最后一项已经从 partial message 变成最终消息。
const message = await streamAssistantResponse(
currentContext,
config,
signal,
emit,
streamFunction,
);
// streamAssistantResponse 已维护完整上下文;这里只记录本次增量。
newMessages.push(message);
// provider 将请求失败与取消编码进最终 AssistantMessage。
// 这两种情况直接正常发出结束事件,不再执行工具或读取队列。
if (
message.stopReason === "error" ||
message.stopReason === "aborted"
) {
await emit({ type: "turn_end", message, toolResults: [] });
await emit({ type: "agent_end", messages: newMessages });
return;
}
const toolCalls = message.content.filter(
(c) => c.type === "toolCall",
);
// turn_end 始终需要统一的数组结构,所以无工具时也创建空数组。
const toolResults: ToolResultMessage[] = [];
// 先假设工具链结束。只有确实执行过工具且批次未要求终止,
// 才会在下面把它恢复为 true。
hasMoreToolCalls = false;
if (toolCalls.length > 0) {
// length 表示 assistant 输出被 token 上限截断。
// 此时所有 toolCall 参数都不可信,即使其中某个 JSON 可以解析。
const executedToolBatch =
message.stopReason === "length"
? await failToolCallsFromTruncatedMessage(
toolCalls,
emit,
)
: await executeToolCalls(
currentContext,
message,
config,
signal,
emit,
);
toolResults.push(...executedToolBatch.messages);
// 普通工具结果必须再交给 LLM,总结结果或决定下一步。
// 只有工具批次明确 terminate 时,才不自动继续这一工具链。
hasMoreToolCalls = !executedToolBatch.terminate;
for (const result of toolResults) {
// 完整上下文决定下次发给 LLM 什么;newMessages 只负责
// 本次 Agent Run 的返回值和 agent_end 事件。
currentContext.messages.push(result);
newMessages.push(result);
}
}
// 一个 turn 包含一次 assistant 响应及其引发的整批工具执行,
// 所以 turn_end 必须等所有 toolResults 准备好后再发送。
await emit({ type: "turn_end", message, toolResults });
// 下一轮的 hook 需要同时看到响应、工具结果、完整上下文和运行增量。
lastCompletedTurn = {
message,
toolResults,
context: currentContext,
newMessages,
};
// “优雅停止”发生在当前 turn 完整结束后,优先于队列轮询。
// 因而不会中断已开始的工具,但会结束本次 Agent Run。
if (await config.shouldStopAfterTurn?.(lastCompletedTurn)) {
await emit({ type: "agent_end", messages: newMessages });
return;
}
// 模型响应与工具执行期间新到达的 steering 在这里读取。
// 非空时,下一次 while 条件成立并开始一个新 turn。
pendingMessages =
(await config.getSteeringMessages?.()) || [];
}
// 内层结束,Agent 本来准备停止;再检查延后执行的 follow-up
const followUpMessages =
(await config.getFollowUpMessages?.()) || [];
if (followUpMessages.length > 0) {
// 复用 pendingMessages 的标准注入路径,使 follow-up 同样产生
// message_start/message_end,并进入上下文与本次运行增量。
pendingMessages = followUpMessages;
// 回到外层会重新令 hasMoreToolCalls=true,所以 follow-up
// 至少触发一次 LLM 请求,而不只是被追加到历史。
continue;
}
// 工具链、steering 和 follow-up 都耗尽,才退出整个 Agent Run。
break;
}
// agent_end 是本次 run 的最后一个协议事件,只携带本次新增消息。
await emit({ type: "agent_end", messages: newMessages });
}
1. 内层循环:完成当前任务链
内层条件是:1
while (hasMoreToolCalls || pendingMessages.length > 0)
初始时 hasMoreToolCalls = true,所以至少请求模型一次。
假设用户说“读取 package.json 并告诉我项目名称”,一次典型轨迹是:1
2
3
4
5
6
7
8
9
10
11第一次 LLM 请求
→ assistant 返回 read 工具调用
→ 执行 read
→ 把 toolResult 追加到 currentContext.messages
→ hasMoreToolCalls = true
第二次 LLM 请求
→ 模型看到 read 的结果
→ assistant 返回最终文本,不再调用工具
→ hasMoreToolCalls = false
→ 如果也没有 steering,退出内层循环
所以工具调用后的“下一轮”不是递归,而是同一个内层 while 的下一次迭代。
steering 也在内层处理。它的含义是:当前 Agent 还在工作时,用户发来一条纠偏消息。这条消息不会强行中断正在消费的模型响应或正在执行的工具,而是在当前 turn 完整结束后、下一次 LLM 请求前注入上下文。
2. 外层循环:处理 Agent 停止后的 follow-up
当工具链结束且没有 steering 时,内层循环退出。此时 Agent 原本应该结束,但外层会再查询一次 getFollowUpMessages()。
followUp 与 steering 的差别是:
steering参与当前任务链,尽快在下一 turn 纠偏;followUp等当前任务自然完成后再开始,更像排队的下一项任务。
如果有 follow-up,代码把它放进 pendingMessages,通过 continue 回到外层顶部,再进入内层循环。没有 follow-up 才真正 break 并发出 agent_end。
3. prepareNextTurn 为什么放在请求前
从第二个 turn 开始,代码会调用 prepareNextTurn(lastCompletedTurn)。它允许上层在下一次请求前替换:
- context,例如完成上下文压缩;
- model,例如运行途中切换模型;
- thinking level。
这就是为什么它必须处在 streamAssistantResponse() 之前,而且只在已有 lastCompletedTurn 时执行。
4. length 时为什么不执行工具
如果模型因为输出 token 上限而以 stopReason === "length" 结束,工具参数也可能被截断。即使残缺 JSON 被宽松解析器碰巧修复成功,也不能保证语义完整。因此代码不会执行任何一个工具,而是给每个调用生成错误结果,让模型在下一轮重新发起完整调用。
这是一个很容易被简化版 Agent Loop 忽略的安全细节。
十、一次 LLM 流式响应是怎样被消费的
streamAssistantResponse() 是 Agent Core 与模型流的直接边界: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
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137async function streamAssistantResponse(
context: AgentContext,
config: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
streamFunction: StreamFn,
): Promise<AssistantMessage> {
// context.messages 是循环内的权威历史。transformContext 不直接写回它,
// 而是生成“仅用于本次请求”的视图,可用于压缩、裁剪或扩展上下文。
// AgentMessage[] → AgentMessage[],此时仍允许包含产品层 custom message。
let messages = context.messages;
if (config.transformContext) {
messages = await config.transformContext(messages, signal);
}
// AgentMessage[] → pi-ai Message[]。UI 通知等模型不认识的消息在这里过滤,
// custom message 也必须在这里变换成 user/assistant/toolResult 之一。
const llmMessages = await config.convertToLlm(messages);
// 把 system prompt、已转换消息和工具 schema 组成一次完整模型输入。
// 注意这里传的是工具“定义”,真正的 execute 函数仍只保留在 Agent Context。
const llmContext: Context = {
systemPrompt: context.systemPrompt,
messages: llmMessages,
tools: context.tools,
};
// 每次请求都重新取 key,兼容会过期和刷新的 OAuth token。
// getApiKey 的动态结果优先;没有结果时才回退到 config.apiKey。
const resolvedApiKey =
(config.getApiKey
? await config.getApiKey(config.model.provider)
: undefined) || config.apiKey;
// streamFunction 返回可异步迭代的 AssistantMessageEventStream,
// 调用到这里并不代表完整回答已经到达。
const response = await streamFunction(config.model, llmContext, {
// config 继承 SimpleStreamOptions,因此请求 hook、thinking、transport 等
// 选项先整体透传,再用下面两个字段覆盖成当前调用的最终值。
...config,
apiKey: resolvedApiKey,
signal,
});
// partialMessage 始终指向“截至当前事件已经拼装出的完整消息”,
// 而不是单独的 token/delta。
let partialMessage: AssistantMessage | null = null;
// 标记 context.messages 中是否已经为这个回答占了最后一个位置。
// 它用于兼容没有 start 事件、直接结束的 provider 流。
let addedPartial = false;
// Provider 事件按产生顺序异步到达;for await 自然形成背压,
// 当前事件处理和 emit 没结束前不会消费下一事件。
for await (const event of response) {
switch (event.type) {
case "start":
// start.partial 通常只有 assistant 消息的骨架和空 content。
partialMessage = event.partial;
// 先把骨架放到历史末尾,为后面的 delta 保留固定槽位。
context.messages.push(partialMessage);
addedPartial = true;
// 展开成新对象,避免监听者持有随后会变化的同一对象引用。
await emit({
type: "message_start",
message: { ...partialMessage },
});
break;
case "text_start":
case "text_delta":
case "text_end":
case "thinking_start":
case "thinking_delta":
case "thinking_end":
case "toolcall_start":
case "toolcall_delta":
case "toolcall_end":
// 所有内容事件走同一条归并路径。event.partial 已由 pi-ai
// 累积好当前文本、thinking 或 toolCall,不需要 Core 自己拼 token。
if (partialMessage) {
partialMessage = event.partial;
// 替换最后一项而非 push:一个回答无论包含多少 delta,
// 在上下文中始终只占一条 assistant message。
context.messages[context.messages.length - 1] =
partialMessage;
// 原始 provider 事件一并向上传递,TUI 可区分文本、思考和工具参数。
await emit({
type: "message_update",
assistantMessageEvent: event,
message: { ...partialMessage },
});
}
break;
case "done":
case "error": {
// result() 给出标准化后的最终 AssistantMessage,包含 usage、
// stopReason、errorMessage 以及最终 content。
const finalMessage = await response.result();
if (addedPartial) {
// 正常流:用最终对象替换之前反复更新的临时槽位。
context.messages[context.messages.length - 1] =
finalMessage;
} else {
// 防御非常规 provider:没有 start 事件时在结束处补入历史。
context.messages.push(finalMessage);
}
if (!addedPartial) {
// 保持 Agent 事件协议完整:即使 provider 没 start,
// 上层仍会按顺序收到 message_start → message_end。
await emit({
type: "message_start",
message: { ...finalMessage },
});
}
// message_end 才表示这条消息可以进入 Agent.state.messages。
await emit({ type: "message_end", message: finalMessage });
return finalMessage;
}
}
}
// 有些实现通过结束 AsyncIterable 表示完成,不额外发 done。
// 这个兜底分支与上面的 done/error 分支保持完全相同的落库和事件语义。
const finalMessage = await response.result();
if (addedPartial) {
context.messages[context.messages.length - 1] = finalMessage;
} else {
context.messages.push(finalMessage);
await emit({
type: "message_start",
message: { ...finalMessage },
});
}
await emit({ type: "message_end", message: finalMessage });
return finalMessage;
}
这里可以看到三层事件:
- Provider 流事件:
text_delta、thinking_delta、toolcall_delta等; - Agent 事件:
message_start、message_update、message_end; - 更外层 UI 订阅 Agent 事件,刷新终端内容和会话状态。
收到 start 时,代码先把 partial message 放进 context.messages。后续每个 delta 都不是继续 push 新消息,而是替换数组最后一项。否则每个 token 都会变成一条历史消息。
收到 done 或 error 后,再通过 response.result() 取得最终对象,用它替换 partial message。末尾的兜底逻辑用于处理流正常耗尽但没有显式 done/error 事件的情况。
十一、工具调用:选择、校验、执行与回填
模型响应完成后,主循环从 message.content 中筛出工具调用。实际执行入口如下: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
42async function executeToolCalls(
currentContext: AgentContext,
assistantMessage: AssistantMessage,
config: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
): Promise<ExecutedToolCallBatch> {
// 再从最终 assistant message 中提取一次,确保执行的是已经完成流式拼装的调用。
const toolCalls = assistantMessage.content.filter(
(c) => c.type === "toolCall",
);
// 工具自身可声明 sequential。例如修改同一状态的工具不能安全并发。
// 只要批次中有一个这样的工具,整批都按串行处理,避免部分交错。
const hasSequentialToolCall = toolCalls.some(
(tc) =>
currentContext.tools?.find((t) => t.name === tc.name)
?.executionMode === "sequential",
);
// 全局策略和工具局部声明任一要求串行,就走串行实现。
if (config.toolExecution === "sequential" || hasSequentialToolCall) {
return executeToolCallsSequential(
currentContext,
assistantMessage,
toolCalls,
config,
signal,
emit,
);
}
// 默认并行路径仍会先顺序执行参数校验和 beforeToolCall,
// 只有通过预检的 tool.execute 才会并发。
return executeToolCallsParallel(
currentContext,
assistantMessage,
toolCalls,
config,
signal,
emit,
);
}
默认策略是并行,但只要本批次任意工具声明 executionMode === "sequential",整批就改为串行。这样有副作用或有顺序依赖的工具可以主动要求保守执行。
每个工具真正运行前还要经过 prepareToolCall():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
82async function prepareToolCall(
currentContext: AgentContext,
assistantMessage: AssistantMessage,
toolCall: AgentToolCall,
config: AgentLoopConfig,
signal: AbortSignal | undefined,
): Promise<PreparedToolCall | ImmediateToolCallOutcome> {
// 模型只能按名称调用当前上下文公开的工具,不能凭空调用任意函数。
const tool = currentContext.tools?.find(
(t) => t.name === toolCall.name,
);
if (!tool) {
// “工具不存在”也规范化成即时结果,让 LLM 有机会改正名称。
return {
kind: "immediate",
result: createErrorToolResult(
`Tool ${toolCall.name} not found`,
),
isError: true,
};
}
try {
// prepareArguments 可做兼容性归一化;随后必须按工具 schema 校验。
// validatedArgs 才是最终允许交给 execute 的参数。
const preparedToolCall = prepareToolCallArguments(tool, toolCall);
const validatedArgs = validateToolArguments(tool, preparedToolCall);
if (config.beforeToolCall) {
// hook 能看到 assistant 原消息、原 toolCall、已校验参数和上下文,
// 可用于权限确认、策略限制、日志或参数审计。
const beforeResult = await config.beforeToolCall(
{
assistantMessage,
toolCall,
args: validatedArgs,
context: currentContext,
},
signal,
);
// hook 等待期间用户可能取消,因此返回后立即再次检查 signal。
if (signal?.aborted) {
return {
kind: "immediate",
result: createErrorToolResult("Operation aborted"),
isError: true,
};
}
if (beforeResult?.block) {
// block 不抛异常,而是产生模型可读的错误 toolResult。
const result = createErrorToolResult(
beforeResult.reason || "Tool execution was blocked",
);
if (beforeResult.terminate === true) {
// terminate 会参与整个批次的终止判断,而不是只标记 UI。
result.terminate = true;
}
return { kind: "immediate", result, isError: true };
}
}
// prepared 是预检与执行之间的边界:只有这种结果会调用 tool.execute。
return {
kind: "prepared",
toolCall,
tool,
args: validatedArgs,
};
} catch (error) {
// 参数预处理、schema 校验和 hook 抛错都降级为即时错误结果,
// 保证 Agent Loop 仍能把错误放回对话,而不是破坏事件序列。
return {
kind: "immediate",
result: createErrorToolResult(
error instanceof Error ? error.message : String(error),
),
isError: true,
};
}
}
它完成四项检查:
- 工具是否存在;
- 是否需要预处理参数;
- 参数是否符合工具 schema;
beforeToolCallhook 是否阻止执行。
工具报错通常也不会炸掉整个 Agent Loop,而是转成 isError: true 的 toolResult。这样模型下一轮能看见错误,并决定修正参数、换一种方法或直接向用户解释失败。
工具执行结束后生成标准消息:1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17function createToolResultMessage(
finalized: FinalizedToolCallOutcome,
): ToolResultMessage {
return {
role: "toolResult",
toolCallId: finalized.toolCall.id,
toolName: finalized.toolCall.name,
content: finalized.result.content ?? [],
details: finalized.result.details,
usage: finalized.result.usage,
...(finalized.result.addedToolNames?.length
? { addedToolNames: finalized.result.addedToolNames }
: {}),
isError: finalized.isError,
timestamp: Date.now(),
};
}
随后 runLoop() 把它加入 currentContext.messages。下一次调用 streamAssistantResponse() 时,convertToLlm() 会把这个 toolResult 转成模型协议需要的工具结果消息。至此完成闭环:1
2
3
4
5
6
7
8assistant toolCall
→ 参数准备与校验
→ beforeToolCall
→ tool.execute
→ afterToolCall
→ ToolResultMessage
→ 加入上下文
→ 下一次 LLM 请求
十二、事件如何更新 Agent 状态
主循环只负责 emit,真正维护 Agent.state 的是 agent.ts 中的 processEvents():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
66private async processEvents(event: AgentEvent): Promise<void> {
// 先归并 Agent 自己的可观察状态,再通知外部 listener。
// 因此外部收到事件时,读取 state 已经能看到与该事件一致的新值。
switch (event.type) {
case "message_start":
// streamingMessage 是正在生成或刚开始发布的临时消息,
// 此时还不能放入正式历史,否则取消或 delta 会留下重复记录。
this._state.streamingMessage = event.message;
break;
case "message_update":
// 每个 delta 都用新的完整 partial snapshot 覆盖旧值。
this._state.streamingMessage = event.message;
break;
case "message_end":
// 只有结束事件才把消息追加到 Agent 的权威 transcript。
// agent-loop 内部的 currentContext 是独立数组,所以这里不是重复 push。
this._state.streamingMessage = undefined;
this._state.messages.push(event.message);
break;
case "tool_execution_start": {
// 创建新 Set 而不是原地 add,保证观察 state 引用变化的 UI 能刷新。
const pendingToolCalls = new Set(
this._state.pendingToolCalls,
);
pendingToolCalls.add(event.toolCallId);
this._state.pendingToolCalls = pendingToolCalls;
break;
}
case "tool_execution_end": {
// 与 start 对称;并行模式下多个工具会按实际完成顺序依次移除。
const pendingToolCalls = new Set(
this._state.pendingToolCalls,
);
pendingToolCalls.delete(event.toolCallId);
this._state.pendingToolCalls = pendingToolCalls;
break;
}
case "turn_end":
// 普通工具错误保存在 toolResult.isError;这里只记录 assistant 级错误。
if (event.message.role === "assistant" && event.message.errorMessage) {
this._state.errorMessage = event.message.errorMessage;
}
break;
case "agent_end":
// agent_end 是事件流终点,但 activeRun 要等所有 listener 完成后
// 才会由 finishRun 清理,因此这里仅清空临时消息。
this._state.streamingMessage = undefined;
break;
}
// 所有事件都必须发生在一次 activeRun 内,并共享这次运行的取消信号。
const signal = this.activeRun?.abortController.signal;
if (!signal) {
throw new Error("Agent listener invoked outside active run");
}
// 串行 await 保证 listener 看到严格事件顺序,同时也会形成背压。
for (const listener of this.listeners) {
await listener(event, signal);
}
}
这是一种很清楚的单向数据流:主循环产生事实事件,Agent 根据事件归并状态,TUI 和 Session 再订阅状态变化。主循环不直接操作界面,界面也不用理解 provider 的原始数据格式。
还有一个细节:listener 是逐个 await 的。这保证事件顺序稳定,但也意味着某个很慢的 listener 会对整个流形成背压。调试“模型已经返回但界面更新很慢”时,除了查网络层,也要检查事件订阅者。
十三、适合实际调试的一组断点
我最后把断点分成四组,按问题选择,而不是每次全部打开。
模型没有出现或认证失败
1 | packages/coding-agent/src/core/sdk.ts |
重点观察 agentDir、modelsPath、选中的 model.provider/model.id 和认证来源。不要在调试控制台或截图中暴露真实 key。
请求参数、代理地址或 header 不符合预期
1 | packages/coding-agent/src/core/sdk.ts |
Agent 为什么继续或停止
1 | packages/agent/src/agent-loop.ts |
每次停在循环顶部时观察:1
2
3
4
5
6hasMoreToolCalls
pendingMessages.length
lastCompletedTurn
message.stopReason
toolCalls.length
followUpMessages.length
工具为什么没有执行
1 | packages/agent/src/agent-loop.ts |
重点排查工具不存在、schema 校验失败、beforeToolCall 阻止、AbortSignal 取消和 stopReason === "length" 这几种分支。
十四、总结
Pi Agent 的核心并不是一个很玄学的黑盒。把 UI、会话和 provider 适配暂时拿掉以后,它的主干可以概括为:1
2
3
4
5
6
7准备上下文
→ 请求模型并消费流事件
→ 得到完整 assistant message
→ 如果有工具调用,执行并追加 toolResult
→ 带着新上下文再次请求模型
→ 当前任务链结束后检查 follow-up
→ 发出 agent_end
但真正值得学习的地方恰恰是这些“简单主干”周围没有被省略的工程细节:
- 用
StreamFn做依赖注入,让 Agent Core 与 provider 解耦; - 用 partial message 原位替换保证流式状态不污染历史;
- 用内外两层循环区分工具链、steering 和 follow-up;
- 工具调用先校验、可拦截、可并行,也能把失败反馈给模型;
- token 截断时拒绝执行可能不完整的工具参数;
- API key 每轮动态解析,适配会过期的认证;
- 通过统一事件流连接 Agent 状态、Session 和 TUI。
完成这条调用链以后,再继续阅读上下文压缩、Session JSONL、扩展系统和具体 Provider,就不会只是看到一堆互相跳转的 TypeScript 文件,而是能判断每一层在整个 Agent 生命周期中的位置。