02 — LLM 抽象层 (pi-ai)

问题

OpenAI、Anthropic、Google 等 LLM 供应商的 API 各不相同:消息格式不同、流式事件不同、工具调用格式不同、认证方式不同。如果 Agent 代码直接调用某个 provider 的 SDK,切换 provider 需要大面积改码。

pi-ai 的任务是提供一个统一抽象层,让上层 Agent 代码只面对一套接口,无需关心底层是哪个 provider。

核心类型

Api 与 Provider

packages/ai/src/types.ts:17-29 定义了已知的 API 类型:

1
2
3
4
5
6
7
8
9
10
11
12
export type KnownApi =
| "openai-completions"
| "openai-responses"
| "anthropic-messages"
| "bedrock-converse-stream"
| "google-generative-ai"
| "google-vertex"
| "mistral-conversations"
| "pi-messages"
// ...

export type Api = KnownApi | (string & {}); // 允许自定义 API

每个 provider 实现一个 Api(如 anthropic-messages),对应 api/ 目录下的具体实现文件(packages/ai/src/api/anthropic-messages.ts)。

Provider 接口

每个 provider 需要实现 Provider 接口,提供 stream() 方法:

1
2
Provider<Model<Api>>
└─ stream(model, context, options) → AssistantMessageEventStream

AssistantMessageEventStreampackages/ai/src/utils/event-stream.ts)是统一的事件流,无论底层 provider 返回什么格式,最终都归一化为这套事件。

AssistantMessageEventStream 事件

流中的事件类型统一为:

  • start — assistant 消息开始
  • update — 流式内容增量(文本、thinking、toolCall)
  • end — 消息结束,携带 stopReasonstoplengthtoolUseerroraborteddeferred
  • usage — token 用量和成本
  • diagnostic — 诊断信息(错误、恢复等)

关键设计:无论 provider 返回 SSE(OpenAI)、message_delta(Anthropic)还是其他格式,api/ 层的适配器都将其转换成同一套事件。上层只需消费这套事件流。

模型管理

Models 类

Modelspackages/ai/src/models.ts)是模型管理的核心类。它管理:

  1. 已注册的 provider — 通过 registerProvider() 注册
  2. 模型目录 — 通过 ModelCatalogmodel-catalog.ts)维护可用模型列表
  3. 模型存储 — 通过 ModelsStoremodels-store.ts)持久化模型数据
  4. 认证 — 集成 CredentialStore 管理 API key 和 OAuth token

模型刷新

refresh() 方法(models.ts:391)负责从各 provider 动态拉取最新模型列表:

  • 并发地对每个 provider 执行刷新
  • 使用 generation counter 和 AbortController 来取消过期的刷新
  • OAuth token 过期时在 store lock 下自动刷新(resolveRefreshCredential, models.ts:453

成本计算

calculateCost()models.ts:891)支持阶梯定价:根据 inputTokensAbove 阈值匹配最高档位。Anthropic 的 1 小时缓存写入有 2x 输入费率的特殊处理。

认证体系

packages/ai/src/auth/ 目录实现了完整的认证抽象。

三种认证方式

类型 接口 适用场景
API Key ApiKeyAuth OpenAI、DeepSeek 等直接使用 API key
OAuth OAuthAuth GitHub Copilot、OpenAI Codex、Anthropic Pro/Max
环境感知 特殊处理 Google Vertex(ADC 文件)、Amazon Bedrock(AWS Profile/IAM)

ApiKeyAuth

auth/types.ts:170-199

1
2
3
4
5
6
interface ApiKeyAuth {
name: string;
login?(interaction: AuthInteraction): Promise<...>; // 交互式设置
check?(input: AuthInput): boolean; // 无副作用检查
resolve(input: AuthInput): Promise<AuthResolution>; // 解析凭证
}

resolve() 的返回值表示是否已配置。envApiKeyAuth()auth/helpers.ts:9)是标准实现:存储的 key 优先,否则读环境变量。

OAuthAuth

auth/types.ts:206-230

1
2
3
4
5
6
interface OAuthAuth {
name: string;
login(interaction: AuthInteraction): Promise<...>;
refresh(credential, signal): Promise<...>; // 网络调用,刷新 token
toAuth(credential): AuthResolution; // 无副作用,从凭证推导请求 auth
}

关键设计refreshtoAuth 分离。refresh 是有副作用的网络调用(刷新过期 token),toAuth 是纯函数(从已有凭证推导 auth header)。这让 Models 类可以在锁的保护下执行 refresh,避免并发刷新。

CredentialStore

auth/types.ts:65-94:每个 provider 对应一个凭证,modify() 是唯一的写入路径(串行化的 read-modify-write)。InMemoryCredentialStore 通过 promise chain 串行化写入。

OAuth 流程

auth/oauth/ 目录实现了多个 OAuth 流程:

  • PKCE 流程pkce.ts):Web Crypto 生成 code_verifier/challenge
  • 设备码流程device-code.ts
  • 本地回调服务器anthropic.ts):使用 node:http 启动回调服务器接收授权码

lazyOAuth()auth/helpers.ts:40)使用动态 import 延迟加载 OAuth 实现,避免在不需要 OAuth 时加载相关代码。

流式调用

stream() 的 lazy 模式

Models.stream()models.ts:672):

