DurableAgent
Breaking changes may occur without a major version bump until the API is stable.
DurableAgent wraps an existing Agent with durable execution and resumable streams. It runs the agentic loop so that a client can disconnect and reconnect without missing events, and it streams those events over PubSub. Use it when a run must outlive a single request or survive a dropped connection.
Create one with the createDurableAgent factory, or use createEventedAgent for fire-and-forget execution on the built-in workflow engine. For Inngest-powered execution, use createInngestAgent from @mastra/inngest.
Usage exampleDirect link to Usage example
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 },
})
Stream a response and read the result. The cleanup function unsubscribes from PubSub when you are done with the run:
const { output, runId, cleanup } = await durableAgent.stream('Hello!')
const text = await output.text
cleanup()
Using the durable config flagDirect link to using-the-durable-config-flag
Set durable: true on AgentConfig and the agent is automatically wrapped with createDurableAgent when it's attached to a Mastra instance. Use an object to forward advanced options such as cache, pubsub, maxSteps, cleanupTimeoutMs, shouldCache, or shouldPersistSnapshot.
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') returns the wrapped DurableAgent. Standalone agents (constructed but never registered on a Mastra instance) don't become durable. The wrapping is applied at registration.
createDurableAgent(options)Direct link to createdurableagentoptions
Wraps an Agent with durable execution and resumable streams. This is the recommended way to create a DurableAgent.
import { createDurableAgent } from '@mastra/core/agent/durable'
const durableAgent = createDurableAgent({ agent })
Returns: DurableAgent
ParametersDirect link to Parameters
agent:
id?:
name?:
cache?:
false to disable caching, which makes streams non-resumable.pubsub?:
maxSteps?:
shouldCache?:
false to publish a topic straight to the underlying PubSub without recording it; subscribers of that topic receive live events only and cannot resume from an offset. Useful for trading replay for minimum publish latency on hot topics when the cache is remote (for example, cross-region Redis). Run-local topics are always excluded, regardless of this option.shouldPersistSnapshot?:
pending, paused, and suspended (required for human-in-the-loop resume) and persists running checkpoints only when the Mastra instance is configured with recovery.durableAgents: 'auto'. Pass a predicate that includes running to keep crash-recovery checkpoints for manual listActiveRuns() / recoverActiveRuns() without enabling automatic recovery. Mastra logs a warning when a predicate excludes suspended or paused, or excludes running while recovery.durableAgents is 'auto'.createEventedAgent(options)Direct link to createeventedagentoptions
Wraps an Agent with fire-and-forget durable execution on the built-in evented workflow engine. Like createDurableAgent, it returns a result you stream from, but the run is started without being awaited: execution is driven by events on the Mastra instance's PubSub, so the run progresses independently of the caller. Register the agent on a Mastra instance with storage to get evented execution; without one, the agent falls back to the default in-process engine and logs a warning. It doesn't accept id or name overrides.
import { createEventedAgent } from '@mastra/core/agent/durable'
const eventedAgent = createEventedAgent({ agent })
Returns: EventedAgent (a subclass of DurableAgent)
ParametersDirect link to Parameters
agent:
cache?:
false to disable caching.pubsub?:
maxSteps?:
shouldCache?:
false to publish a topic straight to the underlying PubSub without recording it; subscribers of that topic receive live events only and cannot resume from an offset. Useful for trading replay for minimum publish latency on hot topics when the cache is remote (for example, cross-region Redis). Run-local topics are always excluded, regardless of this option.shouldPersistSnapshot?:
createDurableAgent, but ignored: the evented engine requires the full snapshot set (pending, paused, suspended, running) because the initial running write creates the base row that suspend-merges and multi-worker coordination build on. A warning is logged if set.Constructor parametersDirect link to Constructor parameters
The DurableAgent class accepts the same options as createDurableAgent, plus cleanupTimeoutMs. Prefer the factory unless you need to subclass.
agent:
id?:
name?:
cache?:
false to disable caching.pubsub?:
maxSteps?:
shouldCache?:
false to publish a topic straight to the underlying PubSub without recording it; subscribers of that topic receive live events only and cannot resume from an offset. Useful for trading replay for minimum publish latency on hot topics when the cache is remote (for example, cross-region Redis). Run-local topics are always excluded, regardless of this option.shouldPersistSnapshot?:
pending, paused, and suspended (required for human-in-the-loop resume) and persists running checkpoints only when the Mastra instance is configured with recovery.durableAgents: 'auto'.cleanupTimeoutMs?:
cleanup() call. Auto-cleanup does not fire on suspended events.MethodsDirect link to Methods
ExecutionDirect link to Execution
stream(messages, options?)Direct link to streammessages-options
Streams a response using durable execution. Returns immediately with a result whose output produces events as the run progresses.
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()
Returns: Promise<DurableAgentStreamResult>
resume(runId, resumeData, options?)Direct link to resumerunid-resumedata-options
Resumes a suspended run, for example after a tool approval. Pass the runId from the original stream and the data the run was waiting on. Throws if no registry entry exists for the run.
const { output, cleanup } = await durableAgent.resume(runId, {
approved: true,
})
await output.text
cleanup()
Returns: Promise<DurableAgentStreamResult>
observe(runId, options?)Direct link to observerunid-options
Reconnects to an existing run. Use this after a network disconnection. For the supported general reconnect flow, omit offset to replay all available cached events before receiving live events.
const { output } = await durableAgent.observe(runId, {
onChunk: chunk => console.log(chunk),
})
await output.text
The offset option is an inclusive, zero-based PubSub event index. It counts every cached event published to the run topic, including lifecycle events, so it isn't a chunk count. Public stream chunks and callbacks don't expose the PubSub index. Only pass an offset when another integration already tracks the transport position.
A non-zero offset can skip earlier events. Aggregated values are built only from chunks delivered to that observer. If the offset skips earlier text deltas, output.text is partial and structured output may be incomplete or fail validation or parsing.
When the PubSub transport supports numeric offsets, an offset beyond the highest retained index replays nothing and also skips live events with lower indexes. The observer receives events only when the published index reaches the requested offset. If the index counter restarts below that offset, the observer can wait indefinitely. On a transport without numeric offset support, providing an offset starts with new live events and doesn't replay cached events. Replay is limited to events still available in the configured cache. Events can't be replayed when caching is disabled, shouldCache excludes the run topic, a cache write fails, or retained events have been cleared.
By default observe() waits indefinitely for events. If the process running the run stops unexpectedly, the run stops producing events but never emits a completion event, so the observed stream would wait forever. Pass idleTimeoutMs to end the stream after a bounded period of silence. Before timing out, observe() can consult an optional isAlive check. Return true while work continues, such as during a long-running tool call or a pause for human input. A false result ends the stream with an error, as does omitting isAlive. A transient throw from isAlive is treated as "still alive", so a momentary check failure never ends a live stream.
const { output } = await durableAgent.observe(runId, {
idleTimeoutMs: 30_000,
isAlive: () => runHeartbeat.isFresh(runId),
})
Ending a run on idle timeout runs the same cleanup as a run that errors (see the warning below), so its cached state is released rather than retained. Both options are opt-in. Omit them for the previous wait-indefinitely behavior.
Returns: Promise<DurableAgentStreamResult>
The cleanup() returned by observe() destroys the run's registry entries and cached events. Only call it when you are done with the run. If the run is suspended and you intend to resume later, don't call cleanup(). Let the auto-cleanup timer handle it after the run finishes or errors. Auto-cleanup doesn't fire on suspended events.
prepare(messages, options?)Direct link to preparemessages-options
Prepares a run for durable execution without starting it. Registers the run in the internal registry and returns the serialized workflow input. Use this when you need to control when and how the workflow is triggered.
const { runId, messageId, workflowInput, threadId, resourceId } = await durableAgent.prepare(
'Summarize the document',
{
memory: { threadId: 'thread-1', resourceId: 'user-1' },
},
)
Returns:
interface PrepareResult {
runId: string
messageId: string
workflowInput: any
registryEntry: object
threadId?: string
resourceId?: string
}
CancellationDirect link to Cancellation
abortRunStream(runId)Direct link to abortrunstreamrunid
Stops a run by ID from any process that shares the agent's PubSub. Only the runId from stream() or prepare() is needed, so a later request can stop a run it didn't start. The run ends the same way as abort() on the stream result: the stream emits an abort chunk, onAbort fires, and observe() consumers see the stream end.
const { runId } = await durableAgent.stream('Summarize the document')
// later, from a different request handler
durableAgent.abortRunStream(runId)
The abort request is published over PubSub, so a multi-process deployment needs a shared backend such as RedisStreamsPubSub. The default EventEmitterPubSub only reaches runs in the same process. InngestAgent doesn't implement this method. Cancel an Inngest run with the abort() returned by its stream().
Stopping a durable run through abortRunStream() or abortThreadStream() requires @mastra/core@1.62.0 or later. Earlier versions record the abort without stopping the run.
Returns: boolean. true when this process aborted the run locally or can see it executing. The abort request is published either way.
abortThreadStream({ threadId, resourceId?, expectedRunId? })Direct link to abortthreadstream-threadid-resourceid-expectedrunid-
Aborts the active run on a memory thread with the same abort request as abortRunStream(). The run is resolved from this process's thread runtime, so it must have been started here or observed through subscribeToThread() on this process. The server route POST /agents/:agentId/threads/abort uses this method.
Pass expectedRunId when the request must only stop a specific run. If another queued run becomes active before the request is handled, the method returns false without aborting the successor. Omit expectedRunId to abort whichever run is active when the request is handled.
const aborted = durableAgent.abortThreadStream({
resourceId: 'user-1',
threadId: 'thread-1',
expectedRunId: runId,
})
Returns: boolean. false when this process has no active run recorded for the thread or the active run doesn't match expectedRunId. No abort request is sent in either case.
RecoveryDirect link to Recovery
Recovery re-drives a run from its last persisted snapshot. Mastra treats a sub-agent delegation as one generated agent-<name> tool step and doesn't checkpoint its inner LLM or tool calls independently. If recovery replays that step, the entire sub-agent run can execute again, including already-completed tool side effects. Make tools used by sub-agents idempotent, or deduplicate their external side effects. See crash recovery for operational guidance.
listActiveRuns(options?)Direct link to listactiverunsoptions
Lists this agent's runs whose persisted snapshot is in running status: runs whose agentic loop was mid-execution when the workflow engine last saved state. On a live process they transition to suspended or a terminal status. After a crash or restart they stay running with nothing driving them, which is what recoverActiveRuns() re-drives. Runs started by other durable agents on the same storage aren't included.
running snapshots are only written when recovery.durableAgents is 'auto' or a custom shouldPersistSnapshot includes running. Without one of those, this method returns no runs. See snapshot persistence.
const { runs, total } = await durableAgent.listActiveRuns({ resourceId: 'user-1' })
for (const run of runs) {
await durableAgent.recoverActiveRuns({ runId: run.runId })
}
Reads workflow snapshot storage on the Mastra instance the agent is registered with, and throws if the agent isn't registered with one. Use persistent storage: the default in-memory store loses these rows on the restart that orphans them.
Only runs in running status are returned. Suspended runs keep their snapshots but aren't included, and finished runs aren't retained. Filter by resourceId to scope the list to one user.
Returns:
interface DurableAgentListActiveRunsResult {
runs: Array<{
runId: string
status: 'running'
threadId?: string
resourceId?: string
updatedAt: Date
}>
total: number
}
options.threadId?:
options.resourceId?:
options.fromDate?:
options.toDate?:
options.perPage?:
page is also set. Must be a positive integer.options.page?:
perPage is also set. Must be a non-negative integer.recoverActiveRuns(options?)Direct link to recoveractiverunsoptions
Discovers runs stuck in running status for this agent and re-drives them from the last persisted snapshot. Accepts the same filters and pagination as listActiveRuns(), plus runId. Returns a summary of what was recovered.
const result = await durableAgent.recoverActiveRuns()
// { recovered: [{ runId, status }], succeeded: 2, failed: 0 }
Pass a runId to recover a single known run:
await durableAgent.recoverActiveRuns({ runId: 'run-abc-123' })
Returns:
interface DurableAgentRecoverActiveRunsResult {
recovered: Array<{ runId: string; status: 'success' | 'failed'; error?: Error }>
succeeded: number
failed: number
}
options.runId?:
recover(runId, options?)Direct link to recoverrunid-options
Recovers a single run by ID. Returns a streamable result with the same shape as stream(). Use this when you need to observe the recovery stream in real time.
const { output, cleanup } = await durableAgent.recover('run-abc-123', {
onChunk: chunk => console.log(chunk),
onError: ({ error }) => console.error(error),
})
await output.text
cleanup()
Returns: Promise<DurableAgentStreamResult>
Stream optionsDirect link to Stream options
stream() accepts a DurableAgentStreamOptions object. It supports the agent execution options below, plus lifecycle callbacks.
runId?:
resume(), observe(), or abortRunStream().instructions?:
context?:
memory?:
requestContext?:
maxSteps?:
toolsets?:
clientTools?:
toolChoice?:
activeTools?:
modelSettings?:
Authorization, X-Api-Key, and similar) are stripped from the serialized snapshot before it crosses process boundaries.stopWhen?:
maxSteps only.system?:
requireToolApproval?:
true or false to gate all or none, or a function for per-call policy. Function-form policies live on the in-process run registry; cross-process resumes fall back to a true shadow.autoResumeSuspendedTools?:
resume() call.closeOnSuspend?:
true, close the stream when the run suspends (for example for tool approval) so fullStream, text, and getFullOutput() resolve at the suspension boundary instead of staying open, matching non-durable Agent.stream(). Useful for callers such as AG-UI or A2A that need the stream to end. Defaults to false, which keeps the stream open across suspension so a later resume can continue streaming on the same reader. resume()/resumeStream() always return a fresh stream.toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
structuredOutput?:
untilIdle?:
true for the default 5-minute idle timeout, or { maxIdleMs } to customise. Equivalent to the deprecated streamUntilIdle() method. Also supported on resume().disableBackgroundTasks?:
tracingOptions?:
requestContextKeys forwarded to the agent and model spans. Fully JSON-serializable.actor?:
transform?:
transformToolPayload closure lives on the in-process run registry; only the JSON-safe targets shadow is serialized. Recovery after a process restart cannot restore the closure. For crash-survivable transcript transforms, configure transform.transcript on the tool with createTool() so the durable worker can re-resolve it from the registered tool.prepareStep?:
PrepareStepProcessor at the start of every iteration. Closure-only โ stored on the in-process run registry. Cross-process resumes lose the hook.isTaskComplete?:
onComplete live on the in-process run registry; the JSON-safe primitives (strategy, timeout, parallel, suppressFeedback, scorerNames) are serialized for cross-process observability.delegation?:
onDelegationStart, onDelegationComplete, messageFilter). Callbacks are baked into the sub-agent tool wrappers at prepare time. Cross-process resumes lose the callbacks.versions?:
abortSignal?:
AbortController, so either source can cancel the run. Cross-process resumes cannot recover the signal โ pass a fresh one to resume() if you need post-resume abortability.onChunk?:
onStepFinish?:
onFinish?:
onAbort; failed runs call onError.onError?:
onSuspended?:
onAbort?:
abortSignal, result.abort(), abortRunStream(), or abortThreadStream(). Receives the steps completed before the abort and, in text, the assistant text streamed so far.onIterationComplete?:
messageList, finishReason, and isFinal flag. Return continue: true to request another iteration or continue: false to stop. Feedback is added as a synthetic assistant message when the loop runs another iteration. If the model has stopped, feedback alone or continue: false with feedback does not start another iteration. When the loop was already going to continue, continue: false with feedback runs one final feedback-guided iteration before stopping.resume(), observe(), and recover() accept the same lifecycle callbacks (onChunk, onStepFinish, onFinish, onError, onAbort, onSuspended). observe() also accepts an offset to control where replay starts.
DurableAgentStreamResultDirect link to DurableAgentStreamResult
The object returned by stream(), resume(), observe(), and 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 for the full text, or consume output.fullStream.fullStream:
output.fullStream.runId:
resume() or observe() to reconnect.threadId?:
resourceId?:
cleanup:
abort:
AbortController. Surfaces as an AbortError inside the durable LLM-execution step and fires the onAbort callback. Safe to call after the run has finished โ a no-op in that case.