惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
博客园 - Franky
T
Tailwind CSS Blog
Microsoft Azure Blog
Microsoft Azure Blog
The Cloudflare Blog
博客园 - 叶小钗
N
Netflix TechBlog - Medium
罗磊的独立博客
量子位
MyScale Blog
MyScale Blog
A
About on SuperTechFans
Blog — PlanetScale
Blog — PlanetScale
V
Visual Studio Blog
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
GbyAI
GbyAI
B
Blog
腾讯CDC
爱范儿
爱范儿
Recent Announcements
Recent Announcements
有赞技术团队
有赞技术团队
F
Fortinet All Blogs
雷峰网
雷峰网
G
Google Developers Blog
Google DeepMind News
Google DeepMind News

Nicksxs's Blog

6GB显存能跑35B MoE吗:FreeToken极限配置、性能压榨与实操教程 聊聊 Java 的类加载机制一 从ASM到ByteKit再到Arthas:一条Java字节码增强链路 8GB显存如何跑35B:深入理解FreeToken的专家缓存与CPU-GPU协同 让chatgpt来给我讲解下deepseek harness的插件体系 来学习下pi agent的原理 如何在命令行中输出进度条 联想 G400 安装 Windows 7 折腾记录 记一个gitea推送失败的问题 学一下chrome的扩展开发 看下chrome的内置模型 从 app.test 到小锁:valet 本地 HTTPS 的完整链路 浅析一下jpeg图片格式及其来源 关于github拉取下载加速的另一个方式 关于适合什么模型,推荐下llmfit 看看目前本地能跑什么模型,使用llama.cpp 最近使用vibe coding的一些感悟 使用php的inotify扩展来监听文件变更 一些设计模式的记忆点 使用xiaomi mimo大模型api运行Hermes Agent 结合Obsidian的cli的一体化体验 开始尝试使用obsidian作为笔记软件 学习下大神的知识库 体验下微软开源的Markdown转换工具Markitdown 学习下git的worktree 一些架构师知识点的记录 解答一下关于traefik的一点疑惑 记录一下迁移服务器需要使用的一些命令 较早代iPhone更换新iPhone的一些小指南 如何查看mac的路由表和网关等信息
Pi Agent 源码解析(一):从调试环境到双层主循环
Nicksxs · 2026-09-13 · via Nicksxs's Blog

最近一段时间我一直在使用 Pi。它的交互界面看起来很简单,但代码里把模型适配、会话管理、上下文处理、工具执行和终端 UI 分成了比较清楚的几层,很适合拿来理解一个 Coding Agent 到底是怎么运转的。

前面已经写过一篇 Pi Agent 工具系统:Read、Write、Edit 与 Bash 的实现原理,那篇关注 Agent 的“手脚”,本文继续向下看驱动这些工具反复工作的主循环。

这篇先从源码调试环境开始,然后沿着一次真实请求的调用链,重点看下面几个问题:

  1. Pi 的源码怎样安装、补齐模型数据并启动调试;
  2. 模型的 baseUrlapiKey 从哪里读取;
  3. Agent 在哪里初始化,streamFn 又是在哪里注入的;
  4. agent-loop.ts 中的两层循环分别解决什么问题;
  5. LLM 的流式事件如何变成 Agent 事件;
  6. 模型发起工具调用以后,工具结果如何回到上下文并触发下一轮请求。

本文对应我本地阅读的源码版本:

1
2
3
4
repository: https://github.com/earendil-works/pi
commit: 71dca871bc80b6bc97be37f0ca3189399d651fff
describe: v0.85.1-67-g71dca871b
date: 2026-09-11

Pi 的代码还在快速变化,后续如果发现文件名或行号对不上,应当先确认源码版本。本文代码片段均摘自上述提交;为了方便阅读,我只在少数位置加入了以 // 解读: 开头的注释,没有把真实逻辑改写成伪代码。

一、先看项目分层

仓库根目录的开发文档给出的主要结构如下:

1
2
3
4
5
packages/
ai/ # LLM provider 抽象、模型目录以及各家 API 适配
agent/ # Agent 状态、事件、主循环和工具执行
tui/ # 终端 UI
coding-agent/ # CLI、会话、配置、扩展、内置工具和交互模式

阅读时最容易混淆的是 agentcoding-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
2
node --version
npm --version

我使用的是 Node 22。如果是较早的 22.x,也要确认小版本不低于 22.19.0

2. Clone 与安装依赖

1
2
3
git clone https://github.com/earendil-works/pi.git
cd pi
npm ci

