Skip to main content

DurableAgent

beta

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 example
Direct link to Usage example

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 },
})

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 flag
Direct 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.

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') 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

Parameters
Direct link to Parameters

agent:

Agent
The Agent to wrap with durable execution capabilities. Agent methods delegate to this agent.

id?:

string
= agent.id
ID override.

name?:

string
= agent.name
Name override.

cache?:

MastraServerCache | false
Cache for stored stream events, which enables resumable streams. If omitted, the agent inherits the cache from the Mastra instance, or uses an InMemoryServerCache. Set to false to disable caching, which makes streams non-resumable.

pubsub?:

PubSub
= EventEmitterPubSub
PubSub instance for streaming events.

maxSteps?:

number
Maximum number of steps for the agentic loop.

shouldCache?:

(topic: string) => boolean
Per-topic opt-out of the replay cache. Return 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?:

(params: { stepResults, workflowStatus }) => boolean
Predicate controlling which workflow snapshots the durable run persists. The default always persists 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)

Parameters
Direct link to Parameters

agent:

Agent
The Agent to wrap with evented durable execution capabilities.

cache?:

MastraServerCache | false
Cache for stored stream events, which enables resumable streams. If omitted, the agent inherits the cache from the Mastra instance, or uses an InMemoryServerCache. Set to false to disable caching.

pubsub?:

PubSub
= EventEmitterPubSub
PubSub instance for streaming events.

maxSteps?:

number
Maximum number of steps for the agentic loop.

shouldCache?:

(topic: string) => boolean
Per-topic opt-out of the replay cache. Return 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?:

(params: { stepResults, workflowStatus }) => boolean
Accepted for API symmetry with 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 parameters
Direct link to Constructor parameters

The DurableAgent class accepts the same options as createDurableAgent, plus cleanupTimeoutMs. Prefer the factory unless you need to subclass.

agent:

Agent
The Agent to wrap with durable execution capabilities.

id?:

string
= agent.id
ID override.

name?:

string
= agent.name
Name override.

cache?:

MastraServerCache | false
Cache for stored stream events. If omitted, inherits from the Mastra instance or uses an InMemoryServerCache. Set to false to disable caching.

pubsub?:

PubSub
= EventEmitterPubSub
PubSub instance for streaming events.

maxSteps?:

number
Maximum number of steps for the agentic loop.

shouldCache?:

(topic: string) => boolean
Per-topic opt-out of the replay cache. Return 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?:

(params: { stepResults, workflowStatus }) => boolean
Predicate controlling which workflow snapshots the durable run persists. The default always persists 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?:

number
= 30000
Grace period in milliseconds before registry entries are cleaned up automatically after a stream finishes or errors. Set to 0 to disable auto-cleanup and require a manual cleanup() call. Auto-cleanup does not fire on suspended events.

Methods
Direct link to Methods

Execution
Direct 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>

warning

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
}

Cancellation
Direct 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.

Recovery
Direct link to Recovery

warning

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

string
Only return runs that belong to this memory thread.

options.resourceId?:

string
Only return runs that belong to this memory resource.

options.fromDate?:

Date
Only return runs created at or after this date.

options.toDate?:

Date
Only return runs created at or before this date.

options.perPage?:

number
Runs per page. Applies only when page is also set. Must be a positive integer.

options.page?:

number
Zero-indexed page number. Applies only when 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?:

string
Recover a specific run by ID. When set, the discovery filters and pagination are ignored.

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 options
Direct link to Stream options

stream() accepts a DurableAgentStreamOptions object. It supports the agent execution options below, plus lifecycle callbacks.

runId?:

string
Unique identifier for this run. Use it later with resume(), observe(), or abortRunStream().

instructions?:

AgentExecutionOptions['instructions']
Overrides the agent's default instructions for this run. Accepts a static string or the same dynamic instructions value the agent supports.

context?:

ModelMessage[]
Additional context messages to provide to the agent.

memory?:

object
Memory configuration for conversation persistence and retrieval.

requestContext?:

RequestContext
Request context carrying dynamic configuration and state for this run.

maxSteps?:

number
Maximum number of steps to run for this stream.

toolsets?:

object
Additional tool sets available for this run.

clientTools?:

object
Client-side tools available during execution.

toolChoice?:

