学习Agent开发8 高级Agent框架Mastra基础用法
Mastra是什么
Mastra是一个由ts编写的Agent开发框架ai-sdk langgraph 是低级的框架 因为它们关注的单次toolcall的自动响应以及ai会话本身 而Mastra则封装了ai-sdk的能力 提供了上下文管理、会话持久化、rag、skill、长期记忆等能力
该库采用 Apache License 2.0,商业使用限制宽松,只有使用云服务时才需要付费。
Mastra的上下文管理方式
Mastra提出了名为 Observational Memory 的新型上下文管理方法。
先设定token上限,默认30k;主会话每次达到上限时,就调用另一个小模型gpt-5-mini进行日志总结,日志再次达到上限后再做一轮总结。
系统同时提供名为recall的tool,可以对当前会话及以往其他会话进行精确的消息擦护照,从而确保模型始终拥有充裕上下文。
构造Agent
import { createOpenAI } from "@ai-sdk/openai";import { Agent } from "@mastra/core/agent";import { createTool } from "@mastra/core/tools";import { createSkill } from "@mastra/core/skills";import { Memory } from "@mastra/memory";import { PostgresStore } from "@mastra/pg";import { z } from "zod";const API_KEY = process.env.API_KEY!;const BASE_URL = process.env.BASE_URL!;/** 主会话模型 */const CHAT_MODEL_ID = process.env.CHAT_MODEL_ID!;/** OM模型 */const OM_MODEL_ID = process.env.OM_MODEL_ID!;const POSTGRES_URL = process.env.POSTGRES_URL!;// 模型providerconst provider = createOpenAI({apiKey: API_KEY,baseURL: BASE_URL,});// 数据存储层 当前采用postgresexport const store = new PostgresStore({connectionString: POSTGRES_URL,id: "mastra-demo",});// 消息持久化和上下文管理层const memory = new Memory({storage: store,options: {// om配置observationalMemory: {model: provider.languageModel(OM_MODEL_ID),// om模型的总结范围 仅限当前会话scope: "thread",// 主会话的检索范围 仅限当前会话retrieval: { scope: "thread" },},// 生成标题generateTitle: {model: provider.languageModel(OM_MODEL_ID),instructions: "根据用户第一条消息生成一个简短的会话标题",}},});// 需要审批的toolexport const sendEmail = createTool({id: "sendEmail",description: "发邮件",inputSchema: z.object({body: z.string(),subject: z.string(),to: z.string(),}),outputSchema: z.object({status: z.union([z.literal("success"), z.literal("failed")]),reason: z.string().optional(),}),// 预审批 可以传函数// 拒绝时无法传入理由requireApproval: true,// 执行过程中可以中断 进行二次确认suspendSchema: z.object({body: z.string(),subject: z.string(),to: z.string(),}),// 最后一次恢复得到的值resumeSchema: z.object({approved: z.boolean(),reason: z.string().optional(),}),// 如果抛出异常 其消息会被作为tool-call-errorexecute: async (input, context) => {// input可以信任 但resumeData不可信任const resumeData = context.agent?.resumeData;if (resumeData) {if (!resumeData.approved) {return { status: "failed" as const, reason: resumeData.reason };}} else if (input.to === "[email protected]") {return await context.agent?.suspend(input);}/** 发送邮件 */return { status: "success" as const };},});const testSkill = createSkill({name: "xx",description: "xx",instructions: "skill-body",})export const agent = new Agent({id: "demo",name: "demo",// 系统提示词 允许异步函数返回stringinstructions: "系统提示词",memory,model: provider.languageModel(CHAT_MODEL_ID),tools: { sendEmail },skills: [testSkill],});
关于requireApproval和suspend
不要同时使用两种机制
resumeData负责让挂起的流恢复,而二者底层采用的是同一套行为逻辑
二者同时开启时 会先经过requireApproval的校验 检查resumeData的approved属性 不为true就直接以Tool call was not approved by the user作为toolcall的结果 即使之前已经approval过了
关于Subagent
不会抛出异常,是因为subagent在调用需要resume的tool之后就停止执行;此时已经生成的text会作为result交回主会话
使用Agent
Agent必须先注册到Mastra实例中,之后才能进行多轮对话。
export const mastra = new Mastra({agents: { demo: agent },storage: store,})
agent对话涉及三个id resourceId threadId runId
简单理解当前映射关系:每轮会话的id对应runId,sessionId对应threadId,userId对应resoueceId
发起流式消息(开始会话或多轮对话)
// stream.id就是runIdconst stream = await agent.stream('xx', {memory: { thread: threadId, resource: resourceId },})for await (const c of stream.fullStream){/** */}
从挂起的流恢复
// 通过予审批(拒绝也可以)// approved是一个新的流const approved = agent.approveToolCall({ runId, toolCallId })// resume一个挂起的toolconst resumed = agent.resumeStream(resumeData, { runId, toolCallId })
停止生成
// 停止所有指定thread的runagent.abortThreadStream({resourceId,threadId})// 停止指定runagent.abortRunStream(runId)
sendMessage/sendSignal
agent.sendMessage用于向thread中加入用户消息 可以基于当前会话fork新thread 可以等待当前run生成完毕 适用于一个thread多人参与的场合(比如群聊Agent)
agent.sendSignal它与sendMessage相似,但插入的是系统消息而非用户消息,因此呈现给模型的方式有所不同。
这两个函数自身都不返回流,需要提前订阅相应的流。
const subscription = await agent.subscribeToThread({ resourceId, threadId })for await (const c of subscription.stream) {}subscription.unsubscribe()
resume同上
AgentController
自动会话管理、消息队列和会话fork等能力都包含在mastra预设的Agent运行时AgentController中;ui更新、工具审批等一系列api也由它封装
这个api目前仍处在beta阶段,文档不够准确,使用体验很差,但未来可期。
resource,thread和run
记忆层由thread和resource构成,其中每个thread都必定归属于一个resouce
run属于执行层。每次调用stream都会继承thread中的历史消息并开启新对话,等本轮内容全部生成后,run会被删除。
在tool 审批/suspense时 run挂起 runid不变
thread不是session
thread内存储的是完整的消息 即完整的toolcall+result/text(一次step的完整输出) 流式消息片段不会加入thread
两个tab页同时对话就是会话混乱的一种情形。原因在于mastra默认允许同时从thread发起对话,并不限制由thread产生run;随后新消息又按时间顺序进入thread。
sendMessage/sendSignal
上述stream竞态问题可由sendMessage缓解,但只能在一定程度上解决
agent.sendMessage('Compare that with the previous option.', {resourceId,threadId,// thread存在活跃/挂起run的行为ifActive: {},// thread不存在活跃run的行为(即 thread空闲)ifIdle: {},})
ifAvcive
{behavior?: "deliver" | "persist" | "discard"}
三种处理方式:discard-丢弃;persist-消息加入thread,当前run无法感知;deliver(默认)-放入当前run的下一条消息中
ifIdle
{behavior?: "persist" | "discard" | "wake"}
wake(默认)用于唤起新run,其余处理方式同上。
上述机制仅在单个后端实例中生效。
多实例环境中的竞态问题
官方提供的解决方案是使用@mastra/redis-streams的RedisStreamsPubSub
在new Mastra时传入pubsub属性,用来处理竞态问题
前端展示
next+ai-sdk/react+ai-elements用于展示mastra十分合适,因为mastra本身构建于ai-sdk之上
无法直接使用的原因是,mastra引入ai-sdk时并未采用peerDependencies方式
可以参考官方文档完成基础设施搭建,并进行以下优化。
解决类型错误
api/chat/route.ts POST函数
const stream = await handleChatStream({// 增加versionversion: "v6",});
不要改GET的内容 v5的sdk messages是正确的版本
增加approval和suspend能力
目标是让流式传输与历史消息都能采用相同方式向用户呈现挂起状态,并支持恢复会话。
这两种情况下的消息结构并不相同。
approval:历史消息缺少tool-approval-request;但在流式传输期间,这条消息位于toolcall和data-tool-call-approval之间
suspense的历史消息与流式传输表现一致:data-tool-call-suspended消息始终紧跟toolcall
这里采用的会话恢复方案是
const { sendMessage } = useChat({/* xx */})const handleResume = (runId: string, resumeData: Record<string, unknown>) =>sendMessage(undefined, { body: { runId, resumeData } })
body会被传给后端并用于恢复会话。前文已经说过 approval和suspended本质是一套机制 可以这样写 这种写法还抹平了流式传输和历史会话消息结构的差异
前端进行以下修改
messages.map(message=>{message.parts?.map((part, i)=>{if (part.type === 'data-tool-call-approval') {const data = part.data as ToolCallApprovalDatareturn (<><Button onClick={() => handleResume(data.runId, { approved: false })}>拒绝</Button><Button onClick={() => handleResume(data.runId, { approved: true })}>同意</Button></>)}if (part.type === 'data-tool-call-suspended') {const data = part.data as ToolCallSuspendedData// 同上 resumeData按照tool要求的格式即可}})})
以上代码仅作示例 实际使用时还需要考虑这个挂起是否已经处理过
-
07.29
夜幕之下预抽卡入口位置在哪里
-
07.29
傲视天下礼包码可以在哪里领取
-
07.29
无畏契约张家界地图有什么爆料
-
07.29
挖掘者米娜躲闪钟摆饰品如何获得
-
07.29
流放之路2深渊天赋怎样加点
-
07.29
耀世格斗必练的三大职业有什么推荐
-
-
下载
- |
-
-
下载
- 《行尸走肉第一章》免安装中文汉化硬盘版下载
- 单机|436 MB
- 一款以动作冒险为主题的游戏
-
-
下载
- 《街头霸王X铁拳》免安装中文汉化硬盘版下载
- 单机|111MB
- 一款非常好玩的格斗游戏
-
-
下载
- |
-
-
下载
- 《暗黑破坏神3》免安装繁体中文正式版下载
- 单机|7630 MB
- 一款以角色扮演为主题的游戏
-
-
下载
- 《马克思佩恩3》免安装硬盘版下载
- 单机|27033 MB
- 一款以第三人称射击为主题的游戏