1
2
3
4
5
6
7
stream(model, context, options): AssistantMessageEventStream {
return lazyStream(model, async () => {
const provider = this.requireProvider(model);
const { requestModel, requestOptions } = await this.applyAuth(model, options);
return provider.stream(requestModel, context, requestOptions);
});
}

lazyStream 的关键:认证解析在流被消费时才执行,不是在 stream() 被调用时。如果认证失败,错误不会以 throw 形式出现,而是作为流中的 error 事件。

applyAuth()

models.ts:641:解析 provider 认证 → 合并 headers(模型 headers + auth headers + 调用者 headers + transformHeaders 回调)→ 覆盖 env → 返回 { requestModel, requestOptions }。如果 provider 未配置,抛出 ModelsError("auth", ...)

为什么 StreamFn 契约要求”不能 throw”

packages/agent/src/types.ts:28 明确规定:

1
2
3
4
5
6
7
8
9
10
11
12
/**
* Contract:
* - Must not throw or return a rejected promise for request/model/runtime failures.
* - Must return an AssistantMessageEventStream.
* - Failures must be encoded in the returned stream via protocol events and a
* final AssistantMessage with stopReason "error" or "aborted" and errorMessage.
*/
export type StreamFn = (
model: Model<Api>,
context: Context,
options?: SimpleStreamOptions,
) => AssistantMessageEventStream | Promise<AssistantMessageEventStream>;

原因:Agent loop 是一个持续运行的状态机。如果 streamFn 抛异常,loop 会被意外的异常打断,正在进行的工具调用和持久化状态可能不一致。将错误放入流中,loop 可以用与处理”正常结束”相同的代码路径来处理”错误结束”,保证状态机一致性。

工具调用统一

不同 provider 的 tool calling 格式差异很大:

  • OpenAI:tool_calls 数组,function.arguments 是 JSON 字符串
  • Anthropic:content 数组中的 tool_use block,input 是对象
  • Google:functionCall 字段

pi-aiapi/ 适配器层将这些差异归一化为统一的 AssistantMessage 结构:

1
2
3
4
5
6
7
8
9
interface AssistantMessage {
content: Array<
| { type: "text"; text: string }
| { type: "thinking"; text: string }
| { type: "toolCall"; toolCall: { id: string; name: string; arguments: object } }
>;
stopReason: "stop" | "length" | "toolUse" | "error" | "aborted" | "deferred";
usage?: Usage;
}

上层 Agent loop 只需处理这套统一结构,不需要知道是哪个 provider 返回的。

重试与容错

Provider 重试

packages/ai/src/utils/provider-retry.ts

retryProviderRequest()(line 105)复现了 OpenAI/Anthropic SDK 的重试策略:

  • isRetryableProviderError()(line 23):检查 x-should-retry header、408/409/429/5xx
  • getRetryDelayMs()(line 51):优先使用 retry-after-ms / retry-after header,否则指数退避,上限 60s
  • 中断式退避:等待期间可以被打断(abort signal)

上下文溢出检测

packages/ai/src/utils/overflow.ts

isContextOverflow()(line 134)通过正则匹配检测 ~20 个 provider 的上下文溢出错误信息。还处理两种静默溢出:

  • z.ai 静默溢出usage.input > contextWindow
  • Xiaomi MiMo length 溢出stopReason "length" + 零输出 + 上下文已满

isRecoverableLength()(line 171)判断是否可以通过 compaction 后重试来恢复。

流式 JSON 修复

packages/ai/src/utils/json-parse.tsparseStreamingJson / parseJsonWithRepair 用于解析流式传输中不完整的 tool call 参数 JSON。LLM 流式输出时,arguments JSON 可能是截断的,需要修复后才能使用。

Telemetry 集成

TelemetryContext(来自 pi-telemetry)通过 ProviderRequestOptions.telemetryContexttypes.ts:127)传入每次 provider 请求。Provider 实现可以调用 startSpan() 来追踪请求耗时和属性。

与 Agent 核心的连接

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
Agent Loop (packages/agent/src/agent-loop.ts)

│ 调用 streamFn(model, context, options)


Models.stream() (packages/ai/src/models.ts:672)

├── applyAuth() → 解析认证
├── Provider.stream() → 调用具体 provider API


AssistantMessageEventStream (统一事件流)

├── start → assistant 消息开始
├── update → 流式增量(text/thinking/toolCall)
├── end → stopReason + usage
└── diagnostic → 错误/恢复信息


Agent Loop 消费事件,执行工具调用,进入下一轮

coding-agentsdk.ts 调用 setDefaultStreamFn(streamSimple)(line 37)将 pi-ai 的流式函数注入 Agent 核心。streamSimpleModels.stream() 的简化封装,满足 StreamFn 类型签名。

学习要点

  • 统一抽象的核心是事件流:不同 provider 的 API 差异在 api/ 适配器层被吸收,上层只消费 AssistantMessageEventStream
  • 错误放流中而非抛异常:这是 Agent 系统的关键设计——保证状态机不被意外异常打断
  • 认证的 refresh/toAuth 分离:将有副作用的 token 刷新和无副作用的 auth 推导分离,支持锁保护的并发安全刷新
  • 溢出检测是 provider-specific 的:没有统一的错误码标准,只能通过正则匹配各 provider 的错误信息来检测
  • lazyStream 实现延迟认证:认证在流被消费时才执行,避免 stream() 调用时的阻塞

上一篇:01-整体架构与包依赖
下一篇:03-Agent 运行时核心-pi-agent-core