'auto' | 'none' | 'required' | { type: 'tool'; toolName: string }
Tool selection strategy.

activeTools?:

string[]
Restricts execution to the named subset of the agent's tools.

modelSettings?:

object
Model-specific settings such as temperature. Credential-bearing headers (Authorization, X-Api-Key, and similar) are stripped from the serialized snapshot before it crosses process boundaries.

stopWhen?:

AgentExecutionOptions['stopWhen']
Predicate or composition that ends the agentic loop early. The closure rides on the in-process run registry; cross-process resumes degrade to maxSteps only.

system?:

string | string[]
Additional system message appended after the agent instructions and before user messages.

requireToolApproval?:

boolean | ((args: { toolName: string; args: unknown; requestContext: RequestContext; workspace?: string }) => boolean | Promise<boolean>)
Require approval for tool calls. Pass 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?:

boolean
Automatically resume tools that suspended, instead of waiting for an external resume() call.

closeOnSuspend?:

boolean
= false
When 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?:

number
Maximum number of tool calls to execute concurrently.

includeRawChunks?:

boolean
Include raw provider chunks in the stream output.

maxProcessorRetries?:

number
Maximum number of processor retries per generation.

structuredOutput?:

object
Structured output configuration.

untilIdle?:

boolean | { maxIdleMs?: number }
When set, keeps the stream open across background-task continuations until the agent is idle. Pass true for the default 5-minute idle timeout, or { maxIdleMs } to customise. Equivalent to the deprecated streamUntilIdle() method. Also supported on resume().

disableBackgroundTasks?:

boolean
Disable background-task dispatch for this run. Background-eligible tools execute inline instead.

tracingOptions?:

AgentExecutionOptions['tracingOptions']
Tracing metadata, tags, trace ID, parent span ID, and requestContextKeys forwarded to the agent and model spans. Fully JSON-serializable.

actor?:

AgentExecutionOptions['actor']
Per-call actor signal forwarded to FGA checks and tool execution.

transform?:

AgentExecutionOptions['transform']
Per-invocation tool payload transform policy. The 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?:

AgentExecutionOptions['prepareStep']
Per-step preparation hook invoked as a PrepareStepProcessor at the start of every iteration. Closure-only โ€” stored on the in-process run registry. Cross-process resumes lose the hook.

isTaskComplete?:

AgentExecutionOptions['isTaskComplete']
Per-call completion policy. Scorer instances and onComplete live on the in-process run registry; the JSON-safe primitives (strategy, timeout, parallel, suppressFeedback, scorerNames) are serialized for cross-process observability.

delegation?:

AgentExecutionOptions['delegation']
Sub-agent delegation hooks (onDelegationStart, onDelegationComplete, messageFilter). Callbacks are baked into the sub-agent tool wrappers at prepare time. Cross-process resumes lose the callbacks.

versions?:

object
Version overrides for sub-agent delegation.

abortSignal?:

AbortSignal
External abort signal. Forwarded to the durable run's internal 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?:

(chunk: ChunkType) => void | Promise<void>
Called for each streamed chunk.

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
Called when a step in the agentic loop finishes.

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
Called when the run completes successfully. Aborted runs call onAbort; failed runs call onError.

onError?:

(error: Error) => void | Promise<void>
Called when the run errors.

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
Called when the run suspends, for example for tool approval.

onAbort?:

AgentExecutionOptions['onAbort']
Called when the run is aborted via abortSignal, result.abort(), abortRunStream(), or abortThreadStream(). Receives the steps completed before the abort and, in text, the assistant text streamed so far.

onIterationComplete?:

AgentExecutionOptions['onIterationComplete']
Called after every agentic-loop iteration with the latest 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.

DurableAgentStreamResult
Direct 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:

MastraModelOutput
The streaming output. Await output.text for the full text, or consume output.fullStream.

fullStream:

ReadableStream
The full event stream, delegating to output.fullStream.

runId:

string
The unique run ID. Pass it to resume() or observe() to reconnect.

threadId?:

string
Thread ID when using memory.

resourceId?:

string
Resource ID when using memory.

cleanup:

() => void
Unsubscribes from PubSub and clears registry entries for the run. Call it when done with the run.

abort:

() => void
Aborts the run by flipping the internal 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.