@mariozechner/pi-agent-core
支持工具执行和事件流的有状态 Agent。基于 @mariozechner/pi-ai 构建。
Quick Start
import { Agent } from "@mariozechner/pi-agent-core";
import { getModel } from "@mariozechner/pi-ai";
const agent = new Agent({
initialState: {
systemPrompt: "You are a helpful assistant.",
model: getModel("anthropic", "claude-sonnet-4-20250514"),
},
});
agent.subscribe((event) => {
if (event.type === "message_update" && event.assistantMessageEvent.type === "text_delta") {
// 只流式输出新的文本块
process.stdout.write(event.assistantMessageEvent.delta);
}
});
await agent.prompt("Hello!");
核心概念
AgentMessage 与 LLM 消息
Agent 使用 AgentMessage,这是一种灵活的类型,可以包括:
- 标准 LLM 消息(
user、assistant、toolResult) - 通过声明合并的自定义应用特定消息类型
LLM 只理解 user、assistant 和 toolResult。convertToLlm 函数通过在每次 LLM 调用前过滤和转换消息来弥合这一差距。
消息流
AgentMessage[] → transformContext() → AgentMessage[] → convertToLlm() → Message[] → LLM
(可选) (必需)
- transformContext:修剪旧消息,注入外部上下文
- convertToLlm:过滤仅用于 UI 的消息,将自定义类型转换为 LLM 格式
事件流
Agent 发出事件以更新 UI。理解事件序列有助于构建响应式界面。
prompt()事件流
当你调用 prompt("Hello") 时:
prompt("Hello")
├─ agent_start
├─ turn_start
├─ message_start { message: userMessage } // 你的提示
├─ message_end { message: userMessage }
├─ message_start { message: assistantMessage } // LLM 开始响应
├─ message_update { message: partial... } // 流式块
├─ message_update { message: partial... }
├─ message_end { message: assistantMessage } // 完整响应
├─ turn_end { message, toolResults: [] }
└─ agent_end { messages: [...] }
带工具调用
如果Agent调用工具,循环会继续:
prompt("Read config.json")
├─ agent_start
├─ turn_start
├─ message_start/end { userMessage }
├─ message_start { assistantMessage with toolCall }
├─ message_update...
├─ message_end { assistantMessage }
├─ tool_execution_start { toolCallId, toolName, args }
├─ tool_execution_update { partialResult } // 如果工具流式输出
├─ tool_execution_end { toolCallId, result }
├─ message_start/end { toolResultMessage }
├─ turn_end { message, toolResults: [toolResult] }
│
├─ turn_start // 下一轮
├─ message_start { assistantMessage } // LLM 响应工具结果
├─ message_update...
├─ message_end
├─ turn_end
└─ agent_end
工具执行模式可配置:
parallel(默认):按顺序预检工具调用,并发执行允许的工具,按助手源顺序发出最终的tool_execution_end和toolResult消息sequential:逐个执行工具调用,匹配历史行为
beforeToolCall 钩子在 tool_execution_start 和参数验证解析之后运行。它可以阻止执行。afterToolCall 钩子在工具执行完成后、发出 tool_execution_end 和最终工具结果消息事件之前运行。
当你使用 Agent 类时,助手 message_end 处理被视为工具预检开始前的屏障。这意味着 beforeToolCall 看到的 Agent 状态已经包含了请求工具调用的助手消息。
continue() 事件序列
continue() 从现有上下文恢复,而不添加新消息。用于错误后的重试。
// 错误后,从当前状态重试
await agent.continue();
上下文中的最后一条消息必须是 user 或 toolResult(不能是 assistant)。
事件类型
| 事件 | 描述 |
|---|---|
agent_start |
Agent 开始处理 |
agent_end |
Agent 完成,包含所有新消息 |
turn_start |
新轮次开始(一次 LLM 调用 + 工具执行) |
turn_end |
轮次完成,包含助手消息和工具结果 |
message_start |
任何消息开始(user、assistant、toolResult) |
message_update |
仅助手。 包含带增量的 assistantMessageEvent |
message_end |
消息完成 |
tool_execution_start |
工具开始 |
tool_execution_update |
工具流式进度 |
tool_execution_end |
工具完成 |
Agent配置选项
const agent = new Agent({
// 初始状态
initialState: {
systemPrompt: string,
model: Model<any>,
thinkingLevel: "off" | "minimal" | "low" | "medium" | "high" | "xhigh",
tools: AgentTool<any>[],
messages: AgentMessage[],
},
// 将 AgentMessage[] 转换为 LLM Message[](自定义消息类型必需)
convertToLlm: (messages) => messages.filter(...),
// 在 convertToLlm 之前转换上下文(用于修剪、压缩)
transformContext: async (messages, signal) => pruneOldMessages(messages),
// 引导模式:"one-at-a-time"(默认)或 "all"
steeringMode: "one-at-a-time",
// 跟进模式:"one-at-a-time"(默认)或 "all"
followUpMode: "one-at-a-time",
// 自定义流函数(用于代理后端)
streamFn: streamProxy,
// 用于提供商缓存的会话 ID
sessionId: "session-123",
// 动态 API 密钥解析(用于过期的 OAuth 令牌)
getApiKey: async (provider) => refreshToken(),
// 工具执行模式:"parallel"(默认)或 "sequential"
toolExecution: "parallel",
// 参数验证后预检每个工具调用。可以阻止执行。
beforeToolCall: async ({ toolCall, args, context }) => {
if (toolCall.name === "bash") {
return { block: true, reason: "bash is disabled" };
}
},
// 在发出最终工具事件之前后处理每个工具结果。
afterToolCall: async ({ toolCall, result, isError, context }) => {
if (!isError) {
return { details: { ...result.details, audited: true } };
}
},
// 基于令牌的提供商的自定义思考预算
thinkingBudgets: {
minimal: 128,
low: 512,
medium: 1024,
high: 2048,
},
});
Agent 状态
interface AgentState {
systemPrompt: string;
model: Model<any>;
thinkingLevel: ThinkingLevel;
tools: AgentTool<any>[];
messages: AgentMessage[];
isStreaming: boolean;
streamMessage: AgentMessage | null; // 流式传输期间的当前部分消息
pendingToolCalls: Set<string>;
error?: string;
}
通过 agent.state 访问。在流式传输期间,streamMessage 包含部分的助手消息。
方法
prompt
// 文本提示
await agent.prompt("Hello");
// 带图片
await agent.prompt("What's in this image?", [
{ type: "image", data: base64Data, mimeType: "image/jpeg" }
]);
// 直接使用 AgentMessage
await agent.prompt({ role: "user", content: "Hello", timestamp: Date.now() });
// 从当前上下文继续(最后一条消息必须是 user 或 toolResult)
await agent.continue();
状态管理
agent.setSystemPrompt("New prompt");
agent.setModel(getModel("openai", "gpt-4o"));
agent.setThinkingLevel("medium");
agent.setTools([myTool]);
agent.setToolExecution("sequential");
agent.setBeforeToolCall(async ({ toolCall }) => undefined);
agent.setAfterToolCall(async ({ toolCall, result }) => undefined);
agent.replaceMessages(newMessages);
agent.appendMessage(message);
agent.clearMessages();
agent.reset(); // 清除所有内容
会话和思考预算
agent.sessionId = "session-123";
agent.thinkingBudgets = {
minimal: 128,
low: 512,
medium: 1024,
high: 2048,
};
控制
agent.abort(); // 取消当前操作
await agent.waitForIdle(); // 等待完成
事件
const unsubscribe = agent.subscribe((event) => {
console.log(event.type);
});
unsubscribe();
引导和跟进
引导消息让你在工具运行时中断 Agent。跟进消息让你在 Agent 本应停止后排队工作。
agent.setSteeringMode("one-at-a-time");
agent.setFollowUpMode("one-at-a-time");
// 当 Agent 正在运行工具时
agent.steer({
role: "user",
content: "Stop! Do this instead.",
timestamp: Date.now(),
});
// 在 Agent 完成当前工作后
agent.followUp({
role: "user",
content: "Also summarize the result.",
timestamp: Date.now(),
});
const steeringMode = agent.getSteeringMode();
const followUpMode = agent.getFollowUpMode();
agent.clearSteeringQueue();
agent.clearFollowUpQueue();
agent.clearAllQueues();
使用 clearSteeringQueue、clearFollowUpQueue 或 clearAllQueues 来丢弃排队的消息。
当轮次完成后检测到引导消息时:
- 当前助手消息中的所有工具调用已经完成
- 引导消息被注入
- LLM 在下一轮响应
跟进消息仅在不再有工具调用且没有引导消息时检查。如果有任何排队的消息,它们会被注入并运行另一轮。
自定义消息类型
通过声明合并扩展 AgentMessage:
declare module "@mariozechner/pi-agent-core" {
interface CustomAgentMessages {
notification: { role: "notification"; text: string; timestamp: number };
}
}
// 现在有效
const msg: AgentMessage = { role: "notification", text: "Info", timestamp: Date.now() };
在 convertToLlm 中处理自定义类型:
const agent = new Agent({
convertToLlm: (messages) => messages.flatMap(m => {
if (m.role === "notification") return []; // 过滤掉
return [m];
}),
});
工具
使用 AgentTool 定义工具:
import { Type } from "@sinclair/typebox";
const readFileTool: AgentTool = {
name: "read_file",
label: "Read File", // 用于 UI 显示
description: "Read a file's contents",
parameters: Type.Object({
path: Type.String({ description: "File path" }),
}),
execute: async (toolCallId, params, signal, onUpdate) => {
const content = await fs.readFile(params.path, "utf-8");
// 可选:流式进度
onUpdate?.({ content: [{ type: "text", text: "Reading..." }], details: {} });
return {
content: [{ type: "text", text: content }],
details: { path: params.path, size: content.length },
};
},
};
agent.setTools([readFileTool]);
错误处理
当工具失败时抛出错误。不要将错误消息作为内容返回。
execute: async (toolCallId, params, signal, onUpdate) => {
if (!fs.existsSync(params.path)) {
throw new Error(`File not found: ${params.path}`);
}
// 仅在成功时返回内容
return { content: [{ type: "text", text: "..." }] };
}
抛出的错误会被 Agent 捕获,并以 isError: true 作为工具错误报告给 LLM。
使用代理
import { Agent, streamProxy } from "@mariozechner/pi-agent-core";
const agent = new Agent({
streamFn: (model, context, options) =>
streamProxy(model, context, {
...options,
authToken: "...",
proxyUrl: "https://your-server.com",
}),
});
底层API
import { agentLoop, agentLoopContinue } from "@mariozechner/pi-agent-core";
const context: AgentContext = {
systemPrompt: "You are helpful.",
messages: [],
tools: [],
};
const config: AgentLoopConfig = {
model: getModel("openai", "gpt-4o"),
convertToLlm: (msgs) => msgs.filter(m => ["user", "assistant", "toolResult"].includes(m.role)),
toolExecution: "parallel",
beforeToolCall: async ({ toolCall, args, context }) => undefined,
afterToolCall: async ({ toolCall, result, isError, context }) => undefined,
};
const userMessage = { role: "user", content: "Hello", timestamp: Date.now() };
for await (const event of agentLoop([userMessage], context, config)) {
console.log(event.type);
}
// 从现有上下文继续
for await (const event of agentLoopContinue(context, config)) {
console.log(event.type);
}