DurableAgent
DurableAgent enveloppe un Agent existant pour lui fournir une exécution durable et des flux pouvant être repris. Il exécute la boucle de l’Agent afin qu’un client puisse se déconnecter puis se reconnecter sans manquer d’événements, et diffuse ces événements via PubSub. Utilisez-le lorsqu’une exécution doit se poursuivre au-delà d’une seule requête ou survivre à une perte de connexion.
Créez-en un avec la fabrique createDurableAgent, ou utilisez createEventedAgent pour une exécution sans attente de résultat sur le moteur de Workflow intégré. Pour une exécution propulsée par Inngest, utilisez createInngestAgent depuis @mastra/inngest.
Exemple d’utilisationLien direct vers Exemple d’utilisation
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 },
})
Diffusez une réponse et lisez le résultat. La fonction cleanup se désabonne de PubSub lorsque vous avez terminé l’exécution :
const { output, runId, cleanup } = await durableAgent.stream('Hello!')
const text = await output.text
cleanup()
Utiliser l’option de configuration durableLien direct vers using-the-durable-config-flag
Définissez durable: true dans AgentConfig : l’Agent est alors automatiquement enveloppé avec createDurableAgent lorsqu’il est rattaché à une instance de Mastra. Utilisez un objet pour transmettre des options avancées telles que cache, pubsub, maxSteps ou 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') renvoie le DurableAgent enveloppé. Les Agents autonomes (construits mais jamais enregistrés auprès d’une instance de Mastra) ne deviennent pas durables. L’enveloppement est effectué lors de l’enregistrement.
createDurableAgent(options)Lien direct vers createdurableagentoptions
Enveloppe un Agent pour lui fournir une exécution durable et des flux pouvant être repris. C’est la méthode recommandée pour créer un DurableAgent.
import { createDurableAgent } from '@mastra/core/agent/durable'
const durableAgent = createDurableAgent({ agent })
Renvoie : DurableAgent
ParamètresLien direct vers Paramètres
agent:
id?:
name?:
cache?:
false pour désactiver la mise en cache, ce qui empêche la reprise des flux.pubsub?:
maxSteps?:
createEventedAgent(options)Lien direct vers createeventedagentoptions
Enveloppe un Agent afin d’assurer une exécution durable sans attente de résultat sur le moteur de Workflow intégré. Comme createDurableAgent, cette fonction renvoie un résultat à partir duquel diffuser le flux, mais le Workflow sous-jacent s’exécute de manière non bloquante (via startAsync) au lieu d’aller jusqu’à son terme avant que le flux ne soit raccordé. Utilisez-la lorsque l’exécution doit progresser indépendamment de l’appelant. Elle n’accepte pas le remplacement de id ni de name.
import { createEventedAgent } from '@mastra/core/agent/durable'
const eventedAgent = createEventedAgent({ agent })
Renvoie : EventedAgent (une sous-classe de DurableAgent)
ParamètresLien direct vers Paramètres
agent:
cache?:
false pour désactiver la mise en cache.pubsub?:
maxSteps?:
Paramètres du constructeurLien direct vers Paramètres du constructeur
La classe DurableAgent accepte les mêmes options que createDurableAgent, auxquelles s’ajoute cleanupTimeoutMs. Privilégiez la fabrique, sauf si vous devez créer une sous-classe.
agent:
id?:
name?:
cache?:
false pour désactiver la mise en cache.pubsub?:
maxSteps?:
cleanupTimeoutMs?:
cleanup(). Le nettoyage automatique ne se déclenche pas pour les événements suspendus.MéthodesLien direct vers Méthodes
ExécutionLien direct vers Exécution
stream(messages, options?)Lien direct vers streammessages-options
Diffuse une réponse au moyen d’une exécution durable. Renvoie immédiatement un résultat dont la propriété output produit des événements à mesure que l’exécution progresse.
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()
Renvoie : Promise<DurableAgentStreamResult>
resume(runId, resumeData, options?)Lien direct vers resumerunid-resumedata-options
Reprend une exécution suspendue, par exemple après l’approbation d’un Tool. Transmettez le runId du flux d’origine ainsi que les données attendues par l’exécution. Lève une exception si aucune entrée de registre n’existe pour cette exécution.
const { output, cleanup } = await durableAgent.resume(runId, {
approved: true,
})
await output.text
cleanup()
Renvoie : Promise<DurableAgentStreamResult>
observe(runId, options?)Lien direct vers observerunid-options
Se reconnecte à une exécution existante en rejouant les événements mis en cache avant de transmettre les événements en direct. Utilisez cette méthode après une déconnexion réseau. Transmettez offset pour commencer la relecture à partir d’une position connue.
const { output, cleanup } = await durableAgent.observe(runId, {
offset: 0,
onChunk: chunk => console.log(chunk),
})
await output.text
Par défaut, observe() attend indéfiniment les événements. Si le processus qui exécute l’opération s’arrête de façon inattendue, celle-ci cesse de produire des événements sans jamais émettre d’événement de fin ; le flux observé attendrait donc indéfiniment. Transmettez idleTimeoutMs pour limiter cette attente : après ce nombre de millisecondes sans activité, le flux prend fin. Une vérification facultative isAlive est d’abord consultée. Renvoyez true tant que l’exécution est toujours en cours de traitement (par exemple lors d’un appel de Tool de longue durée ou d’une exécution mise en pause dans l’attente d’une intervention humaine) afin de continuer à attendre. Renvoyer false, ou omettre isAlive, met fin au flux avec une erreur. Une exception temporaire levée par isAlive est interprétée comme « toujours active » ; l’échec ponctuel d’une vérification ne met donc jamais fin à un flux actif.
const { output } = await durableAgent.observe(runId, {
idleTimeoutMs: 30_000,
isAlive: () => runHeartbeat.isFresh(runId),
})
Mettre fin à une exécution après expiration du délai d’inactivité déclenche le même nettoyage qu’une exécution qui rencontre une erreur (voir l’avertissement ci-dessous) : son état mis en cache est donc libéré au lieu d’être conservé. Ces deux options sont facultatives. Omettez-les pour conserver le comportement antérieur d’attente indéfinie.
Renvoie : Promise<DurableAgentStreamResult>
La fonction cleanup() renvoyée par observe() détruit les entrées de registre et les événements mis en cache de l’exécution. Appelez-la uniquement lorsque vous avez terminé cette exécution. Si celle-ci est suspendue et que vous prévoyez de la reprendre plus tard, n’appelez pas cleanup(). Laissez le minuteur de nettoyage automatique s’en charger une fois l’exécution terminée ou en erreur. Le nettoyage automatique ne se déclenche pas pour les événements suspendus.
prepare(messages, options?)Lien direct vers preparemessages-options
Prépare une exécution durable sans la démarrer. Enregistre l’exécution dans le registre interne et renvoie l’entrée sérialisée du Workflow. Utilisez cette méthode lorsque vous devez contrôler le moment et la manière dont le Workflow est déclenché.
const { runId, messageId, workflowInput, threadId, resourceId } = await durableAgent.prepare(
'Summarize the document',
{
memory: { threadId: 'thread-1', resourceId: 'user-1' },
},
)
Renvoie :
interface PrepareResult {
runId: string
messageId: string
workflowInput: any
registryEntry: object
threadId?: string
resourceId?: string
}
RécupérationLien direct vers Récupération
recoverActiveRuns(options?)Lien direct vers recoveractiverunsoptions
Détecte les exécutions de cet Agent bloquées avec le statut running et les relance à partir du dernier instantané persistant. Récupère jusqu’à options.limit exécutions (valeur par défaut : 100). Renvoie un récapitulatif des éléments récupérés.
const result = await durableAgent.recoverActiveRuns()
// { recovered: [{ runId, status }], succeeded: 2, failed: 0 }
Transmettez un runId pour récupérer une seule exécution connue :
await durableAgent.recoverActiveRuns({ runId: 'run-abc-123' })
Renvoie :
interface DurableAgentRecoverActiveRunsResult {
recovered: Array<{ runId: string; status: 'success' | 'failed'; error?: Error }>
succeeded: number
failed: number
}
options.runId?:
options.limit?:
options.createdBefore?:
recover(runId, options?)Lien direct vers recoverrunid-options
Récupère une seule exécution à partir de son ID. Renvoie un résultat diffusable ayant la même structure que stream(). Utilisez cette méthode lorsque vous devez observer le flux de récupération en temps réel.
const { output, cleanup } = await durableAgent.recover('run-abc-123', {
onChunk: chunk => console.log(chunk),
onError: ({ error }) => console.error(error),
})
await output.text
cleanup()
Renvoie : Promise<DurableAgentStreamResult>
Options de fluxLien direct vers Options de flux
stream() accepte un objet DurableAgentStreamOptions. Il prend en charge les options d’exécution de l’Agent ci-dessous, ainsi que les fonctions de rappel du cycle de vie.
runId?:
resume() ou observe().instructions?:
context?:
memory?:
requestContext?:
maxSteps?:
toolsets?:
clientTools?:
toolChoice?:
activeTools?:
modelSettings?:
Authorization, X-Api-Key et similaires) sont retirés de l’instantané sérialisé avant son passage d’un processus à un autre.stopWhen?:
maxSteps.system?:
requireToolApproval?:
true ou false pour les soumettre tous à approbation ou n’en soumettre aucun, ou une fonction pour définir une stratégie par appel. Les stratégies sous forme de fonction résident dans le registre des exécutions du processus ; les reprises entre processus se rabattent sur une copie true.autoResumeSuspendedTools?:
resume().toolCallConcurrency?:
includeRawChunks?:
maxProcessorRetries?:
structuredOutput?:
untilIdle?:
true pour appliquer le délai d’inactivité par défaut de 5 minutes, ou { maxIdleMs } pour le personnaliser. Équivaut à la méthode obsolète streamUntilIdle(). Également pris en charge par resume().disableBackgroundTasks?:
tracingOptions?:
requestContextKeys transmis aux spans de l’Agent et du modèle. Entièrement sérialisable en JSON.actor?:
transform?:
transformToolPayload réside dans le registre des exécutions du processus ; seule la copie targets, compatible avec JSON, est sérialisée.prepareStep?:
PrepareStepProcessor au début de chaque itération. Cette fonction est stockée uniquement dans le registre des exécutions du processus. Les reprises entre processus perdent donc ce hook.isTaskComplete?:
onComplete résident dans le registre des exécutions du processus ; les primitives compatibles avec JSON (strategy, timeout, parallel, suppressFeedback, scorerNames) sont sérialisées pour permettre l’observabilité entre processus.delegation?:
onDelegationStart, onDelegationComplete, messageFilter). Les fonctions de rappel sont intégrées aux enveloppes de Tool des sous-Agents lors de la préparation. Les reprises entre processus perdent ces fonctions de rappel.versions?:
abortSignal?:
AbortController interne de l’exécution durable, afin que l’une ou l’autre source puisse annuler l’exécution. Les reprises entre processus ne peuvent pas récupérer ce signal — transmettez-en un nouveau à resume() si vous devez pouvoir annuler après la reprise.onChunk?:
onStepFinish?:
onFinish?:
onError?:
onSuspended?:
onAbort?:
abortSignal ou result.abort().onIterationComplete?:
messageList, finishReason et de l’indicateur isFinal. Sert uniquement à l’observation sur les Agents durables : renvoyer continue: false ou un retour n’influence pas la boucle.resume() et observe() acceptent les mêmes fonctions de rappel du cycle de vie (onChunk, onStepFinish, onFinish, onError, onSuspended). observe() accepte également un offset pour définir le point de départ de la relecture.
DurableAgentStreamResultLien direct vers DurableAgentStreamResult
Objet renvoyé par stream(), resume(), observe() et 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 pour obtenir le texte complet, ou consommez output.fullStream.fullStream:
output.fullStream.runId:
resume() ou observe() pour vous reconnecter.threadId?:
resourceId?:
cleanup:
abort:
AbortController interne. Se manifeste sous la forme d’une AbortError dans l’étape d’exécution durable du LLM et déclenche la fonction de rappel onAbort. Peut être appelée sans risque après la fin de l’exécution — elle est alors sans effet.