DurableAgent
DurableAgent 使用持久执行和可恢复 stream 封装现有 Agent。它运行 agentic loop,使 client 可以断开并重新连接而不会错过事件,并通过 PubSub 以 streaming 方式传输这些事件。当运行必须比单个请求持续更久,或需要在连接中断后继续运行时,请使用它。
可使用 createDurableAgent 工厂创建实例;若要在内置 Workflow 引擎上触发后即无需等待地执行,请使用 createEventedAgent。对于由 Inngest 驱动的执行,请使用 @mastra/inngest 中的 createInngestAgent。
使用示例使用示例的直接链接
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 订阅:
const { output, runId, cleanup } = await durableAgent.stream('Hello!')
const text = await output.text
cleanup()
使用 durable 配置标志using-the-durable-config-flag的直接链接
在 AgentConfig 中设置 durable: true 后,将 Agent 附加到 Mastra 实例时,系统会自动使用 createDurableAgent 封装它。可使用对象传递 cache、pubsub、maxSteps 或 cleanupTimeoutMs 等高级选项。
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)createdurableagentoptions的直接链接
使用持久执行和可恢复 stream 封装 Agent。这是创建 DurableAgent 的推荐方式。
import { createDurableAgent } from '@mastra/core/agent/durable'
const durableAgent = createDurableAgent({ agent })
返回: DurableAgent
参数参数的直接链接
agent:
id?:
name?:
cache?:
false 可禁用 cache,此时 stream 不可恢复。pubsub?:
maxSteps?:
createEventedAgent(options)createeventedagentoptions的直接链接
使用内置 Workflow 引擎上触发后即无需等待的持久执行来封装 Agent。与 createDurableAgent 一样,它会返回可供 streaming 的结果;但底层 Workflow 会以非阻塞方式(通过 startAsync)运行,而不是在连接 stream 之前运行至完成。希望运行独立于调用方推进时,请使用它。该函数不接受 id 或 name 覆盖。
import { createEventedAgent } from '@mastra/core/agent/durable'
const eventedAgent = createEventedAgent({ agent })
返回: EventedAgent (a subclass of DurableAgent)
参数参数的直接链接
agent:
cache?:
false 可禁用 cache。pubsub?:
maxSteps?:
构造函数参数构造函数参数的直接链接
DurableAgent 类接受与 createDurableAgent 相同的选项,另加 cleanupTimeoutMs。除非需要创建子类,否则优先使用工厂函数。
agent:
id?:
name?:
cache?:
false 可禁用 cache。pubsub?:
maxSteps?:
cleanupTimeoutMs?:
cleanup()。暂停事件不会触发自动清理。方法方法的直接链接
执行执行的直接链接
stream(messages, options?)streammessages-options的直接链接
使用持久执行以 streaming 方式传输响应。该方法立即返回结果,其 output 会随着运行推进而生成事件。
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?)resumerunid-resumedata-options的直接链接
恢复已暂停的运行,例如在 Tool 审批后恢复。传入原始 stream 的 runId 以及运行所等待的数据。如果该运行不存在 registry 条目,则抛出错误。
const { output, cleanup } = await durableAgent.resume(runId, {
approved: true,
})
await output.text
cleanup()
返回: Promise<DurableAgentStreamResult>
observe(runId, options?)observerunid-options的直接链接
重新连接到现有运行,在传递实时事件之前重放已缓存的事件。网络连接中断后可使用此方法。传入 offset 可从已知位置开始重放。
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。
const { output } = await durableAgent.observe(runId, {
idleTimeoutMs: 30_000,
isAlive: () => runHeartbeat.isFresh(runId),
})
因空闲超时结束运行时,会执行与运行出错时相同的清理(请参阅下方警告),因此其缓存状态会被释放而不是保留。这两个选项都需要显式启用。省略它们即可保留此前的无限期等待行为。
返回: Promise<DurableAgentStreamResult>
observe() 返回的 cleanup() 会销毁运行的 registry 条目和缓存事件。仅在不再需要该运行时调用它。如果运行已暂停且打算稍后恢复,请勿调用 cleanup()。让自动清理计时器在运行完成或出错后进行处理。暂停事件不会触发自动清理。
prepare(messages, options?)preparemessages-options的直接链接
为运行准备持久执行,但不启动运行。该方法会在内部 registry 中注册运行,并返回序列化的 Workflow 输入。需要控制 Workflow 的触发时间和方式时,请使用此方法。
const { runId, messageId, workflowInput, threadId, resourceId } = await durableAgent.prepare(
'Summarize the document',
{
memory: { threadId: 'thread-1', resourceId: 'user-1' },
},
)
返回:
interface PrepareResult {
runId: string
messageId: string
workflowInput: any
registryEntry: object
threadId?: string
resourceId?: string
}
恢复恢复的直接链接
recoverActiveRuns(options?)recoveractiverunsoptions的直接链接
查找此 Agent 中停滞在 running 状态的运行,并从上次持久化的 snapshot 重新驱动它们。最多恢复 options.limit 个运行(默认:100),并返回恢复结果摘要。
const result = await durableAgent.recoverActiveRuns()
// { recovered: [{ runId, status }], succeeded: 2, failed: 0 }
传入 runId 可恢复单个已知运行:
await durableAgent.recoverActiveRuns({ runId: 'run-abc-123' })
返回:
interface DurableAgentRecoverActiveRunsResult {
recovered: Array<{ runId: string; status: 'success' | 'failed'; error?: Error }>
succeeded: number
failed: number
}
options.runId?:
options.limit?:
options.createdBefore?:
recover(runId, options?)recoverrunid-options的直接链接
按 ID 恢复单个运行。返回与 stream() 结构相同、可供 streaming 的结果。需要实时观察恢复 stream 时,请使用此方法。
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 选项的直接链接
stream() 接受 DurableAgentStreamOptions 对象。它支持以下 Agent 执行选项以及生命周期回调。
runId?:
resume() 或 observe()。instructions?:
context?:
memory?:
requestContext?:
maxSteps?:
toolsets?:
clientTools?:
toolChoice?:
activeTools?:
modelSettings?:
Authorization、X-Api-Key 等)。stopWhen?:
maxSteps。system?:
requireToolApproval?:
true 或 false 可审批全部或不审批任何调用,也可传入函数以设置逐次调用策略。函数形式的策略保存在进程内运行 registry 中;跨进程恢复时会回退到值为 true 的 shadow。autoResumeSuspendedTools?:
resume()。toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
structuredOutput?:
untilIdle?:
true 使用默认的 5 分钟空闲超时,或传入 { maxIdleMs } 进行自定义。等同于已弃用的 streamUntilIdle() 方法。resume() 也支持此选项。disableBackgroundTasks?:
tracingOptions?:
requestContextKeys。完全支持 JSON 序列化。actor?:
transform?:
transformToolPayload closure 保存在进程内运行 registry 中;仅序列化 JSON 安全的 targets shadow。prepareStep?:
PrepareStepProcessor 调用的逐步骤准备 hook。仅包含 closure,存储在进程内运行 registry 中。跨进程恢复会丢失此 hook。isTaskComplete?:
onComplete 保存在进程内运行 registry 中;JSON 安全的 primitive(strategy、timeout、parallel、suppressFeedback、scorerNames)会被序列化,以支持跨进程可观测性。delegation?:
onDelegationStart、onDelegationComplete、messageFilter)。准备阶段会将回调写入子 Agent Tool wrapper。跨进程恢复会丢失这些回调。versions?:
abortSignal?:
AbortController,因此任一来源都可取消运行。跨进程恢复无法恢复此 signal;如果恢复后仍需支持中止,请向 resume() 传入新的 signal。onChunk?:
onStepFinish?:
onFinish?:
onError?:
onSuspended?:
onAbort?:
abortSignal 或 result.abort() 中止运行时调用。onIterationComplete?:
messageList、finishReason 和 isFinal 标志。对于持久 Agent,此回调仅用于观察:返回 continue: false 或 feedback 不会影响循环。resume() 和 observe() 接受相同的生命周期回调(onChunk、onStepFinish、onFinish、onError、onSuspended)。observe() 还接受 offset,用于控制重放的起始位置。
DurableAgentStreamResultDurableAgentStreamResult的直接链接
stream()、resume()、observe() 和 recover() 返回的对象。
interface DurableAgentStreamResult<OUTPUT = undefined> {
output: MastraModelOutput<OUTPUT>
readonly fullStream: ReadableStream<any>
runId: string
threadId?: string
resourceId?: string
cleanup: () => void
abort: () => void
}
output:
output.text 可获取完整文本,也可消费 output.fullStream。fullStream:
output.fullStream。runId:
resume() 或 observe() 可重新连接。threadId?:
resourceId?:
cleanup:
abort:
AbortController 中止运行。在持久 LLM 执行步骤中表现为 AbortError,并触发 onAbort 回调。运行完成后调用也是安全的,此时不会执行任何操作。