DurableAgent
DurableAgent 會以持久執行與可恢復 stream 包裝現有的 Agent。它會執行 Agent 迴圈,讓 client 中斷連線後仍可重新連線而不漏掉 event,並透過 PubSub 串流這些 event。若 run 必須比單一 request 存續更久,或要在連線中斷後繼續,請使用此類別。
使用 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 },
})
串流回應並讀取結果。run 使用完畢後,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 可停用快取,stream 將無法恢復。pubsub?:
maxSteps?:
createEventedAgent(options)「createeventedagentoptions」的直接連結
在內建 Workflow 引擎上,以傳送後不等待結果的持久執行方式包裝 Agent。和 createDurableAgent 一樣,它會回傳可供串流的結果;但底層 Workflow 會以非阻塞方式(透過 startAsync)執行,而不是先執行至完成再連接 stream。若希望 run 獨立於呼叫端繼續進行,請使用此函式。它不接受 id 或 name 覆寫值。
import { createEventedAgent } from '@mastra/core/agent/durable'
const eventedAgent = createEventedAgent({ agent })
回傳:EventedAgent(DurableAgent 的 subclass)
參數「參數」的直接連結
agent:
cache?:
false 可停用快取。pubsub?:
maxSteps?:
constructor 參數「constructor 參數」的直接連結
DurableAgent 類別接受與 createDurableAgent 相同的選項,另加 cleanupTimeoutMs。除非需要建立 subclass,否則建議使用 factory。
agent:
id?:
name?:
cache?:
false 可停用快取。pubsub?:
maxSteps?:
cleanupTimeoutMs?:
cleanup()。自動清理不會在暫停 event 時觸發。方法「方法」的直接連結
執行「執行」的直接連結
stream(messages, options?)「streammessages-options」的直接連結
使用持久執行串流回應。立即回傳結果;隨 run 進行,結果的 output 會產生 event。
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」的直接連結
恢復已暫停的 run,例如 Tool 核准後。請傳入原始 stream 的 runId,以及 run 等待的資料。若 registry 中沒有該 run 的項目,便會擲回錯誤。
const { output, cleanup } = await durableAgent.resume(runId, {
approved: true,
})
await output.text
cleanup()
回傳: Promise<DurableAgentStreamResult>
observe(runId, options?)「observerunid-options」的直接連結
重新連線至現有 run;先重播快取 event,再傳送即時 event。網路中斷後請使用此方法。傳入 offset 可從已知位置開始重播。
const { output, cleanup } = await durableAgent.observe(runId, {
offset: 0,
onChunk: chunk => console.log(chunk),
})
await output.text
observe() 預設會無限期等待 event。若執行 run 的處理程序意外停止,run 會停止產生 event,但絕不會發出完成 event,因此受觀察的 stream 會永遠等待。傳入 idleTimeoutMs 可限制等待時間:靜默達到指定毫秒數後,stream 便會結束。系統會先執行選用的 isAlive 檢查。run 仍在處理中時(例如長時間執行的 Tool 呼叫,或 run 暫停並等待人工輸入),請回傳 true 以繼續等待。回傳 false 或省略 isAlive,stream 會以錯誤結束。isAlive 暫時擲回錯誤會視為「仍在運作」,因此短暫的檢查失敗不會結束即時 stream。
const { output } = await durableAgent.observe(runId, {
idleTimeoutMs: 30_000,
isAlive: () => runHeartbeat.isFresh(runId),
})
因閒置逾時而結束 run 時,會執行與 run 發生錯誤時相同的清理作業(請參閱下方警告),因此會釋放快取狀態,而非保留。兩個選項皆需明確啟用。省略即可使用原本的無限期等待行為。
回傳: Promise<DurableAgentStreamResult>
observe() 回傳的 cleanup() 會刪除 run 的 registry 項目與快取 event。只有在 run 使用完畢後才能呼叫。若 run 已暫停且您打算稍後恢復,請勿呼叫 cleanup();請讓自動清理計時器在 run 完成或發生錯誤後處理。自動清理不會在暫停 event 時觸發。
prepare(messages, options?)「preparemessages-options」的直接連結
準備 run 以供持久執行,但不會啟動。此方法會在內部 registry 註冊 run,並回傳序列化的 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 狀態的 run,並從最後一個持久化快照重新驅動。最多復原 options.limit 個 run(預設:100),並回傳復原摘要。
const result = await durableAgent.recoverActiveRuns()
// { recovered: [{ runId, status }], succeeded: 2, failed: 0 }
傳入 runId 可復原單一已知 run:
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 復原單一 run。回傳結構與 stream() 相同的可串流結果。需要即時觀察復原 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 執行選項,以及生命週期 callback。
runId?:
resume() 或 observe() 使用。instructions?:
context?:
memory?:
requestContext?:
maxSteps?:
toolsets?:
clientTools?:
toolChoice?:
activeTools?:
modelSettings?:
Authorization、X-Api-Key 等)。stopWhen?:
maxSteps。system?:
requireToolApproval?:
true 或 false 可控管全部或完全不控管;也可傳入函式,為每次呼叫設定原則。函式形式的原則存放在處理程序內 run registry;跨處理程序恢復時會退回使用 true shadow。autoResumeSuspendedTools?:
resume() 呼叫。toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
structuredOutput?:
untilIdle?:
true 可使用預設 5 分鐘閒置逾時,或傳入 { maxIdleMs } 自訂。等同已棄用的 streamUntilIdle() 方法。resume() 也支援此選項。disableBackgroundTasks?:
tracingOptions?:
requestContextKeys。可完整序列化為 JSON。actor?:
transform?:
transformToolPayload closure 存放在處理程序內 run registry;只有 JSON-safe 的 targets shadow 會序列化。prepareStep?:
PrepareStepProcessor 叫用的步驟準備 hook。僅限 closure,並儲存在處理程序內 run registry。跨處理程序恢復時會遺失此 hook。isTaskComplete?:
onComplete 存放在處理程序內 run registry;JSON-safe primitive(strategy、timeout、parallel、suppressFeedback、scorerNames)會序列化,以供跨處理程序可觀測性使用。delegation?:
onDelegationStart、onDelegationComplete、messageFilter)。準備時會將 callback 寫入子 Agent Tool wrapper。跨處理程序恢復時會遺失 callback。versions?:
abortSignal?:
AbortController,因此任一來源都能取消 run。跨處理程序恢復無法復原 signal;若需要在恢復後中止,請將新的 signal 傳給 resume()。onChunk?:
onStepFinish?:
onFinish?:
onError?:
onSuspended?:
onAbort?:
abortSignal 或 result.abort() 中止時呼叫。onIterationComplete?:
messageList、finishReason 與 isFinal 旗標。在持久 Agent 上僅供觀察:回傳 continue: false 或 feedback 不會影響迴圈。resume() 與 observe() 接受相同的生命週期 callback(onChunk、onStepFinish、onFinish、onError、onSuspended)。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:
output.text 可取得完整文字,或取用 output.fullStream。fullStream:
output.fullStream。runId:
resume() 或 observe() 可重新連線。threadId?:
resourceId?:
cleanup:
abort:
AbortController 中止 run。會在持久 LLM 執行步驟內顯示為 AbortError,並觸發 onAbort callback。run 完成後也可安全呼叫;此時不會執行任何作業。