ARTICLE · 1157712
Pi源码阅读「2」
前言
我们都说pi是一个极简的agent内核,只是因为原生tools少,设计loop简单吗?并不是这样
Pi插件系统
初始化
看看是什么加载插件的吧。在loader.ts中,我们可以看见插件的加载方法
/** * Discover extensions in a directory. * * Discovery rules: * 1. Direct files: `extensions/*.ts` or `*.js` → load * 2. Subdirectory with index: `extensions/* /index.ts` or `index.js` → load * 3. Subdirectory with package.json: `extensions/* /package.json` with "pi" field → load what it declares * * No recursion beyond one level. Complex packages must use package.json manifest. */functiondiscoverExtensionsInDir(dir: string): string[] {if (!fs.existsSync(dir)) {return []; }constdiscovered: string[] = [];try {const entries = fs.readdirSync(dir, { withFileTypes: true });for (const entry of entries) {const entryPath = path.join(dir, entry.name);// 1. Direct files: *.ts or *.jsif ((entry.isFile() || entry.isSymbolicLink()) && isExtensionFile(entry.name)) { discovered.push(entryPath);continue; }// 2 & 3. Subdirectoriesif (entry.isDirectory() || entry.isSymbolicLink()) {const entries = resolveExtensionEntries(entryPath);if (entries) { discovered.push(...entries); } } } } catch {return []; }return discovered;}resolveExtensionEntries跳转
/** * Resolve extension entry points from a directory. * * Checks for: * 1. package.json with "pi.extensions" field -> returns declared paths * 2. index.ts or index.js -> returns the index file * * Returns resolved paths or null if no entry points found. */functionresolveExtensionEntries(dir: string): string[] | null {// Check for package.json with "pi" field firstconst packageJsonPath = path.join(dir, "package.json");if (fs.existsSync(packageJsonPath)) {const manifest = readPiManifest(packageJsonPath);if (manifest?.extensions?.length) {constentries: string[] = [];for (const extPath of manifest.extensions) {const resolvedExtPath = path.resolve(dir, extPath);if (fs.existsSync(resolvedExtPath)) { entries.push(resolvedExtPath); } }if (entries.length > 0) {return entries; } } }// Check for index.ts or index.jsconst indexTs = path.join(dir, "index.ts");const indexJs = path.join(dir, "index.js");if (fs.existsSync(indexTs)) {return [indexTs]; }if (fs.existsSync(indexJs)) {return [indexJs]; }returnnull;}其实No recursion beyond one level. Complex packages must use package.json manifest.已经写明白了设计边界,同时标明了加载优先级项目级先于全局级,且按path.resolve去重。
导入是在
asyncfunctionloadExtensionModule(extensionPath: string, cacheToken?: ExtensionCacheToken) {if (isCurrentCacheToken(cacheToken)) {const cachedFactory = extensionCache.get(extensionPath);if (cachedFactory) {return cachedFactory; } }const createJitiImpl = awaitgetCreateJiti();// Compiled binaries and the bundled Node distribution use embedded modules.// Source TypeScript reuses host modules and root tsconfig paths. Unbundled// Node builds use dist aliases and do not need the bundled virtual modules.const resolutionOptions = usesEmbeddedModules ? { virtualModules: awaitgetVirtualModules(), tryNative: false } : isTypeScriptSourceRuntime ? { virtualModules: awaitgetVirtualModules(), tsconfigPaths: true } : { alias: getAliases() };const jiti = createJitiImpl(import.meta.url, {moduleCache: false, ...resolutionOptions, });constmodule = await jiti.import(extensionPath, { default: true });const factory = moduleasExtensionFactory;if (typeof factory !== "function") {returnundefined; }if (isCurrentCacheToken(cacheToken)) { extensionCache.set(extensionPath, factory); }return factory;}看这里插件状态有三种
const isNodeSeaBinary = ("sea"in process.features && process.features.sea === true) || process.getBuiltinModule("node:sea")?.isSea() === true;const isTypeScriptSourceRuntime = !isBunBinary && path.extname(fileURLToPath(import.meta.url)) === ".ts";const usesEmbeddedModules = isBunBinary || isNodeSeaBinary || isBundledNode;SEA是nodejs的单文件可执行程序,类似pkg?第二个就是ts源码直接跑,第三个是二进制文件,三种情况全包含
看这里的兼容层核心是在
functiongetAliases(): Record<string, string> {if (_aliases) return _aliases;const __dirname = path.dirname(fileURLToPath(import.meta.url));const packageIndex = path.resolve(__dirname, "../..", "index.js");const typeboxEntry = require.resolve("typebox");const typeboxCompileEntry = require.resolve("typebox/compile");const typeboxValueEntry = require.resolve("typebox/value");const packagesRoot = path.resolve(__dirname, "../../../../");const resolveWorkspaceOrImport = (workspaceRelativePath: string, specifier: string): string => {const workspacePath = path.join(packagesRoot, workspaceRelativePath);if (fs.existsSync(workspacePath)) {return workspacePath; }returnfileURLToPath(import.meta.resolve(specifier)); };const piCodingAgentEntry = packageIndex;const piAgentCoreEntry = resolveWorkspaceOrImport("agent/dist/index.js", "@earendil-works/pi-agent-core");const piTuiEntry = resolveWorkspaceOrImport("tui/dist/index.js", "@earendil-works/pi-tui");// Extensions resolve the pi-ai root to the compat entrypoint (a strict// superset of the core entrypoint): existing extensions using the old// global API keep working at runtime until compat is removed.const piAiCompatEntry = resolveWorkspaceOrImport("ai/dist/compat.js", "@earendil-works/pi-ai/compat");const piAiOauthEntry = resolveWorkspaceOrImport("ai/dist/oauth.js", "@earendil-works/pi-ai/oauth");const piAiProvidersEntry = resolveWorkspaceOrImport("ai/dist/providers/all.js","@earendil-works/pi-ai/providers/all", ); _aliases = {"@earendil-works/pi-coding-agent": piCodingAgentEntry,"@earendil-works/pi-agent-core": piAgentCoreEntry,"@earendil-works/pi-tui": piTuiEntry,"@earendil-works/pi-ai/providers/all": piAiProvidersEntry,"@earendil-works/pi-ai/compat": piAiCompatEntry,"@earendil-works/pi-ai/oauth": piAiOauthEntry,"@earendil-works/pi-ai": piAiCompatEntry,"@mariozechner/pi-coding-agent": piCodingAgentEntry,"@mariozechner/pi-agent-core": piAgentCoreEntry,"@mariozechner/pi-tui": piTuiEntry,"@mariozechner/pi-ai/providers/all": piAiProvidersEntry,"@mariozechner/pi-ai/compat": piAiCompatEntry,"@mariozechner/pi-ai/oauth": piAiOauthEntry,"@mariozechner/pi-ai": piAiCompatEntry,typebox: typeboxEntry,"typebox/compile": typeboxCompileEntry,"typebox/value": typeboxValueEntry,"@sinclair/typebox": typeboxEntry,"@sinclair/typebox/compile": typeboxCompileEntry,"@sinclair/typebox/value": typeboxValueEntry, };return _aliases;}做了模块配置和兼容处理
createExtensionRuntime函数是初始化的地方,主要是抛出异常
createExtensionAPI函数是api层面了
on(event: string, handler: HandlerFn): () =>void {assertActive();constregisteredHandler: HandlerFn = (...args) =>handler(...args);const list = extension.handlers.get(event) ?? []; list.push(registeredHandler); extension.handlers.set(event, list);return() => {const handlers = extension.handlers.get(event);if (!handlers) return;const handlerIndex = handlers.indexOf(registeredHandler);if (handlerIndex === -1) return; handlers.splice(handlerIndex, 1);if (handlers.length === 0) extension.handlers.delete(event); }; },registerTool(tool: ToolDefinition): void {assertActive();if (typeof tool.parameters !== "object" || tool.parameters === null || Array.isArray(tool.parameters)) {thrownewError(`Tool "${tool.name}" registered by extension "${extension.path}" must define an object parameter schema.`, ); } extension.tools.set(tool.name, {definition: tool,sourceInfo: extension.sourceInfo, }); runtime.refreshTools(); },这里是注册类方法直接写扩展对象
asyncfunctioninitializeExtension(factory: ExtensionFactory,extensionPath: string,resolvedPath: string,cwd: string,eventBus: EventBus,runtime: ExtensionRuntime,): Promise<Extension> {const extension = createExtension(extensionPath, resolvedPath);const load = createExtensionAPI(extension, runtime, cwd, eventBus);try {awaitfactory(load.api); load.commit(); } catch (error) { load.discard();throw error; }time(`${extensionPath} factory`, "extensions");return extension;}把全部的串联起来,异常会被全部拿到,一个扩展加载失败不影响其他扩展,错误进LoadExtensionsResult.errors
跑起来
真正的跑起来是在runner.ts
for (const { name, config, extensionPath } ofthis.runtime.pendingProviderRegistrations) {try {if (providerActions?.registerProvider) { providerActions.registerProvider(name, config); } else {this.modelRegistry.registerProvider(name, config); } } catch (err) {this.emitError({ extensionPath,event: "register_provider",error: err instanceofError ? err.message : String(err),stack: err instanceofError ? err.stack : undefined, }); } }this.runtime.pendingProviderRegistrations = [];for (const { provider, extensionPath } ofthis.runtime.pendingNativeProviderRegistrations) {try {if (providerActions?.registerNativeProvider) { providerActions.registerNativeProvider(provider); } else {this.modelRegistry.registerProvider(provider); } } catch (err) {this.emitError({ extensionPath,event: "register_provider",error: err instanceofError ? err.message : String(err),stack: err instanceofError ? err.stack : undefined, }); } }this.runtime.pendingNativeProviderRegistrations = [];每个provider单独try/catch,一个扩展注册的provider配置非法,不会挡住后面的
冲完队列之后马上把方法体换掉了
// From this point on, provider registration/unregistration takes effect immediately// without requiring a /reload.this.runtime.registerProvider = (name, config) => {if (providerActions?.registerProvider) { providerActions.registerProvider(name, config);return; }this.modelRegistry.registerProvider(name, config);};所以整个runtime是三个阶段的:加载期全是抛异常的桩(createExtensionRuntime),bindCore的时候冲队列并且换成真实现,之后就立即生效了
两阶段提交
前面说注册类方法直接写扩展对象,那registerProvider这种需要全局状态的怎么办?排队
constapplyRuntimeChange = (change: () => void) => {if (state === "loading") pendingRuntimeChanges.push(change);elsechange();};state是个三态"loading" | "active" | "failed",registerProvider/registerMcpServer/registerVirtualModel还有flag的默认值全都走这个
为什么要这么搞?因为扩展工厂是允许async的(文档里说了可以fetch配置再注册provider)。如果允许工厂执行到一半直接改全局状态,失败的时候就得回滚,而回滚一个已经生效的provider注册很麻烦。排队+commit就让"加载失败"永远是干净的
commit: () => {if (state !== "loading") return; runtime.assertActive();for (const [name, value] of pendingFlagValues) {if (!runtime.flagValues.has(name)) runtime.flagValues.set(name, value); }for (const apply of pendingRuntimeChanges) apply(); state = "active";clearPending();},discard: () => {if (state !== "loading") return; state = "failed";for (const unsubscribe of loadingUnsubscribers) unsubscribe();clearPending();},failed状态的assertActive还会换一条错误消息,区分"加载失败"和"ctx过期"
constassertActive = () => {if (state === "failed") {thrownewError(`Extension "${extension.path}" failed to load and its API is no longer active.`); } runtime.assertActive();};Context是怎么造的
扩展拿到的ctx全是用getter写的,runner.ts的createContext
createContext(): ExtensionContext {const runner = this;const getModel = this.getModel;return {getui() { runner.assertActive(); return runner.uiContext; },getmode() { runner.assertActive(); return runner.mode; },getcwd() { runner.assertActive(); return runner.cwd; },getsignal() { runner.assertActive(); returngetSignalFn(); },isIdle: () => { runner.assertActive(); return runner.isIdleFn(); },abort: () => { runner.assertActive(); runner.abortFn(); },getContextUsage: () => { runner.assertActive(); return runner.getContextUsageFn(); }, ... };}每一个属性访问都先assertActive(),而且大部分是通过runner.xxxFn()运行时读的,不是捕获时的快照。这就是为什么扩展热重载之后,老ctx再碰一下就报错
然后命令专用的那个context,注释直接说明了不能用展开
createCommandContext(): ExtensionCommandContext {// Use property descriptors instead of object spread so the guarded getters from// createContext() stay lazy. A spread would eagerly read them once and freeze the// old values into the returned object, bypassing stale-instance checks.const context = Object.defineProperties( {},Object.getOwnPropertyDescriptors(this.createContext()), ) asExtensionCommandContext;这个注释是全文件最值钱的一段。要是写成{ ...this.createContext() },getter会被求值一遍,把当时的uiContext/signal/cwd冻进新对象,之后reload完全检测不到。getOwnPropertyDescriptors复制的是描述符本身,getter保持惰性
工具context同理,用defineProperties挂tools(getter,每次现取)和executeTool(value,函数不变)
createToolContext(toolCallId: string, signal: AbortSignal | undefined): ExtensionToolContext {const runner = this;// createContext() returns a fresh object, so adding properties does not affect other contexts.returnObject.defineProperties(this.createContext() asExtensionToolContext, {tools: { get() { runner.assertActive(); return runner.getCallableToolsFn(); } },executeTool: {value: async (name: string, args: unknown, options: ExecuteToolOptions = {}) => { runner.assertActive(); ... }, }, });}顺便,无头模式(比如print)也不是直接给个空对象,而是一整套noOp实现
constnoOpUIContext: ExtensionUIContext = {select: async () => undefined,confirm: async () => false,input: async () => undefined,notify: () => {},custom: async () => undefinedasnever,setTheme: (_theme: string | Theme) => ({ success: false, error: "UI not available" }), ...};注意confirm返回的是false不是抛错。这样写给TUI的扩展在print模式下也不会崩,而且配合扩展里if (!(await ctx.ui.confirm(...))) return { cancel: true }的写法,无头下自动就是"取消",保守默认。hasUI()就是身份比较this.uiContext !== noOpUIContext
事件派发
ExtensionEvent是个41个成员的联合类型(types.ts:1375),每个事件有自己的合流规则。派发前的第一件事是快照
functionsnapshotEventHandlers(extensions: Extension[], event: ExtensionEvent["type"]) {return extensions.map((ext) => ({ ext, handlers: ext.handlers.get(event)?.slice() ?? [] }));}.slice()是必须的。handler是await的,执行期间扩展可能又调pi.on()或者调用之前拿到的注销函数。快照让本次派发的参与者集合在开始时就固定住
合流规则我整理了一下,大概是这样
有12个emitXxx方法,循环骨架全是重复的。代价换来的是每个事件的语义能被文档精确表述,扩展作者不用猜
一个有意思的细节:emitToolCall不try/catch
对比一下emitToolResult里密密麻麻的try/catch,emitToolCall是一个都没有
asyncemitToolCall(event: ToolCallEvent): Promise<ToolCallEventResult | undefined> {const ctx = this.createContext();letresult: ToolCallEventResult | undefined;for (const { handlers } ofsnapshotEventHandlers(this.extensions, "tool_call")) {for (const handler of handlers) {const handlerResult = awaithandler(event, ctx); // 没有 try/catchif (handlerResult) { result = handlerResult asToolCallEventResult;if (result.block) {return result; } } } }return result;}因为下游有兜底。异常从agent-session.ts的_beforeToolCall一路rethrow,最后被agent-loop.ts的prepareToolCall接住
} catch (error) {return {kind: "immediate",result: createErrorToolResult(error instanceofError ? error.message : String(error)),isError: true, };}文档说的"A tool_call handler failure blocks the tool as a fail-safe"就是这么实现的——靠调用链,不靠派发器。权限类扩展挂了等于所有工具调用失败,这是保守但正确的方向
emitToolResult的字段耦合
if (handlerResult.content !== undefined) { currentEvent.content = handlerResult.content;// Structured content that is not replaced along with the content may no longer match it.if (handlerResult.structuredContent === undefined) delete currentEvent.structuredContent; modified = true;}替换content但不给structuredContent,就把它删掉,因为结构化输出可能已经和新文本不符了。这条规则在agent-loop.ts的finalizeExecutedToolCall里有对称的实现
还有modified标志,全部handler都没改就返回undefined,调用方agent-session.ts就能跳过图片归一化那一段
emitContext和prompt cache
这个是整个系统里缓存敏感度最高的一段
asyncemitContext(messages: AgentMessage[]): Promise<AgentMessage[]> {const ctx = this.createContext();let currentMessages = structuredClone(messages);for (const { ext, handlers } ofsnapshotEventHandlers(this.extensions, "context")) {for (const handler of handlers) {try {const visibleMessages = currentMessages.filter((message) => message.role !== "system");const visibleSnapshot = visibleMessages.slice();constevent: ContextEvent = { type: "context", messages: visibleMessages };const handlerResult = (awaithandler(event, ctx)) asContextEventResult | undefined;// Handlers may return a new list or edit event.messages in place.const returned = handlerResult?.messages ?? (sameMessages(visibleMessages, visibleSnapshot) ? undefined : visibleMessages);if (!returned) continue; currentMessages = restoreSystemMessages(currentMessages, visibleSnapshot, returned); } catch (err) {this.emitError({ ... }); } } } ...}context事件的handler只看到对话,看不到system message(提示词和工具声明)。handler返回新数组或者原地改都支持,靠sameMessages判断有没有动过
functionsameMessages(left: AgentMessage[], right: AgentMessage[]): boolean {return left.length === right.length && left.every((message, index) => message === right[index]);}/** * Re-attach the prompt and tool state after a `context` handler. Handlers only see the * conversation; the system messages belong to Pi. An unchanged conversation keeps every * system message in place, so models with mid-conversation support keep their cached * prefix. A changed one gets the replayed prompt sections and tool declarations as one * leading system message, so pruning, windowing, or slicing from a compaction summary * cannot drop them. */functionrestoreSystemMessages(current, visible, returned): AgentMessage[] {if (sameMessages(returned, visible)) return current;const head = getCurrentSystemMessage(current);return head ? [head, ...returned] : returned;}分两支:没改就原样返回,连system message的位置都不动,支持中途system message的provider(Anthropic)因此保住缓存前缀;改了(剪枝、窗口化、摘要切片)就把所有system message重放合并成一条放回头部,因为provider是从头部system message读提示词和初始工具声明的
sameMessages是引用比较不是深比较,slice()只是浅拷贝,所以原地改消息对象属性检测不出来,只能检测增删改顺序。应该是故意的,深比较几十个字段不划算,而增删消息才是handler的主要行为
context_with_system是另一个循环,handler能看到完整transcript,输出直接采用,但会检测头部system message有没有被丢掉并报错(报错但仍然采用handler的结果)
emitBoundary
turn_end和agent_before_settle共用一个实现,handler可以链式追加compaction/context_edit/custom entry草稿还能请求再跑一轮模型请求
for (const { ext, handlers } ofsnapshotEventHandlers(this.extensions, baseEvent.type)) {for (const handler of handlers) {const event = { ...baseEvent, entries, continue: shouldContinue, context };try {const handlerResult = (awaithandler(event, ctx)) asBoundaryResult | undefined;if (handlerResult?.entries !== undefined) entries = handlerResult.entries;if (handlerResult?.continue !== undefined) shouldContinue = handlerResult.continue; } catch (err) {this.emitError({ ... }); }try { context = awaitbuildContext(entries); valid = true; } catch (err) { valid = false;this.emitError({ ..., error: `Invalid boundary entries: ${err.message}` }); } }}return valid ? { entries, continue: shouldContinue, context, valid: true } : { entries: [], continue: false, context, valid: false };每个handler之后立刻用buildContext(entries)验证草稿并把预览回显给下一个handler。草稿非法就整体丢弃所有entries——因为entry草稿是要进transcript的,一部分合法一部分不合法会让会话处于难以推理的状态
生命周期:stale与reload
invalidate和assertActive是一对
invalidate(message = "This extension ctx is stale after session replacement or reload. ..."): void {if (!this.staleMessage) {this.staleMessage = message;this.runtime.invalidate(message); }}privateassertActive(): void {if (this.staleMessage) {thrownewError(this.staleMessage); }}runtime那边还会顺手拆掉所有event bus订阅,防止热重载之后老扩展的监听泄漏
reload在agent-session.ts
asyncreload(options?): Promise<void> {const oldRunner = this._extensionRunner;const previousFlagValues = oldRunner.getFlagValues();awaitemitSessionShutdownEvent(oldRunner, { type: "session_shutdown", reason: "reload" }); oldRunner.invalidate();// ... settingsManager.reload() / resetApiProviders() / _resourceLoader.reload()for (const name ofthis.getActiveToolNames()) this._pendingToolNames.add(name);this._buildRuntime({activeToolNames: [...this.getActiveToolNames(), ...addedDefaultTools],flagValues: previousFlagValues,includeAllExtensionTools: true, });if (hasBindings) {await options?.beforeSessionStart?.();awaitthis._extensionRunner.emit({ type: "session_start", reason: "reload" });this._extensionRunner.reportUnhandledMcpServers();awaitthis.extendResourcesFromExtensions("reload"); }}flag值和active tool集合都要跨越reload保留
那agent上的hook为什么不用重装?agent-session.ts里注释写得很清楚
/** * Install tool hooks once on the Agent instance. * The callbacks read `this._extensionRunner` at execution time, so extension reload swaps in the * new runner without reinstalling hooks. */private_installAgentToolHooks(): void {this.agent.beforeToolCall = (context) =>this._beforeToolCall(context);this.agent.afterToolCall = (context) =>this._afterToolCall(context);}每次现取runner,runner换了hook自动跟着换。和之前讲的transformContext链式包装是同一类技巧
顺带一个零扩展时的快路径
privateasync_beforeToolCall({ toolCall, args }, parentToolCallId?): Promise<BeforeToolCallResult | undefined> {const runner = this._extensionRunner;if (!runner.hasHandlers("tool_call")) {returnundefined; } ...}41个事件×每次工具调用,没有扩展的时候一次遍历就返回了
跨扩展通信
core/event-bus.ts全文就33行,套了个EventEmitter
exportfunctioncreateEventBus(): EventBusController {const emitter = newEventEmitter();return {emit: (channel, data) => { emitter.emit(channel, data); },on: (channel, handler) => {constsafeHandler = async (data: unknown) => {try { awaithandler(data); }catch (err) { console.error(`Event handler error (${channel}):`, err); } }; emitter.on(channel, safeHandler);return() => emitter.off(channel, safeHandler); },clear: () => { emitter.removeAllListeners(); }, };}同步emit异步handler,handler错误吞掉只log——跨扩展通信不该因为一个订阅者出错就影响emit的一方
在loader.ts里挂上生命周期追踪,pi.events.on返回的注销函数会被trackEventBusSubscription包一层,保证幂等,并且invalidate的时候统一清理
events: {emit(channel, data) { assertActive(); eventBus.emit(channel, data); },on(channel, handler) {assertActive();const unsubscribe = runtime.trackEventBusSubscription(eventBus.on(channel, handler));if (state === "loading") loadingUnsubscribers.push(unsubscribe);return unsubscribe; },},看examples/extensions/event-bus.ts就很清楚,扩展把ctx存起来,在bus回调里用
letcurrentCtx: ExtensionContext | undefined;pi.on("session_start", async (_event, ctx) => { currentCtx = ctx; });pi.events.on("my:notification", (data) => {const { message, from } = data as { message: string; from: string }; currentCtx?.ui.notify(`Event from ${from}: ${message}`, "info");});ctx在捕获的时候不检查,currentCtx?.ui.notify触发getter的时候才检查。reload之后如果这个扩展没被重新加载,notify会抛stale错误,而不是静默操作旧UI。这个就是全getter设计的用武之地
最后
写这种扩展框架最难的不是功能,是边界。比如tool_call的handler故意不try/catch,靠下游prepareToolCall兜底成错误结果;比如context_with_system的handler把头部system message搞丢了,只报错不强制恢复。这些决定的共同点是"报错但尊重handler的输出",因为扩展作者比框架更清楚自己要干什么,框架的职责是说清楚"这样做的后果是什么",而不是替他做决定
后面还没有读的部分应该是压缩和tree结构,但是我感觉前者设计的一般(,后者想和其他的插件一起搞,所以就先不写了,先研究agent开发了