ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Mastra Durable Agents 深度解析:让 Agent 循环与流式输出具备崩溃恢复与断线重连能力

2026/9/13 11:52:01 拓冰建站 浏览量
Mastra Durable Agents 深度解析:让 Agent 循环与流式输出具备崩溃恢复与断线重连能力 Mastra Durable Agents 深度解析让 Agent 循环与流式输出具备崩溃恢复与断线重连能力【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastraMastra 的 Durable Agent持久化 Agent能力是构建在packages/core之上的一套让 AI Agent 执行循环可跨崩溃存活的架构方案服务端在 10 步 Agent 循环的任何一步崩溃都能从断点而非起点恢复客户端在网络抖动或页面刷新导致断流后也能通过缓存历史事件从丢失位置无缝续流。本文以仓库中的 durable-agent.md 探索文档为主线结合packages/core、packages/server与workflows/inngest的真实实现完整拆解三层持久化模式、运行注册表、PubSub 事件流、可恢复流式输出、Inngest 分布式执行以及配套 HTTP 端点读完即可在自己的 Mastra 应用中原样落地这套方案。这套方案解决什么问题服务端视角当 AI Agent 在一次会话中多次调用工具时每一步都是潜在故障点。如果服务器在 10 步 Agent 循环进行到一半时崩溃传统实现会丢掉全部进度、必须从头开始。Durable Agent 会对整个执行过程做 checkpoint若崩溃发生在第 7 步之后重启后从第 7 步继续而不是从 0 开始。客户端视角客户端在流式传输中途断开网络闪断、页面刷新流就丢了。可恢复流Resumable Stream让客户端重连后能从缓存的历史事件中回放错过的内容从离开的位置精确续接。三种递增的持久化模式Durable Agent 架构支持三种逐级增强的持久化模式对应 探索文档 中的 Three PatternscreateDurableAgent— 仅可恢复流。执行仍停留在 HTTP 请求内。适合需要重连支持但不需要持久化执行能力的场景。createEventedAgent— 可恢复流 基于工作流引擎的即发即忘fire-and-forget执行。LLM 循环在本地工作流中运行与 HTTP 请求解耦。适合单实例部署上的长耗时操作。createInngestAgent— 可恢复流 Inngest 驱动的持久化执行。LLM 循环作为 Inngest 函数运行具备 checkpoint、重试与跨进程流式能力。适合生产级分布式系统。三者之间是逐层叠加的关系EventedAgent继承DurableAgentInngestAgent复用同一套准备preparation与流式适配stream adapter逻辑只是把执行引擎替换为 Inngest。核心基础设施DurableAgent 与四类支撑组件DurableAgent 类packages/core/src/agent/durable/durable-agent.tsDurableAgent继承Agent将准备与执行分离stream()— 一次调用完成准备与执行返回带流式回调的MastraModelOutput。prepare()— 生成可序列化的工作流输入非持久化部分并填充运行注册表run registry。resume()— 恢复一个挂起的工作流执行例如工具审批通过后。observe()— 通过runId重连到已存在的流可选从某个 offset 开始借助CachingPubSub回放错过的事件。从当前源码看这一 API 面已进一步扩充generate()/resumeGenerate()将运行收敛为单个FullOutput挂起时finishReason为suspended、recover()从持久化的工作流快照重建不可序列化状态并返回新的流式结果、recoverActiveRuns()批量重启孤立运行、listActiveRuns()分页查询当前running状态的运行等能力均已落地。stream()返回的DurableAgentStreamResult包含outputMastraModelOutput、fullStreamReadableStream兼容服务端、runId、threadId/resourceId使用 memory 时、cleanup()取消 PubSub 订阅以及abort(reason?)翻转内部AbortController并跨进程发布 abort 请求安全地幂等。运行注册表packages/core/src/agent/durable/run-registry.ts工作流状态必须可序列化而工具对象、SaveQueueManager、实时MessageList等无法序列化。运行注册表以runId为键保存这些运行时状态RunRegistry— 每个DurableAgent实例级注册表提供register/get/getTools/getModel/getSaveQueueManager/cleanup等操作。ExtendedRunRegistry— 额外维护MessageList引用与 threadId/resourceId 记忆信息供回调读取消息状态。globalRunRegistry— 模块级TTLCache来自isaacs/ttlcache让工作流步骤也能按runId拿到运行时状态。探索文档最初记录的规格是10 分钟 TTL、1000 条上限当前实现已演进为RUN_REGISTRY_TTL_MS 2 * 60 * 60 * 10002 小时滑动 TTLupdateAgeOnGet: true与max: 1000的硬上限并在dispose回调中调用entry.cleanup()。源码注释明确强调TTL 只是崩溃/泄漏的安全网而非活动超时——运行中的工具调用、子 Agent 或模型调用可能长时间无人读取该条目此时通过markRunActive()引用计数 每 60 秒一次的runActivityHeartbeat心跳保持条目存活避免把运行中的实时运行时状态提前清掉。PubSub 事件流系统packages/core/src/events/闭包无法序列化因此跨进程流式依赖 PubSub 事件分发PubSub抽象基类— 定义 subscribe / publish / unsubscribe 契约。EventEmitterPubSub— 基于 Node.jsEventEmitter的内存实现。CachingPubSub— 装饰器为任意 PubSub 叠加事件缓存与回放能力提供subscribeWithReplay()与subscribeFromOffset()是实现可恢复流的钥匙。subscribeFromOffset()采用四阶段引导先订阅 live 事件并缓冲、按序投递缓存历史、排空缓冲并去重、切换到基于 index 水位线的直通模式。InngestPubSub位于workflows/inngest/src/pubsub.ts— 基于 Inngest Realtime 的实现publish()走inngest.realtime.publish()subscribe()走inngest/realtime主题agent.stream.{runId}映射到 Inngest 通道agent:{runId}的agent-streamtopic工作流事件workflow.events.v2.{runId}映射到通道workflow:{workflowId}:{runId}的watchtopic。CachingPubSub的缓存后端抽象为MastraServerCache默认InMemoryServerCacheRedis 等实现通过独立包提供cache配置传false则完全禁用缓存流不可恢复。流适配器packages/core/src/agent/durable/stream-adapter.tscreateDurableAgentStream()订阅 PubSub 通道把事件翻译为 chunk 并驱动MastraModelOutput的回调onChunk、onStepFinish、onFinish、onError、onSuspended、onAbort、onIterationComplete。它以cancelled标志处理订阅与取消的竞态支持从 offset 的按索引回放supportsOffsets为 true 时走subscribeFromOffset否则退化为从最新开始实时订阅还内置了 idle 看门狗——idleTimeoutMsisAlive活性探针避免生产者进程崩溃后从未发布终态事件导致 observe 永久悬挂。流事件类型定义于 constants.ts事件含义chunk流式数据块文本、工具调用等step-startAgent 循环中新步骤开始step-finish步骤结束finish执行成功完成error执行出错suspended工作流挂起如等待工具审批abort被 abortSignal 中止iteration-complete单次循环迭代完成可观测性钩子流事件发布在主题agent.stream.{runId}AGENT_STREAM_TOPIC控制消息如 abort 请求走独立的agent.control.{runId}AGENT_CONTROL_TOPIC两者分离避免相互干扰。持久化 Agent 循环可复用的工作流步骤DurableAgent把 Agent 循环建模为可 checkpoint 的工作流位于packages/core/src/agent/durable/workflows/steps/createDurableLLMExecutionStep()llm-execution.ts— 从工作流状态反序列化 messageList按modelConfig从 Mastra 解析模型执行 LLM 调用通过 PubSub 发射 chunk把 messageList 序列化回输出。createDurableToolCallStep()tool-call.ts— 从注册表解析工具检查审批要求需要审批时先 flush 消息再挂起执行工具支持执行中挂起回调与后台任务执行通过 PubSub 发射结果/错误。createDurableLLMMappingStep()llm-mapping.ts— 把工具调用映射为单个工具调用步骤的输入。createDurableScorerStep()scorer-execution.ts— 每步之后运行配置的 scorers 对输出评分。当前目录还扩展出is-task-complete.ts任务完成判定、goal.ts目标活动追踪、background-task-check.ts后台任务等待、signal-drain.ts信号排空等步骤。主工作流由create-durable-agentic-workflow.ts的dowhile(shouldContinue)循环组装步骤 ID 定义在DurableStepIds如durable-llm-execution、durable-tool-call、durable-agentic-loop。Inngest 集成workflows/inngestcreateInngestAgent()create-inngest-agent.ts把普通 Agent 包装成 Inngest 驱动的持久化 Agentconst agent new Agent({ id: my-agent, model: openai(gpt-4), ... }); const durableAgent createInngestAgent({ agent, inngest }); const mastra new Mastra({ agents: { myAgent: durableAgent } });返回的InngestAgent是一个 Proxy 对象stream/resume/observe/prepare/generate/resumeGenerate为持久化执行面其余 Agent 方法listTools、getMemory等转发给底层 Agent。缓存解析顺序为用户提供 mastra.serverCacheInMemoryServerCache兜底跨进程 observe 必须提供共享缓存后端如 Redis。内部 pubsub 默认InngestPubSub再包一层CachingPubSub若用户已传CachingPubSub则复用避免双重包裹。执行流程stream()先走prepareForDurableExecution生成可序列化workflowInput安装本段AbortController把非序列化状态写入globalRunRegistry然后createDurableAgentStream建立订阅ready之后再通过inngest.send()触发事件workflow.{AGENTIC_LOOP}。步骤 worker 与调用方通常不在同一进程因此abort()除了翻转本地 controller 外还会通过 PubSub 发布 abort 请求让执行 worker 优雅收尾并发射终态流事件。resume()在全新进程中也确保重建 registry 条目带isPlaceholder标记由resolveRuntimeDependencies按需从 Mastra 实例重建运行时状态。服务端 HTTP 端点packages/server两个核心端点定义在 agents.tsPOST /agents/:agentId/observeOBSERVE_AGENT_STREAM_ROUTE— 按runId重连到已有流body 支持offset做部分回放有 offset 走subscribeFromOffset否则subscribeWithReplay返回 SSE 事件流。为持久化 Agent 专门读取其自身的CachingPubSub而非mastra.pubsub内置 5 分钟 idle 超时防止 agent 崩溃后订阅泄漏。POST /agents/:agentId/resume-streamRESUME_STREAM_ROUTE— 携带新数据如工具审批结果恢复挂起的执行返回续跑的 SSE 流。此外还有POST /agents/:agentId/resume-stream-until-idleL3011与untilIdle选项保持外层流在后台任务续跑期间不关闭。Mastra 集成自动注册Mastra.addAgent()识别DurableAgentLike见packages/core/src/mastra自动注册底层 agent通过getDurableWorkflows()自动注册关联工作流无需手动注册工作流。普通 Agent 还有更省事的声明式入口在AgentConfig上配置durable选项types.ts 中的AgentDurableOptiontrue表示注册到 Mastra 时用默认参数自动包裹成createDurableAgent传对象则转发cache/pubsub/maxSteps/cleanupTimeoutMs/shouldCache/id/nameconst agent new Agent({ id: my-agent, instructions: You are helpful, model: openai(gpt-4), durable: true, // 注册到 Mastra 时自动包装为 durable agent }); const mastra new Mastra({ agents: { myAgent: agent } });关键配置参数一览DurableAgentConfig/CreateDurableAgentOptions核心参数参数默认值说明agent—被包裹的Agent实例其余方法均委托给它id/nameagent.id/agent.name标识覆盖pubsubEventEmitterPubSub流式事件通道cache继承 Mastra 或InMemoryServerCache流事件缓存后端false禁用缓存流不可恢复maxSteps工作流默认DurableAgentDefaults.MAX_STEPS 5Agent 循环最大迭代次数cleanupTimeoutMs30_00030 秒流结束后注册表条目的自动清理宽限期0禁用需手动cleanup()shouldCache全部缓存按主题 opt-out 回放缓存热点主题可换最低发布延迟运行局部主题始终排除stream()的执行选项还包括stopWhen停止条件闭包形式驻留进程内注册表跨进程引擎退化为仅maxSteps、requireToolApproval全部/按调用策略、toolCallConcurrency默认 10可审批/挂起的工具集始终串行、structuredOutput、memory、untilIdle、abortSignal等。端到端使用示例以下示例完整来自探索文档并在当前 API 上保持可运行import { Agent } from mastra/core/agent; import { createDurableAgent, createEventedAgent } from mastra/core/agent/durable; import { createInngestAgent } from mastra/inngest; import { Mastra } from mastra/core/mastra; // 创建一个普通 agent const agent new Agent({ id: assistant, name: Assistant, instructions: You are helpful, model: openai(gpt-4o), tools: {/* your tools */}, }); // Pattern 1: 仅可恢复流 const durableAgent createDurableAgent({ agent }); // Pattern 2: 可恢复流 工作流引擎执行 const eventedAgent createEventedAgent({ agent, pubsub }); // Pattern 3: 可恢复流 Inngest 执行 const inngestAgent createInngestAgent({ agent, inngest }); // 注册到 Mastracache 与 pubsub 可继承 const mastra new Mastra({ agents: { assistant: durableAgent }, }); // 流式执行 const { output, runId, cleanup } await durableAgent.stream( [{ role: user, content: Analyze this data and create a report }], { onChunk: chunk { /* 流式推送给客户端 */ }, onStepFinish: step { /* 每个 LLM 步骤之后调用 */ }, onFinish: result { /* 完成时调用 */ }, }, ); const text await output.text; cleanup(); // 断线后重连回放错过的事件 const { output: reconnected } await durableAgent.observe(runId, { offset: lastSeenIndex, // 从该位置开始回放 onChunk: chunk { /* 流式推送给客户端 */ }, }); // 挂起后恢复例如工具审批 const { output: resumed } await durableAgent.resume(runId, { approved: true });架构总览┌─────────────────────────────────────────────────────────────────────┐ │ DurableAgent.stream() │ └─────────────────────────────────────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────┐ │ PREPARATION PHASE (non-durable) │ │ │ │ 1. Resolve tools → store as { id, name, schema } (no execute fn) │ │ 2. Create MessageList, load memory, run input processors │ │ 3. Serialize: { messageListState, toolsMetadata, modelConfig, ... }│ │ 4. Store non-serializable state in per-run registry global │ └─────────────────────────────────────────────────────────────────────┘ │ ┌─────────────┴──────────────┐ ▼ ▼ ┌──────────────────┐ ┌──────────────────────┐ │ LocalExecutor │ │ InngestExecutor │ │ (in-process) │ │ (Inngest function) │ └────────┬─────────┘ └──────────┬───────────┘ │ │ ▼ ▼ ┌─────────────────────────────────────────────────────────────────────┐ │ DURABLE AGENTIC LOOP (workflow) │ │ │ │ ┌─────────────────────────────────────────────────────────────┐ │ │ │ dowhile(shouldContinue) │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌─────────────────────────────────────────────────────┐ │ │ │ │ │ durableLLMExecutionStep │ │ │ │ │ │ - Deserialize messageList from workflow state │ │ │ │ │ │ - Resolve model from mastra via modelConfig │ │ │ │ │ │ - Execute LLM call │ │ │ │ │ │ - Emit chunks via PubSub (agent.stream.{runId}) │ │ │ │ │ │ - Serialize messageList to output │ │ │ │ │ └─────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌─────────────────────────────────────────────────────┐ │ │ │ │ │ foreach(toolCalls) → durableToolCallStep │ │ │ │ │ │ - Resolve tool from registry via toolName │ │ │ │ │ │ - Check approval requirements │ │ │ │ │ │ - If needs approval: flush messages, suspend │ │ │ │ │ │ - Execute tool (with suspend callback for mid-exec) │ │ │ │ │ │ - Emit result/error via PubSub │ │ │ │ │ │ - Background task execution support │ │ │ │ │ └─────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌─────────────────────────────────────────────────────┐ │ │ │ │ │ scorerExecutionStep (optional) │ │ │ │ │ │ - Run configured scorers against output │ │ │ │ │ └─────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ Check stopWhen / maxSteps condition │ │ │ └──────────────────────────────────────────────────────────────┘ │ │ │ │ On finish: run output processors, persist messages to memory │ └─────────────────────────────────────────────────────────────────────┘ │ ▼ ┌─────────────────────────────────────────────────────────────────────┐ │ OUTPUT (resumable streaming) │ │ │ │ CachingPubSub caches events for replay │ │ createDurableAgentStream subscribes to PubSub channel │ │ Drives MastraModelOutput callbacks (onChunk, onStepFinish, etc.) │ │ observe() reconnects with subscribeFromOffset for missed events │ │ resume() creates new stream resumes suspended workflow │ └─────────────────────────────────────────────────────────────────────┘整个设计遵循关注点分离准备阶段非持久化产出可序列化状态执行阶段持久化在工作流中运行并逐步骤 checkpoint流式输出可恢复基于缓存 PubSub 事件。文件结构探索文档给出的布局在当前仓库中依然成立其中__tests__已从文档记录的 22 个文件大幅扩展下面仅列出核心骨架packages/core/src/agent/durable/ ├── index.ts # Exports ├── durable-agent.ts # DurableAgent class继承 Agent ├── evented-agent.ts # EventedAgent class继承 DurableAgent ├── create-durable-agent.ts # createDurableAgent() 工厂 ├── create-evented-agent.ts # createEventedAgent() 工厂 ├── types.ts # DurableAgentState、DurableStepInput 等 ├── constants.ts # AGENT_STREAM_TOPIC、步骤 ID、默认值 ├── run-registry.ts # RunRegistry、ExtendedRunRegistry、globalRunRegistry ├── stream-adapter.ts # PubSub → MastraModelOutput 适配器 emit 辅助函数 ├── preparation.ts # 准备阶段逻辑 ├── workflows/ │ ├── index.ts │ ├── create-durable-agentic-workflow.ts # 主工作流工厂 │ ├── shared/ # execute-tool-calls / iteration-state / schemas │ └── steps/ # llm-execution / tool-call / llm-mapping / scorer-execution 等 ├── utils/ # resolve-runtime / serialize-state 等 └── __tests__/ # 数十个测试文件 packages/core/src/events/ ├── pubsub.ts # 抽象 PubSub 基类 ├── event-emitter/ # EventEmitterPubSub内存实现 ├── caching-pubsub.ts # CachingPubSub为任意 PubSub 叠加回放 ├── topics.ts / codec/ / processor.ts # 主题、编解码、事件处理器 └── index.ts workflows/inngest/src/ ├── durable-agent/ │ ├── create-inngest-agent.ts # createInngestAgent() 工厂 │ └── create-inngest-agentic-workflow.ts # Inngest 工作流 ├── pubsub.ts # InngestPubSub基于 Inngest Realtime ├── execution-engine.ts / serve.ts # Inngest 执行引擎与 serve 处理器 └── index.ts关键设计决策三层持久化—DurableAgent可恢复流→EventedAgent 工作流执行→InngestAgent 分布式持久化每层基于上一层扩展。直接工作流执行—DurableAgent通过workflow.createRun()run.start()直接运行EventedAgent重写为不 await 的run.start()实现即发即忘并在到达非挂起终态后deleteRunSnapshots清理快照InngestAgent重写为通过 Inngest 分发。PubSub 承载流式— 闭包无法序列化故用 PubSub 跨进程分发事件CachingPubSub包装任意 PubSub 增加事件缓存与回放。关注点分离— 准备非持久化产序列化状态执行持久化跑在工作流中流式可恢复用缓存 PubSub 事件。工具注册— 工具含不可序列化的execute函数准备阶段抽取元数据id、name、schema序列化真实工具对象存于 per-run 注册表、执行时解析。挂起前持久化消息— 任何工作流挂起工具审批、执行中挂起前消息通过SaveQueueManagerflush 到 memory与非持久化 Agent 行为保持一致。带 TTL 的全局运行注册表—globalRunRegistry使用TTLCache当前为 2 小时滑动 TTL、1000 条上限让工作流步骤可访问非序列化运行状态条目通过dispose回调自动清理。工作流自动注册—getDurableWorkflows()让 Mastra 在添加 durable agent 时自动注册工作流无需手工接线。测试覆盖多执行引擎共享测试套件探索文档记录了packages/core/src/agent/durable/__tests__/下 22 个测试文件覆盖核心操作stream / prepare / resume、工具执行单个/多个/并发/基于工作流、工具审批与执行中挂起、memory 集成、结构化输出、usage 追踪、UI 消息、推理模式、图片输入、模型 fallback、请求上下文传播、scorers、停止条件、可恢复流、缓存 TTL 等。当前仓库中该目录已扩展为数十个测试文件含recover-run.test.ts、recover-active-runs.test.ts、resumable-streams.test.ts、resume-api.test.ts、observe-idle-timeout.test.ts、durable-agent-tool-approval.test.ts、parity.test.ts等并新增了untilIdle、后台任务、跨进程 abort、委托、运行恢复等场景。值得关注的是 workflows/_test-utils/ 中的共享测试工厂createDurableAgentTestSuite同一套测试可分别跑在EventEmitterPubSub快速测试与 Inngest集成测试通过eventPropagationDelay模拟事件传播延迟上从侧面印证了同一套持久化语义、可插拔执行引擎的设计目标。适用前提与边界三层模式按需选择单实例、仅需重连体验选createDurableAgent需要长任务脱离 HTTP 请求执行选createEventedAgent生产分布式、需要跨进程 checkpoint 与重试选createInngestAgent需自行接入 Inngest 服务。跨进程 observe 必须有共享缓存后端如 Redis纯InMemoryServerCache只保证单进程内的回放。闭包形态的配置stopWhen函数、prepareStep、transformToolPayload、委托回调等驻留在进程内运行注册表Inngest worker 重启后跨进程恢复会退化为对应默认行为需要跨进程语义的配置请使用可序列化的数据形态。服务端observe端点要求认证requiresAuth: true并内置 5 分钟 idle 超时兜底防止生产者崩溃导致订阅泄漏。【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastra创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考