DurableAgent
DurableAgent 以持久執行及可恢復串流包裝現有的 Agent。它會執行 Agent 迴圈,讓客戶端即使中斷連線再重新連線,也不會錯過事件,並透過 PubSub 串流傳送這些事件。當一次執行必須超越單一請求的生命週期,或需要在連線中斷後繼續時,便應使用它。
你可使用 createDurableAgent factory 建立它;如要在內置 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 },
})
串流傳送回應並讀取結果。完成該次執行後,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 的直接連結
以持久執行及可恢復串流包裝 Agent。這是建立 DurableAgent 的建議方式。
import { createDurableAgent } from '@mastra/core/agent/durable'
const durableAgent = createDurableAgent({ agent })
傳回:DurableAgent
參數參數 的直接連結
agent:
id?:
name?:
cache?:
false 可停用快取,令串流無法恢復。pubsub?:
maxSteps?:
createEventedAgent(options)createeventedagentoptions 的直接連結
在內置 Workflow 引擎上,以「發出後不理」式持久執行包裝 Agent。它與 createDurableAgent 一樣會傳回可供串流讀取的結果,但底層 Workflow 會以非阻塞方式(透過 startAsync)執行,而非在接通串流前一直執行至完成。當你希望該次執行可獨立於呼叫者繼續進行時,便應使用它。它不接受 id 或 name 覆寫。
import { createEventedAgent } from '@mastra/core/agent/durable'
const eventedAgent = createEventedAgent({ agent })
傳回:EventedAgent(DurableAgent 的子類別)
參數參數 的直接連結
agent:
cache?:
false 可停用快取。pubsub?:
maxSteps?:
建構函式參數建構函式參數 的直接連結
DurableAgent 類別接受與 createDurableAgent 相同的選項,另加 cleanupTimeoutMs。除非需要建立子類別,否則建議使用 factory。
agent:
id?:
name?:
cache?:
false 可停用快取。pubsub?:
maxSteps?:
cleanupTimeoutMs?:
cleanup()。自動清理不會在暫停事件上觸發。方法方法 的直接連結
執行執行 的直接連結
stream(messages, options?)streammessages-options 的直接連結
使用持久執行串流傳送回應。此方法會立即傳回結果,而結果的 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 獲批後。請傳入原始串流的 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() 會無限期等候事件。如果執行該次運行的程序意外停止,該次執行會停止產生事件,卻不會發出完成事件,因此被觀察的串流會永遠等候。傳入 idleTimeoutMs 可限制等候時間:靜默達指定毫秒數後,串流便會結束。系統會先查詢選用的 isAlive 檢查。當該次執行仍在處理中(例如正在進行長時間的 Tool 呼叫,或已暫停以等候人手輸入)時,傳回 true 可繼續等候。傳回 false 或省略 isAlive,串流便會以錯誤結束。isAlive 暫時擲回錯誤時,系統會視為「仍然存活」,因此短暫的檢查失敗不會結束仍在運作的串流。
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 狀態的執行,並由上次保存的快照重新驅動。最多復原 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() 相同的可串流結果。當你需要即時觀察復原串流時,便應使用此方法。
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() 接受 DurableAgentStreamOptions 物件。它支援以下 Agent 執行選項及生命週期 callback。
runId?:
resume() 或 observe() 使用。instructions?:
context?:
memory?:
requestContext?:
maxSteps?:
toolsets?:
clientTools?:
toolChoice?:
activeTools?:
modelSettings?:
Authorization、X-Api-Key 及同類 header)。stopWhen?:
maxSteps。system?:
requireToolApproval?:
true 或 false 以管制全部或完全不管制,亦可傳入函式以設定逐次呼叫政策。函式形式的政策存放於程序內的執行 registry;跨程序恢復時會退回使用值為 true 的影子設定。autoResumeSuspendedTools?:
resume() 呼叫。toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
structuredOutput?:
untilIdle?:
true 可使用預設 5 分鐘閒置逾時,亦可傳入 { maxIdleMs } 自訂。等同已棄用的 streamUntilIdle() 方法。resume() 亦支援此選項。disableBackgroundTasks?:
tracingOptions?:
requestContextKeys。完全可序列化為 JSON。actor?:
transform?:
transformToolPayload closure 存放於程序內的執行 registry;只會序列化可安全用於 JSON 的 targets 影子設定。prepareStep?:
PrepareStepProcessor 叫用的逐步準備 hook。只限 closure,並儲存於程序內的執行 registry。跨程序恢復會失去此 hook。isTaskComplete?:
onComplete 存放於程序內的執行 registry;可安全用於 JSON 的 primitive(strategy、timeout、parallel、suppressFeedback、scorerNames)會序列化,以供跨程序觀察。delegation?:
onDelegationStart、onDelegationComplete、messageFilter)。準備時,callback 會嵌入子 Agent 的 Tool wrapper。跨程序恢復會失去這些 callback。versions?:
abortSignal?:
AbortController,因此任何一方都可取消執行。跨程序恢復無法復原訊號;如需在恢復後保留中止能力,請向 resume() 傳入新訊號。onChunk?:
onStepFinish?:
onFinish?:
onError?:
onSuspended?:
onAbort?:
abortSignal 或 result.abort() 中止執行時呼叫。onIterationComplete?:
messageList、finishReason 及 isFinal 旗標。在持久 Agent 上只供觀察:傳回 continue: false 或意見回饋不會影響迴圈。resume() 及 observe() 接受相同的生命週期 callback(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 callback。執行完成後呼叫亦屬安全;在此情況下不會進行任何操作。