Aller au contenu principal

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’utilisation
Lien direct vers Exemple d’utilisation

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

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 durable
Lien 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.

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') 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ètres
Lien direct vers Paramètres

agent:

Agent
L’Agent à envelopper avec des capacités d’exécution durable. Les méthodes de l’Agent délèguent leurs opérations à cet Agent.

id?:

string
= agent.id
Remplacement de l’ID.

name?:

string
= agent.name
Remplacement du nom.

cache?:

MastraServerCache | false
Cache des événements de flux stockés, qui permet de reprendre les flux. Si cette option est omise, l’Agent hérite du cache de l’instance Mastra ou utilise un InMemoryServerCache. Définissez-la sur false pour désactiver la mise en cache, ce qui empêche la reprise des flux.

pubsub?:

PubSub
= EventEmitterPubSub
Instance PubSub utilisée pour diffuser les événements.

maxSteps?:

number
Nombre maximal d’étapes de la boucle de l’Agent.

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ètres
Lien direct vers Paramètres

agent:

Agent
L’Agent à envelopper avec des capacités d’exécution durable pilotée par événements.

cache?:

MastraServerCache | false
Cache des événements de flux stockés, qui permet de reprendre les flux. Si cette option est omise, l’Agent hérite du cache de l’instance Mastra ou utilise un InMemoryServerCache. Définissez-la sur false pour désactiver la mise en cache.

pubsub?:

PubSub
= EventEmitterPubSub
Instance PubSub utilisée pour diffuser les événements.

maxSteps?:

number
Nombre maximal d’étapes de la boucle de l’Agent.

Paramètres du constructeur
Lien 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:

Agent
L’Agent à envelopper avec des capacités d’exécution durable.

id?:

string
= agent.id
Remplacement de l’ID.

name?:

string
= agent.name
Remplacement du nom.

cache?:

MastraServerCache | false
Cache des événements de flux stockés. Si cette option est omise, il est hérité de l’instance Mastra ou un InMemoryServerCache est utilisé. Définissez-la sur false pour désactiver la mise en cache.

pubsub?:

PubSub
= EventEmitterPubSub
Instance PubSub utilisée pour diffuser les événements.

maxSteps?:

number
Nombre maximal d’étapes de la boucle de l’Agent.

cleanupTimeoutMs?:

number
= 30000
Délai de grâce, en millisecondes, avant le nettoyage automatique des entrées du registre lorsqu’un flux se termine ou rencontre une erreur. Définissez cette valeur sur 0 pour désactiver le nettoyage automatique et exiger un appel manuel à cleanup(). Le nettoyage automatique ne se déclenche pas pour les événements suspendus.

Méthodes
Lien direct vers Méthodes

Exécution
Lien 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>

attention

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ération
Lien 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?:

string
Récupère une exécution particulière à partir de son ID. Lorsque cette option est définie, les filtres de détection sont ignorés.

options.limit?:

number
Nombre maximal d’exécutions actives à détecter. La valeur par défaut est 100.

options.createdBefore?:

Date
Récupère uniquement les exécutions créées avant cette date.

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

string
Identifiant unique de cette exécution. Utilisez-le ensuite avec resume() ou observe().

instructions?:

AgentExecutionOptions['instructions']
Remplace les instructions par défaut de l’Agent pour cette exécution. Accepte une chaîne statique ou la même valeur d’instructions dynamiques que celle prise en charge par l’Agent.

context?:

ModelMessage[]
Messages de contexte supplémentaires à fournir à l’Agent.

memory?:

object
Configuration de la mémoire pour la persistance et la récupération des conversations.

requestContext?:

RequestContext
Contexte de requête transportant la configuration dynamique et l’état de cette exécution.

maxSteps?:

number
Nombre maximal d’étapes à exécuter pour ce flux.

toolsets?:

object
Ensembles de Tools supplémentaires disponibles pour cette exécution.

clientTools?:

object
Tools côté client disponibles pendant l’exécution.

toolChoice?:

'auto' | 'none' | 'required' | { type: 'tool'; toolName: string }
Stratégie de sélection des Tools.

activeTools?:

string[]
Limite l’exécution au sous-ensemble nommé des Tools de l’Agent.

modelSettings?:

object
Paramètres propres au modèle, tels que la température. Les en-têtes contenant des identifiants (Authorization, X-Api-Key et similaires) sont retirés de l’instantané sérialisé avant son passage d’un processus à un autre.

