跳到主要内容

DurableAgent

DurableAgent 使用持久执行和可恢复 stream 封装现有 Agent。它运行 agentic loop,使 client 可以断开并重新连接而不会错过事件,并通过 PubSub 以 streaming 方式传输这些事件。当运行必须比单个请求持续更久,或需要在连接中断后继续运行时,请使用它。

可使用 createDurableAgent 工厂创建实例;若要在内置 Workflow 引擎上触发后即无需等待地执行,请使用 createEventedAgent。对于由 Inngest 驱动的执行,请使用 @mastra/inngest 中的 createInngestAgent

使用示例
使用示例的直接链接

src/mastra/index.ts
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 封装它。可使用对象传递 cachepubsubmaxStepscleanupTimeoutMs 等高级选项。

src/mastra/index.ts
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:

Agent
要使用持久执行能力封装的 Agent。Agent 方法会委托给此 Agent。

id?:

string
= agent.id
覆盖 ID。

name?:

string
= agent.name
覆盖名称。

cache?:

MastraServerCache | false
用于存储 stream 事件的 cache,从而支持可恢复 stream。如果省略,Agent 会从 Mastra 实例继承 cache,或使用 InMemoryServerCache。设为 false 可禁用 cache,此时 stream 不可恢复。

pubsub?:

PubSub
= EventEmitterPubSub
用于以 streaming 方式传输事件的 PubSub 实例。

maxSteps?:

number
agentic loop 的最大步骤数。

createEventedAgent(options)
createeventedagentoptions的直接链接

使用内置 Workflow 引擎上触发后即无需等待的持久执行来封装 Agent。与 createDurableAgent 一样,它会返回可供 streaming 的结果;但底层 Workflow 会以非阻塞方式(通过 startAsync)运行,而不是在连接 stream 之前运行至完成。希望运行独立于调用方推进时,请使用它。该函数不接受 idname 覆盖。

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
= EventEmitterPubSub
用于以 streaming 方式传输事件的 PubSub 实例。

maxSteps?:

number
agentic loop 的最大步骤数。

构造函数参数
构造函数参数的直接链接

DurableAgent 类接受与 createDurableAgent 相同的选项,另加 cleanupTimeoutMs。除非需要创建子类,否则优先使用工厂函数。

agent:

Agent
要使用持久执行能力封装的 Agent。

id?:

string
= agent.id
覆盖 ID。

name?:

string
= agent.name
覆盖名称。

cache?:

MastraServerCache | false
用于存储 stream 事件的 cache。如果省略,则从 Mastra 实例继承或使用 InMemoryServerCache。设为 false 可禁用 cache。

pubsub?:

PubSub
= EventEmitterPubSub
用于以 streaming 方式传输事件的 PubSub 实例。

maxSteps?:

number
agentic loop 的最大步骤数。

cleanupTimeoutMs?:

number
= 30000
stream 完成或出错后,自动清理 registry 条目前等待的宽限时间(毫秒)。设为 0 可禁用自动清理,并要求手动调用 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?:

string
按 ID 恢复指定运行。设置后,将忽略查找过滤条件。

options.limit?:

number
要查找的活跃运行最大数量。默认为 100。

options.createdBefore?:

Date
仅恢复在此日期之前创建的运行。

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?:

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(AuthorizationX-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<boolean>)
要求审批 Tool 调用。传入 truefalse 可审批全部或不审批任何调用,也可传入函数以设置逐次调用策略。函数形式的策略保存在进程内运行 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(strategytimeoutparallelsuppressFeedbackscorerNames)会被序列化,以支持跨进程可观测性。

delegation?:

AgentExecutionOptions['delegation']
子 Agent 委托 hook(onDelegationStartonDelegationCompletemessageFilter)。准备阶段会将回调写入子 Agent Tool wrapper。跨进程恢复会丢失这些回调。

versions?:

object
子 Agent 委托的版本覆盖。

abortSignal?:

AbortSignal
外部 abort signal。该 signal 会转发到持久运行的内部 AbortController,因此任一来源都可取消运行。跨进程恢复无法恢复此 signal;如果恢复后仍需支持中止,请向 resume() 传入新的 signal。

onChunk?:

(chunk: ChunkType) => void | Promise<void>
每个 streaming chunk 都会调用。

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
agentic loop 中的步骤完成时调用。

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
运行完成时调用。

onError?:

(error: Error) => void | Promise<void>
运行出错时调用。

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
运行暂停时调用,例如等待 Tool 审批时。

onAbort?:

AgentExecutionOptions['onAbort']
通过 abortSignalresult.abort() 中止运行时调用。

onIterationComplete?:

AgentExecutionOptions['onIterationComplete']
每次 agentic loop 迭代后调用,并提供最新的 messageListfinishReasonisFinal 标志。对于持久 Agent,此回调仅用于观察:返回 continue: false 或 feedback 不会影响循环。

resume()observe() 接受相同的生命周期回调(onChunkonStepFinishonFinishonErroronSuspended)。observe() 还接受 offset,用于控制重放的起始位置。

DurableAgentStreamResult
DurableAgentStreamResult的直接链接

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:

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 回调。运行完成后调用也是安全的,此时不会执行任何操作。