> Discover all available pages from the documentation index: https://mastra.zisheng.pro/llms.txt # DurableAgent `DurableAgent` 使用持久执行和可恢复 stream 封装现有 [`Agent`](https://mastra.zisheng.pro/reference/agents/agent)。它运行 agentic loop,使 client 可以断开并重新连接而不会错过事件,并通过 [PubSub](https://mastra.zisheng.pro/docs/server/pubsub) 以 streaming 方式传输这些事件。当运行必须比单个请求持续更久,或需要在连接中断后继续运行时,请使用它。 可使用 [`createDurableAgent`](#createdurableagentoptions) 工厂创建实例;若要在内置 Workflow 引擎上触发后即无需等待地执行,请使用 [`createEventedAgent`](#createeventedagentoptions)。对于由 Inngest 驱动的执行,请使用 `@mastra/inngest` 中的 [`createInngestAgent`](https://mastra.zisheng.pro/reference/agents/inngest-agent)。 ## 使用示例 ```typescript import { Mastra } from '@mastra/core' import { Agent } from '@mastra/core/agent' import { createDurableAgent } from '@mastra/core/agent/durable' const agent = new Agent({ id: 'my-agent', name: 'My Agent', instructions: 'You are a helpful assistant', model: 'openai/gpt-5.6-sol', }) const durableAgent = createDurableAgent({ agent }) export const mastra = new Mastra({ agents: { myAgent: durableAgent }, }) ``` 以 streaming 方式传输响应并读取结果。运行使用完毕后,`cleanup` 函数会取消 PubSub 订阅: ```typescript const { output, runId, cleanup } = await durableAgent.stream('Hello!') const text = await output.text cleanup() ``` ### 使用 `durable` 配置标志 在 `AgentConfig` 中设置 `durable: true` 后,将 Agent 附加到 `Mastra` 实例时,系统会自动使用 `createDurableAgent` 封装它。可使用对象传递 `cache`、`pubsub`、`maxSteps` 或 `cleanupTimeoutMs` 等高级选项。 ```typescript import { Mastra } from '@mastra/core' import { Agent } from '@mastra/core/agent' const myAgent = new Agent({ id: 'my-agent', name: 'My Agent', instructions: 'You are a helpful assistant', model: 'openai/gpt-5.6-sol', durable: true, // or: { maxSteps: 10, cleanupTimeoutMs: 60_000 } }) export const mastra = new Mastra({ agents: { myAgent }, }) ``` `mastra.getAgent('myAgent')` 返回封装后的 `DurableAgent`。独立 Agent(已构造但从未注册到 `Mastra` 实例)不会变为持久 Agent。封装会在注册时应用。 ## `createDurableAgent(options)` 使用持久执行和可恢复 stream 封装 `Agent`。这是创建 `DurableAgent` 的推荐方式。 ```typescript import { createDurableAgent } from '@mastra/core/agent/durable' const durableAgent = createDurableAgent({ agent }) ``` 返回: `DurableAgent` ### 参数 **agent** (`Agent`): 要使用持久执行能力封装的 Agent。Agent 方法会委托给此 Agent。 **id** (`string`): 覆盖 ID。 (Default: `agent.id`) **name** (`string`): 覆盖名称。 (Default: `agent.name`) **cache** (`MastraServerCache | false`): 用于存储 stream 事件的 cache,从而支持可恢复 stream。如果省略,Agent 会从 Mastra 实例继承 cache,或使用 InMemoryServerCache。设为 false 可禁用 cache,此时 stream 不可恢复。 **pubsub** (`PubSub`): 用于以 streaming 方式传输事件的 PubSub 实例。 (Default: `EventEmitterPubSub`) **maxSteps** (`number`): agentic loop 的最大步骤数。 ## `createEventedAgent(options)` 使用内置 Workflow 引擎上触发后即无需等待的持久执行来封装 `Agent`。与 `createDurableAgent` 一样,它会返回可供 streaming 的结果;但底层 Workflow 会以非阻塞方式(通过 `startAsync`)运行,而不是在连接 stream 之前运行至完成。希望运行独立于调用方推进时,请使用它。该函数不接受 `id` 或 `name` 覆盖。 ```typescript import { createEventedAgent } from '@mastra/core/agent/durable' const eventedAgent = createEventedAgent({ agent }) ``` 返回: `EventedAgent` (a subclass of `DurableAgent`) ### 参数 **agent** (`Agent`): 要使用事件驱动持久执行能力封装的 Agent。 **cache** (`MastraServerCache | false`): 用于存储 stream 事件的 cache,从而支持可恢复 stream。如果省略,Agent 会从 Mastra 实例继承 cache,或使用 InMemoryServerCache。设为 false 可禁用 cache。 **pubsub** (`PubSub`): 用于以 streaming 方式传输事件的 PubSub 实例。 (Default: `EventEmitterPubSub`) **maxSteps** (`number`): agentic loop 的最大步骤数。 ## 构造函数参数 `DurableAgent` 类接受与 `createDurableAgent` 相同的选项,另加 `cleanupTimeoutMs`。除非需要创建子类,否则优先使用工厂函数。 **agent** (`Agent`): 要使用持久执行能力封装的 Agent。 **id** (`string`): 覆盖 ID。 (Default: `agent.id`) **name** (`string`): 覆盖名称。 (Default: `agent.name`) **cache** (`MastraServerCache | false`): 用于存储 stream 事件的 cache。如果省略,则从 Mastra 实例继承或使用 InMemoryServerCache。设为 false 可禁用 cache。 **pubsub** (`PubSub`): 用于以 streaming 方式传输事件的 PubSub 实例。 (Default: `EventEmitterPubSub`) **maxSteps** (`number`): agentic loop 的最大步骤数。 **cleanupTimeoutMs** (`number`): stream 完成或出错后,自动清理 registry 条目前等待的宽限时间(毫秒)。设为 0 可禁用自动清理,并要求手动调用 cleanup()。暂停事件不会触发自动清理。 (Default: `30000`) ## 方法 ### 执行 #### `stream(messages, options?)` 使用持久执行以 streaming 方式传输响应。该方法立即返回结果,其 `output` 会随着运行推进而生成事件。 ```typescript const { output, runId, cleanup } = await durableAgent.stream('Hello!', { onChunk: chunk => console.log(chunk), onFinish: result => console.log('done', result), }) const text = await output.text cleanup() ``` 返回: [`Promise`](#durableagentstreamresult) #### `resume(runId, resumeData, options?)` 恢复已暂停的运行,例如在 Tool 审批后恢复。传入原始 stream 的 `runId` 以及运行所等待的数据。如果该运行不存在 registry 条目,则抛出错误。 ```typescript const { output, cleanup } = await durableAgent.resume(runId, { approved: true, }) await output.text cleanup() ``` 返回: [`Promise`](#durableagentstreamresult) #### `observe(runId, options?)` 重新连接到现有运行,在传递实时事件之前重放已缓存的事件。网络连接中断后可使用此方法。传入 `offset` 可从已知位置开始重放。 ```typescript const { output, cleanup } = await durableAgent.observe(runId, { offset: 0, onChunk: chunk => console.log(chunk), }) await output.text ``` 默认情况下,`observe()` 会无限期等待事件。如果执行运行的进程意外停止,运行会停止生成事件但永远不会发出完成事件,因此被观察的 stream 会一直等待。传入 `idleTimeoutMs` 可限制等待时间:静默达到指定毫秒数后,stream 将结束。系统会先调用可选的 `isAlive` 检查。当运行仍在处理中(例如正在执行长时间运行的 Tool 调用,或暂停等待人工输入)时返回 `true`,即可继续等待。返回 `false` 或省略 `isAlive` 会使 stream 以错误结束。`isAlive` 暂时抛出异常会被视为“仍然存活”,因此短暂的检查失败不会结束活跃 stream。 ```typescript const { output } = await durableAgent.observe(runId, { idleTimeoutMs: 30_000, isAlive: () => runHeartbeat.isFresh(runId), }) ``` 因空闲超时结束运行时,会执行与运行出错时相同的清理(请参阅下方警告),因此其缓存状态会被释放而不是保留。这两个选项都需要显式启用。省略它们即可保留此前的无限期等待行为。 返回: `Promise` > **注意:** `observe()` 返回的 `cleanup()` 会销毁运行的 registry 条目和缓存事件。仅在不再需要该运行时调用它。如果运行已暂停且打算稍后恢复,请勿调用 `cleanup()`。让自动清理计时器在运行完成或出错后进行处理。暂停事件不会触发自动清理。 #### `prepare(messages, options?)` 为运行准备持久执行,但不启动运行。该方法会在内部 registry 中注册运行,并返回序列化的 Workflow 输入。需要控制 Workflow 的触发时间和方式时,请使用此方法。 ```typescript const { runId, messageId, workflowInput, threadId, resourceId } = await durableAgent.prepare( 'Summarize the document', { memory: { threadId: 'thread-1', resourceId: 'user-1' }, }, ) ``` 返回: ```typescript interface PrepareResult { runId: string messageId: string workflowInput: any registryEntry: object threadId?: string resourceId?: string } ``` ### 恢复 #### `recoverActiveRuns(options?)` 查找此 Agent 中停滞在 `running` 状态的运行,并从上次持久化的 snapshot 重新驱动它们。最多恢复 `options.limit` 个运行(默认:100),并返回恢复结果摘要。 ```typescript const result = await durableAgent.recoverActiveRuns() // { recovered: [{ runId, status }], succeeded: 2, failed: 0 } ``` 传入 `runId` 可恢复单个已知运行: ```typescript await durableAgent.recoverActiveRuns({ runId: 'run-abc-123' }) ``` 返回: ```typescript interface DurableAgentRecoverActiveRunsResult { recovered: Array<{ runId: string; status: 'success' | 'failed'; error?: Error }> succeeded: number failed: number } ``` **options.runId** (`string`): 按 ID 恢复指定运行。设置后,将忽略查找过滤条件。 **options.limit** (`number`): 要查找的活跃运行最大数量。默认为 100。 **options.createdBefore** (`Date`): 仅恢复在此日期之前创建的运行。 #### `recover(runId, options?)` 按 ID 恢复单个运行。返回与 `stream()` 结构相同、可供 streaming 的结果。需要实时观察恢复 stream 时,请使用此方法。 ```typescript const { output, cleanup } = await durableAgent.recover('run-abc-123', { onChunk: chunk => console.log(chunk), onError: ({ error }) => console.error(error), }) await output.text cleanup() ``` 返回: [`Promise`](#durableagentstreamresult) ## Stream 选项 `stream()` 接受 `DurableAgentStreamOptions` 对象。它支持以下 Agent 执行选项以及生命周期回调。 **runId** (`string`): 此运行的唯一标识符。之后可将其用于 resume() 或 observe()。 **instructions** (`AgentExecutionOptions['instructions']`): 覆盖 Agent 针对此运行的默认 instructions。接受静态字符串或 Agent 支持的同类动态 instructions 值。 **context** (`ModelMessage[]`): 提供给 Agent 的其他上下文消息。 **memory** (`object`): 用于持久化和检索对话的 Memory 配置。 **requestContext** (`RequestContext`): 携带此运行动态配置和状态的 RequestContext。 **maxSteps** (`number`): 此 stream 最多运行的步骤数。 **toolsets** (`object`): 此运行可用的其他 Tool 集。 **clientTools** (`object`): 执行期间可用的 client 端 Tool。 **toolChoice** (`'auto' | 'none' | 'required' | { type: 'tool'; toolName: string }`): Tool 选择策略。 **activeTools** (`string[]`): 将执行限制为 Agent Tool 中指定名称的子集。 **modelSettings** (`object`): Model 专属设置,例如 temperature。序列化 snapshot 跨越进程边界之前,会移除包含凭据的 header(Authorization、X-Api-Key 等)。 **stopWhen** (`AgentExecutionOptions['stopWhen']`): 提前结束 agentic loop 的 predicate 或组合条件。closure 保存在进程内运行 registry 中;跨进程恢复时将降级为仅使用 maxSteps。 **system** (`string | string[]`): 追加在 Agent instructions 之后、用户消息之前的其他 system 消息。 **requireToolApproval** (`boolean | ((args: { toolName: string; args: unknown; requestContext: RequestContext; workspace?: string }) => boolean | Promise)`): 要求审批 Tool 调用。传入 true 或 false 可审批全部或不审批任何调用,也可传入函数以设置逐次调用策略。函数形式的策略保存在进程内运行 registry 中;跨进程恢复时会回退到值为 true 的 shadow。 **autoResumeSuspendedTools** (`boolean`): 自动恢复已暂停的 Tool,而不是等待外部调用 resume()。 **toolCallConcurrency** (`number`): 并发执行的 Tool 调用最大数量。 **includeRawChunks** (`boolean`): 在 stream 输出中包含 Provider 的原始 chunk。 **maxProcessorRetries** (`number`): 每次生成中 Processor 的最大重试次数。 **structuredOutput** (`object`): 结构化输出配置。 **untilIdle** (`boolean | { maxIdleMs?: number }`): 设置后,在后台任务继续运行期间保持 stream 开启,直到 Agent 空闲。传入 true 使用默认的 5 分钟空闲超时,或传入 { maxIdleMs } 进行自定义。等同于已弃用的 streamUntilIdle() 方法。resume() 也支持此选项。 **disableBackgroundTasks** (`boolean`): 禁用此运行的后台任务分派。可在后台运行的 Tool 将改为内联执行。 **tracingOptions** (`AgentExecutionOptions['tracingOptions']`): 转发到 Agent 和 Model span 的 tracing metadata、tag、trace ID、parent span ID 及 requestContextKeys。完全支持 JSON 序列化。 **actor** (`AgentExecutionOptions['actor']`): 转发到 FGA 检查和 Tool 执行的逐次调用 actor signal。 **transform** (`AgentExecutionOptions['transform']`): 逐次调用的 Tool payload 转换策略。transformToolPayload closure 保存在进程内运行 registry 中;仅序列化 JSON 安全的 targets shadow。 **prepareStep** (`AgentExecutionOptions['prepareStep']`): 每次迭代开始时作为 PrepareStepProcessor 调用的逐步骤准备 hook。仅包含 closure,存储在进程内运行 registry 中。跨进程恢复会丢失此 hook。 **isTaskComplete** (`AgentExecutionOptions['isTaskComplete']`): 逐次调用的完成策略。Scorer 实例和 onComplete 保存在进程内运行 registry 中;JSON 安全的 primitive(strategy、timeout、parallel、suppressFeedback、scorerNames)会被序列化,以支持跨进程可观测性。 **delegation** (`AgentExecutionOptions['delegation']`): 子 Agent 委托 hook(onDelegationStart、onDelegationComplete、messageFilter)。准备阶段会将回调写入子 Agent Tool wrapper。跨进程恢复会丢失这些回调。 **versions** (`object`): 子 Agent 委托的版本覆盖。 **abortSignal** (`AbortSignal`): 外部 abort signal。该 signal 会转发到持久运行的内部 AbortController,因此任一来源都可取消运行。跨进程恢复无法恢复此 signal;如果恢复后仍需支持中止,请向 resume() 传入新的 signal。 **onChunk** (`(chunk: ChunkType) => void | Promise`): 每个 streaming chunk 都会调用。 **onStepFinish** (`(result: AgentStepFinishEventData) => void | Promise`): agentic loop 中的步骤完成时调用。 **onFinish** (`(result: AgentFinishEventData) => void | Promise`): 运行完成时调用。 **onError** (`(error: Error) => void | Promise`): 运行出错时调用。 **onSuspended** (`(data: AgentSuspendedEventData) => void | Promise`): 运行暂停时调用,例如等待 Tool 审批时。 **onAbort** (`AgentExecutionOptions['onAbort']`): 通过 abortSignal 或 result.abort() 中止运行时调用。 **onIterationComplete** (`AgentExecutionOptions['onIterationComplete']`): 每次 agentic loop 迭代后调用,并提供最新的 messageList、finishReason 和 isFinal 标志。对于持久 Agent,此回调仅用于观察:返回 continue: false 或 feedback 不会影响循环。 `resume()` 和 `observe()` 接受相同的生命周期回调(`onChunk`、`onStepFinish`、`onFinish`、`onError`、`onSuspended`)。`observe()` 还接受 `offset`,用于控制重放的起始位置。 ## DurableAgentStreamResult `stream()`、`resume()`、`observe()` 和 `recover()` 返回的对象。 ```typescript interface DurableAgentStreamResult { output: MastraModelOutput readonly fullStream: ReadableStream runId: string threadId?: string resourceId?: string cleanup: () => void abort: () => void } ``` **output** (`MastraModelOutput`): streaming 输出。等待 output.text 可获取完整文本,也可消费 output.fullStream。 **fullStream** (`ReadableStream`): 完整事件 stream,委托给 output.fullStream。 **runId** (`string`): 唯一运行 ID。将其传给 resume() 或 observe() 可重新连接。 **threadId** (`string`): 使用 Memory 时的 thread ID。 **resourceId** (`string`): 使用 Memory 时的 resource ID。 **cleanup** (`() => void`): 取消 PubSub 订阅并清除运行的 registry 条目。运行使用完毕后调用。 **abort** (`() => void`): 通过触发内部 AbortController 中止运行。在持久 LLM 执行步骤中表现为 AbortError,并触发 onAbort 回调。运行完成后调用也是安全的,此时不会执行任何操作。 ## 相关内容 - [`createInngestAgent()`](https://mastra.zisheng.pro/reference/agents/inngest-agent) - [Agent 类](https://mastra.zisheng.pro/reference/agents/agent) - [PubSub](https://mastra.zisheng.pro/docs/server/pubsub) - [`.getMemory()`](https://mastra.zisheng.pro/reference/agents/getMemory)