stopWhen?:

AgentExecutionOptions['stopWhen']
Prédicat ou composition qui met fin de manière anticipée à la boucle de l’Agent. La fonction correspondante est conservée dans le registre des exécutions du processus ; les reprises entre processus se rabattent uniquement sur maxSteps.

system?:

string | string[]
Message système supplémentaire ajouté après les instructions de l’Agent et avant les messages utilisateur.

requireToolApproval?:

boolean | ((args: { toolName: string; args: unknown; requestContext: RequestContext; workspace?: string }) => boolean | Promise<boolean>)
Exige une approbation pour les appels de Tool. Transmettez 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?:

boolean
Reprend automatiquement les Tools suspendus au lieu d’attendre un appel externe à resume().

toolCallConcurrency?:

number
Nombre maximal d’appels de Tool à exécuter simultanément.

includeRawChunks?:

boolean
Inclut les fragments bruts du Provider dans la sortie du flux.

maxProcessorRetries?:

number
Nombre maximal de nouvelles tentatives du processeur par génération.

structuredOutput?:

object
Configuration de la sortie structurée.

untilIdle?:

boolean | { maxIdleMs?: number }
Lorsque cette option est définie, maintient le flux ouvert pendant les continuations des tâches en arrière-plan jusqu’à ce que l’Agent soit inactif. Transmettez 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?:

boolean
Désactive la répartition des tâches en arrière-plan pour cette exécution. Les Tools éligibles à l’arrière-plan sont alors exécutés directement.

tracingOptions?:

AgentExecutionOptions['tracingOptions']
Métadonnées de traçage, tags, ID de trace, ID du span parent et requestContextKeys transmis aux spans de l’Agent et du modèle. Entièrement sérialisable en JSON.

actor?:

AgentExecutionOptions['actor']
Signal d’acteur propre à chaque appel, transmis aux vérifications FGA et à l’exécution des Tools.

transform?:

AgentExecutionOptions['transform']
Stratégie de transformation de la charge utile du Tool propre à chaque invocation. La fonction transformToolPayload réside dans le registre des exécutions du processus ; seule la copie targets, compatible avec JSON, est sérialisée.

prepareStep?:

AgentExecutionOptions['prepareStep']
Hook de préparation propre à chaque étape, invoqué comme 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?:

AgentExecutionOptions['isTaskComplete']
Stratégie d’achèvement propre à chaque appel. Les instances de Scorer et 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?:

AgentExecutionOptions['delegation']
Hooks de délégation aux sous-Agents (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?:

object
Remplacements de version pour la délégation aux sous-Agents.

abortSignal?:

AbortSignal
Signal d’annulation externe. Transmis à l’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?:

(chunk: ChunkType) => void | Promise<void>
Appelé pour chaque fragment diffusé.

onStepFinish?:

(result: AgentStepFinishEventData) => void | Promise<void>
Appelé lorsqu’une étape de la boucle de l’Agent se termine.

onFinish?:

(result: AgentFinishEventData) => void | Promise<void>
Appelé lorsque l’exécution se termine.

onError?:

(error: Error) => void | Promise<void>
Appelé lorsque l’exécution rencontre une erreur.

onSuspended?:

(data: AgentSuspendedEventData) => void | Promise<void>
Appelé lorsque l’exécution est suspendue, par exemple pour l’approbation d’un Tool.

onAbort?:

AgentExecutionOptions['onAbort']
Appelé lorsque l’exécution est annulée via abortSignal ou result.abort().

onIterationComplete?:

AgentExecutionOptions['onIterationComplete']
Appelé après chaque itération de la boucle de l’Agent avec les dernières valeurs de 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.

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

MastraModelOutput
Sortie diffusée en continu. Attendez output.text pour obtenir le texte complet, ou consommez output.fullStream.

fullStream:

ReadableStream
Flux complet des événements, délégué à output.fullStream.

runId:

string
ID unique de l’exécution. Transmettez-le à resume() ou observe() pour vous reconnecter.

threadId?:

string
ID du fil lors de l’utilisation de la mémoire.

resourceId?:

string
ID de la ressource lors de l’utilisation de la mémoire.

cleanup:

() => void
Se désabonne de PubSub et efface les entrées de registre de l’exécution. Appelez cette fonction lorsque vous avez terminé l’exécution.

abort:

() => void
Annule l’exécution en faisant basculer l’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.