这里补充一下 npm installnpm 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
30
async 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
21
const 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");


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
11
try {
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;
  • 模型的 idproviderapi 是否匹配;
  • baseUrlreasoning、输入模态、上下文窗口、最大输出和价格字段是否合法。

看到下面的输出才说明这一步完成:

1
Generated model data is valid.

三、配置模型和密钥

Pi 默认读取的全局目录是:

1
~/.pi/agent/

经常使用 Pi 的机器上一般已经有 auth.jsonmodels.json 和相关设置,所以直接从源码启动时也能读取已有配置。sdk.ts 中创建模型运行时的代码如下:

1
2
3
4
const 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
5
const 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
19
packages/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
17
packages/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
#!/usr/bin/env node
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
7
if (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
17
import {
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";




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
17
import 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
114
agent = new Agent({


initialState: {
systemPrompt: "",
model,
thinkingLevel,
tools: [],
},



convertToLlm: convertToLlmWithBlockImages,



streamFn: async (model, context, options) => {

const providerRetrySettings = settingsManager.getProviderRetrySettings();
const httpIdleTimeoutMs = settingsManager.getHttpIdleTimeoutMs();



const effectiveTimeoutMs =
httpIdleTimeoutMs === 0 ? 2147483647 : httpIdleTimeoutMs;


const timeoutMs =
options?.timeoutMs ??
providerRetrySettings.timeoutMs ??
effectiveTimeoutMs;


const websocketConnectTimeoutMs =
options?.websocketConnectTimeoutMs ??
settingsManager.getWebSocketConnectTimeoutMs();



const headerRunner = extensionRunnerRef.current;


return modelRuntime.streamSimple(model, context, {

...options,
timeoutMs,
websocketConnectTimeoutMs,


maxRetries:
options?.maxRetries ?? providerRetrySettings.maxRetries,
maxRetryDelayMs:
options?.maxRetryDelayMs ??
providerRetrySettings.maxRetryDelayMs,



transformHeaders: async (requestHeaders) => {
const headers = mergeProviderAttributionHeaders(
model,
settingsManager,
options?.sessionId,
requestHeaders,
);
return headerRunner?.hasHandlers("before_provider_headers")
? headerRunner.emitBeforeProviderHeaders(headers ?? {})
: (headers ?? {});
},
});
},



onPayload: async (payload, _model) => {
const runner = extensionRunnerRef.current;
if (!runner?.hasHandlers("before_provider_request")) {
return payload;
}
return runner.emitBeforeProviderRequest(payload);
},


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: sessionManager.getSessionId(),



transformContext: async (messages) => {
const runner = extensionRunnerRef.current;
if (!runner) return messages;
return runner.emitContext(messages);
},


steeringMode: settingsManager.getSteeringMode(),
followUpMode: settingsManager.getFollowUpMode(),


transport: settingsManager.getTransport(),
thinkingBudgets: settingsManager.getThinkingBudgets(),
maxRetryDelayMs: settingsManager.getProviderRetrySettings().maxRetryDelayMs,
});

这个包装函数不只是把请求转发出去,它还在真正请求前补上:

  • HTTP 空闲超时和 WebSocket 连接超时;
  • 最大重试次数、最大重试等待时间;
  • provider attribution headers;
  • 扩展提供的 before_provider_headers hook;
  • 请求 payload 与响应 headers 的扩展 hook。

因此 streamFn 是 Agent Core 与真实模型运行时之间的注入点,也是调试“最终请求为什么长这样”的最佳断点之一。

3. StreamFn 的契约

它不是一个普通的 Promise<AssistantMessage>,而是返回事件流:

1
2
3
4
5
export 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
14
streamSimple(
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 的 streamSimpleagent-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
16
private 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
42
async prompt(message: AgentMessage | AgentMessage[]): Promise<void>;
async prompt(input: string, images?: ImageContent[]): Promise<void>;
async prompt(
input: string | AgentMessage | AgentMessage[],
images?: ImageContent[],
): Promise<void> {

if (this.activeRun) {
throw new Error(
"Agent is already processing a prompt. Use steer() or followUp() " +
"to queue messages, or wait for completion.",
);
}


const messages = this.normalizePromptInput(input, images);
await this.runPromptMessages(messages);
}

private async runPromptMessages(
messages: AgentMessage[],
options: { skipInitialSteeringPoll?: boolean } = {},
): Promise<void> {


await this.runWithLifecycle(async (signal) => {
await runAgentLoop(

messages,

this.createContextSnapshot(),

this.createLoopConfig(options),

(event) => this.processEvents(event),

signal,

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
38
export async function runAgentLoop(
prompts: AgentMessage[],
context: AgentContext,
config: AgentLoopConfig,
emit: AgentEventSink,
signal: AbortSignal | undefined,
streamFn: StreamFn,
): Promise<AgentMessage[]> {

const newMessages: AgentMessage[] = [...prompts];



const currentContext: AgentContext = {
...context,
messages: [...context.messages, ...prompts],
};


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 });
}


await runLoop(
currentContext,
newMessages,
config,
signal,
emit,
streamFn ?? getDefaultStreamFn(),
);

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
178
async function runLoop(
initialContext: AgentContext,
newMessages: AgentMessage[],
initialConfig: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
streamFunction: StreamFn,
): Promise<void> {

let currentContext = initialContext;
let config = initialConfig;

let lastCompletedTurn: PrepareNextTurnContext | undefined;



let pendingMessages: AgentMessage[] =
(await config.getSteeringMessages?.()) || [];


while (true) {


let hasMoreToolCalls = true;


while (hasMoreToolCalls || pendingMessages.length > 0) {


if (lastCompletedTurn) {
const nextTurnSnapshot =
await config.prepareNextTurn?.(lastCompletedTurn);
if (nextTurnSnapshot) {

currentContext =
nextTurnSnapshot.context ?? currentContext;
config = {
...config,
model: nextTurnSnapshot.model ?? config.model,
reasoning:
nextTurnSnapshot.thinkingLevel === undefined
? config.reasoning
: nextTurnSnapshot.thinkingLevel === "off"
? undefined
: nextTurnSnapshot.thinkingLevel,
};
}




if (pendingMessages.length === 0) {
pendingMessages =
(await config.getSteeringMessages?.()) || [];
}

await emit({ type: "turn_start" });
}


if (pendingMessages.length > 0) {
for (const message of pendingMessages) {

await emit({ type: "message_start", message });
await emit({ type: "message_end", message });
currentContext.messages.push(message);
newMessages.push(message);
}

pendingMessages = [];
}



const message = await streamAssistantResponse(
currentContext,
config,
signal,
emit,
streamFunction,
);

newMessages.push(message);



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",
);

const toolResults: ToolResultMessage[] = [];


hasMoreToolCalls = false;

if (toolCalls.length > 0) {


const executedToolBatch =
message.stopReason === "length"
? await failToolCallsFromTruncatedMessage(
toolCalls,
emit,
)
: await executeToolCalls(
currentContext,
message,
config,
signal,
emit,
);

toolResults.push(...executedToolBatch.messages);


hasMoreToolCalls = !executedToolBatch.terminate;

for (const result of toolResults) {


currentContext.messages.push(result);
newMessages.push(result);
}
}



await emit({ type: "turn_end", message, toolResults });


lastCompletedTurn = {
message,
toolResults,
context: currentContext,
newMessages,
};



if (await config.shouldStopAfterTurn?.(lastCompletedTurn)) {
await emit({ type: "agent_end", messages: newMessages });
return;
}



pendingMessages =
(await config.getSteeringMessages?.()) || [];
}


const followUpMessages =
(await config.getFollowUpMessages?.()) || [];
if (followUpMessages.length > 0) {


pendingMessages = followUpMessages;


continue;
}


break;
}


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()

followUpsteering 的差别是:

  • 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
137
async function streamAssistantResponse(
context: AgentContext,
config: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
streamFunction: StreamFn,
): Promise<AssistantMessage> {



let messages = context.messages;
if (config.transformContext) {
messages = await config.transformContext(messages, signal);
}



const llmMessages = await config.convertToLlm(messages);



const llmContext: Context = {
systemPrompt: context.systemPrompt,
messages: llmMessages,
tools: context.tools,
};



const resolvedApiKey =
(config.getApiKey
? await config.getApiKey(config.model.provider)
: undefined) || config.apiKey;



const response = await streamFunction(config.model, llmContext, {


...config,
apiKey: resolvedApiKey,
signal,
});



let partialMessage: AssistantMessage | null = null;


let addedPartial = false;



for await (const event of response) {
switch (event.type) {
case "start":

partialMessage = event.partial;

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":


if (partialMessage) {
partialMessage = event.partial;


context.messages[context.messages.length - 1] =
partialMessage;

await emit({
type: "message_update",
assistantMessageEvent: event,
message: { ...partialMessage },
});
}
break;

case "done":
case "error": {


const finalMessage = await response.result();
if (addedPartial) {

context.messages[context.messages.length - 1] =
finalMessage;
} else {

context.messages.push(finalMessage);
}
if (!addedPartial) {


await emit({
type: "message_start",
message: { ...finalMessage },
});
}

await emit({ type: "message_end", message: finalMessage });
return finalMessage;
}
}
}



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;
}

这里可以看到三层事件:

  1. Provider 流事件:text_deltathinking_deltatoolcall_delta 等;
  2. Agent 事件:message_startmessage_updatemessage_end
  3. 更外层 UI 订阅 Agent 事件,刷新终端内容和会话状态。

收到 start 时,代码先把 partial message 放进 context.messages。后续每个 delta 都不是继续 push 新消息,而是替换数组最后一项。否则每个 token 都会变成一条历史消息。

收到 doneerror 后,再通过 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
42
async function executeToolCalls(
currentContext: AgentContext,
assistantMessage: AssistantMessage,
config: AgentLoopConfig,
signal: AbortSignal | undefined,
emit: AgentEventSink,
): Promise<ExecutedToolCallBatch> {

const toolCalls = assistantMessage.content.filter(
(c) => c.type === "toolCall",
);


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,
);
}



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
82
async 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) {

return {
kind: "immediate",
result: createErrorToolResult(
`Tool ${toolCall.name} not found`,
),
isError: true,
};
}

try {


const preparedToolCall = prepareToolCallArguments(tool, toolCall);
const validatedArgs = validateToolArguments(tool, preparedToolCall);

if (config.beforeToolCall) {


const beforeResult = await config.beforeToolCall(
{
assistantMessage,
toolCall,
args: validatedArgs,
context: currentContext,
},
signal,
);


if (signal?.aborted) {
return {
kind: "immediate",
result: createErrorToolResult("Operation aborted"),
isError: true,
};
}

if (beforeResult?.block) {

const result = createErrorToolResult(
beforeResult.reason || "Tool execution was blocked",
);
if (beforeResult.terminate === true) {

result.terminate = true;
}
return { kind: "immediate", result, isError: true };
}
}


return {
kind: "prepared",
toolCall,
tool,
args: validatedArgs,
};
} catch (error) {


return {
kind: "immediate",
result: createErrorToolResult(
error instanceof Error ? error.message : String(error),
),
isError: true,
};
}
}

它完成四项检查:

  1. 工具是否存在;
  2. 是否需要预处理参数;
  3. 参数是否符合工具 schema;
  4. beforeToolCall hook 是否阻止执行。

工具报错通常也不会炸掉整个 Agent Loop,而是转成 isError: truetoolResult。这样模型下一轮能看见错误,并决定修正参数、换一种方法或直接向用户解释失败。

工具执行结束后生成标准消息:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
function 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
8
assistant toolCall
→ 参数准备与校验
→ beforeToolCall
→ tool.execute
→ afterToolCall
→ ToolResultMessage
→ 加入上下文
→ 下一次 LLM 请求

主循环只负责 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
66
private async processEvents(event: AgentEvent): Promise<void> {


switch (event.type) {
case "message_start":


this._state.streamingMessage = event.message;
break;

case "message_update":

this._state.streamingMessage = event.message;
break;

case "message_end":


this._state.streamingMessage = undefined;
this._state.messages.push(event.message);
break;

case "tool_execution_start": {

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;
}

case "turn_end":

if (event.message.role === "assistant" && event.message.errorMessage) {
this._state.errorMessage = event.message.errorMessage;
}
break;

case "agent_end":


this._state.streamingMessage = undefined;
break;
}


const signal = this.activeRun?.abortController.signal;
if (!signal) {
throw new Error("Agent listener invoked outside active run");
}

for (const listener of this.listeners) {
await listener(event, signal);
}
}

这是一种很清楚的单向数据流:主循环产生事实事件,Agent 根据事件归并状态,TUI 和 Session 再订阅状态变化。主循环不直接操作界面,界面也不用理解 provider 的原始数据格式。

还有一个细节:listener 是逐个 await 的。这保证事件顺序稳定,但也意味着某个很慢的 listener 会对整个流形成背压。调试“模型已经返回但界面更新很慢”时,除了查网络层,也要检查事件订阅者。

十三、适合实际调试的一组断点

我最后把断点分成四组,按问题选择,而不是每次全部打开。

模型没有出现或认证失败

1
2
3
4
5
6
7
8
packages/coding-agent/src/core/sdk.ts
ModelRuntime.create(...)
findInitialModel(...)

packages/coding-agent/src/core/model-runtime.ts
create(...)
getAuth(...)
prepareRequest(...)

重点观察 agentDirmodelsPath、选中的 model.provider/model.id 和认证来源。不要在调试控制台或截图中暴露真实 key。

1
2
3
4
5
6
7
packages/coding-agent/src/core/sdk.ts
streamFn: async (...)
onPayload
transformHeaders

packages/coding-agent/src/core/model-runtime.ts
streamSimple(...)

Agent 为什么继续或停止

1
2
packages/agent/src/agent-loop.ts
runLoop(...)

每次停在循环顶部时观察:

1
2
3
4
5
6
hasMoreToolCalls
pendingMessages.length
lastCompletedTurn
message.stopReason
toolCalls.length
followUpMessages.length

工具为什么没有执行

1
2
3
4
5
packages/agent/src/agent-loop.ts
executeToolCalls(...)
prepareToolCall(...)
executePreparedToolCall(...)
finalizeExecutedToolCall(...)

重点排查工具不存在、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 生命周期中的位置。

参考源码