Interface Processor
L’interface Processor définit le contrat de tous les Processors dans Mastra. Les Processors peuvent implémenter une ou plusieurs méthodes pour gérer différentes étapes du pipeline d’exécution de l’Agent.
Quand les méthodes des Processors s’exécutentLien direct vers Quand les méthodes des Processors s’exécutent
Les méthodes des Processors s’exécutent à différents moments du cycle de vie de l’Agent :
┌────────────────────────────────────────────────────────────────────┐
│ Agent Execution Flow │
├────────────────────────────────────────────────────────────────────┤
│ │
│ User Input │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ processInput │ ← Runs ONCE at start │
│ └───────────┬────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────────────┐ │
│ │ Agentic Loop │ │
│ │ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processInputStep │ ← Runs at EACH step │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processLLMRequest │ ← Before provider call │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ LLM Execution ──── API Error? ───┐ │ │
│ │ │ │ │ │
│ │ │ ┌───────────┴──────────┐ │ │
│ │ │ │ processAPIError │ │ │
│ │ │ └──────────────────────┘ │ │
│ │ │ (retry loops back to LLM) │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processOutputStream │ ← Runs on EACH stream chunk │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processLLMResponse │ ← After stream completes │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processOutputStep │ ← Runs after EACH LLM step │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ Tool Execution (if needed) │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processToolResult │ ← Runs per tool, after each │ │
│ │ └───────────┬────────────┘ tool.execute() returns │ │
│ │ │ │ │
│ │ └──────── Loop back if tools called ────────────│ │
│ │ │ │
│ └──────────────────────────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ processOutputResult │ ← Runs ONCE after completion │
│ └────────────────────────┘ │
│ │ │
│ ▼ │
│ Final Response │
│ │
└────────────────────────────────────────────────────────────────────┘
| Méthode | Moment d’exécution | Cas d’utilisation |
|---|---|---|
processInput | Une fois au début, avant la boucle agentique | Valider ou transformer l’entrée utilisateur initiale, ajouter du contexte |
processInputStep | À chaque étape de la boucle agentique, avant chaque appel au LLM | Transformer les messages entre les étapes, gérer les résultats des Tools |
processLLMRequest | Après la conversion de la requête LLM, avant l’appel au Provider | Réécrire le LanguageModelV2Prompt sortant pour l’appel actuel sans persister les modifications |
processAPIError | Lorsqu’un appel à l’API du LLM échoue | Examiner les rejets de l’API, modifier éventuellement l’état ou les messages et demander une nouvelle tentative |
processOutputStream | À chaque fragment diffusé pendant la réponse du LLM | Filtrer ou modifier le contenu diffusé, détecter des motifs en temps réel |
processLLMResponse | Une fois l’étape LLM terminée et les fragments du flux collectés | Capturer ou mettre en cache la réponse complète, exécuter les effets secondaires postérieurs à l’appel associés à processLLMRequest |
processOutputStep | Après chaque réponse du LLM, avant l’exécution des Tools | Valider la qualité de la sortie, mettre en œuvre des garde-fous avec nouvelle tentative |
processToolResult | Pour chaque Tool, après le retour de tool.execute() et avant l’ajout du résultat à la liste des messages | Rechercher une injection de prompt dans la sortie du Tool, masquer les champs sensibles, interrompre en cas de violation de politique |
processOutputResult | Une fois la génération terminée | Post-traiter la réponse finale, journaliser les résultats |
Définition de l’interfaceLien direct vers Définition de l’interface
interface Processor<TId extends string = string, TTripwireMetadata = unknown> {
readonly id: TId
readonly name?: string
readonly description?: string
/** Index of this processor in the workflow (set at runtime when combining processors). */
processorIndex?: number
/** When true, processOutputStream also receives `data-*` chunks. Default: false. */
processDataParts?: boolean
/** Callback invoked when this processor detects a violation, regardless of strategy. */
onViolation?: (violation: ProcessorViolation) => void | Promise<void>
processInput?(
args: ProcessInputArgs<TTripwireMetadata>,
): Promise<ProcessInputResult> | ProcessInputResult
processInputStep?(
args: ProcessInputStepArgs<TTripwireMetadata>,
):
| Promise<ProcessInputStepResult | MessageList | MastraDBMessage[] | undefined | void>
| ProcessInputStepResult
| MessageList
| MastraDBMessage[]
| void
| undefined
processLLMRequest?(
args: ProcessLLMRequestArgs<TTripwireMetadata>,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult
processLLMResponse?(
args: ProcessLLMResponseArgs<TTripwireMetadata>,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult
processAPIError?(
args: ProcessAPIErrorArgs<TTripwireMetadata>,
): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void
processOutputStream?(
args: ProcessOutputStreamArgs<TTripwireMetadata>,
): Promise<ChunkType | null | undefined>
processOutputStep?(args: ProcessOutputStepArgs<TTripwireMetadata>): ProcessorMessageResult
processToolResult?(args: ProcessToolResultArgs<TTripwireMetadata>): ProcessorMessageResult
processOutputResult?(args: ProcessOutputResultArgs<TTripwireMetadata>): ProcessorMessageResult
}
PropriétésLien direct vers Propriétés
id:
name?:
description?:
processorIndex?:
processDataParts?:
data-* émis par les Tools via writer.custom(). La valeur par défaut est false.onViolation?:
Arguments de messageLien direct vers Arguments de message
La plupart des méthodes des Processors reçoivent à la fois messages et messageList. Ces deux valeurs désignent la même conversation sous-jacente, mais l’exposent différemment.
messages ou messageListLien direct vers messages-vs-messagelist
messages: tableau simple d’objetsMastraDBMessage, limité à l’étape actuelle. PourprocessInputetprocessInputStep, il exclut les messages système. PourprocessOutputResultetprocessOutputStep, il inclut la dernière réponse du LLM. Le tableau s’appuie surmessageList; toute modification sur place ducontent.partsd’un message est donc visible par les Processors en aval et lors de la persistance.messageList: instance active deMessageListqui sous-tend l’exécution. Elle expose des vues filtrées (entrée, réponse, éléments mémorisés, totalité), plusieurs formats de sortie (db, ui, core) et des méthodes permettant de modifier la conversation.
Utilisez messages si vous devez uniquement lire, parcourir avec map ou modifier légèrement les champs des messages de l’étape actuelle. Utilisez messageList lorsque vous devez :
- Lire les messages d’une autre étape, par exemple les messages d’entrée pendant le traitement de la sortie.
- Ajouter, supprimer ou remplacer des messages entiers.
- Convertir les messages dans un autre format, tel que des messages UI ou core pour une API tierce.
messages est toujours dérivé de messageList. Modifier messageList constitue donc la méthode canonique pour ajouter, supprimer ou réordonner des messages. Pour les modifications sur place du contenu d’un message, par exemple la réécriture de content.parts, modifier directement messages revient au même. Si vous renvoyez un nouveau tableau depuis messages, Mastra le réconcilie avec messageList pour l’étape actuelle.
PersistanceLien direct vers Persistance
Lorsque la Memory est activée, seul le contenu final de messageList, une fois tous les Processors terminés, est persisté dans le stockage. Les deux styles de retour sont équivalents pour la persistance :
- Lorsque vous modifiez directement
messageListou renvoyez la même instance deMessageList, les mutations enregistrées sont appliquées sur place ; la conversation sauvegardée reflète donc vos modifications. - Lorsque vous renvoyez un
MastraDBMessage[]ou{ messages, systemMessages }, Mastra réconcilie le tableau renvoyé avecmessageListpour l’étape actuelle, en supprimant les messages absents et en remplaçant les messages système.
Renvoyer une autre instance de MessageList constitue une erreur. Modifiez toujours celle qui est transmise à votre Processor.
Lire le texte d’un messageLien direct vers Lire le texte d’un message
MastraDBMessage.content utilise un objet structuré. Les chaînes ne sont pas prises en charge. La méthode canonique pour lire le texte de l’utilisateur ou de l’assistant consiste à utiliser content.parts :
import type { MastraDBMessage } from '@mastra/core/memory'
function getText(message: MastraDBMessage): string {
let text = ''
if (message.content.parts) {
for (const part of message.content.parts) {
if (part.type === 'text' && typeof part.text === 'string') {
text += part.text
}
}
}
// Fallback for legacy messages that only have the flattened `content` string
if (!text && typeof message.content.content === 'string') {
text = message.content.content
}
return text
}
Points essentiels :
message.content.partsest la source principale. Un message peut contenir plusieurs parties, y compris des parties non textuelles telles que des appels de Tools, des résultats de Tools et des fichiers. Filtrez avecpart.type === 'text'avant de lirepart.text.message.content.contentest une chaîne aplatie conservée pour la rétrocompatibilité. Utilisez-la uniquement comme solution de repli lorsquepartsest vide ou absent.message.contentn’est jamais une simple chaîne dansMastraDBMessage. Les anciennes structuresCoreMessagepeuvent être des chaînes, mais les Processors reçoivent toujours unMastraDBMessage.
MéthodesLien direct vers Méthodes
processInputLien direct vers processinput
Traite les messages d’entrée avant leur envoi au LLM. S’exécute une fois au début de l’exécution de l’Agent.
processInput?(args: ProcessInputArgs): Promise<ProcessInputResult> | ProcessInputResult;
ProcessInputArgsLien direct vers processinputargs
messages:
systemMessages:
messageList:
abort:
retry: true pour demander au LLM de réessayer l’étape avec un retour.retryCount:
tracingContext?:
requestContext?:
ProcessInputResultLien direct vers processinputresult
La méthode peut renvoyer l’un des trois types suivants :
MastraDBMessage[]:
MessageList:
{ messages, systemMessages }:
processInputStepLien direct vers processinputstep
Traite les messages d’entrée à chaque étape de la boucle agentique, avant leur envoi au LLM. Contrairement à processInput, qui ne s’exécute qu’une fois au début, cette méthode s’exécute à chaque étape, y compris lors de la poursuite des appels de Tools.
processInputStep?<TTripwireMetadata = unknown>(
args: ProcessInputStepArgs<TTripwireMetadata>,
):
| Promise<ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined>
| ProcessInputStepResult
| MessageList
| MastraDBMessage[]
| void
| undefined;
Ordre d’exécution dans la boucle agentiqueLien direct vers Ordre d’exécution dans la boucle agentique
processInput(une fois au début)processInputStepdepuis inputProcessors (à chaque étape, avant l’appel au LLM)- Callback
prepareStep(s’exécute dans le pipeline processInputStep, après inputProcessors) processLLMRequestdepuis inputProcessors (après la conversion du prompt, avant l’appel au Provider)- Exécution du LLM
processOutputStreamdepuis outputProcessors (à chaque fragment diffusé)processLLMResponsedepuis inputProcessors (une fois le flux terminé, en association avecprocessLLMRequest)processOutputStepdepuis outputProcessors (après la réponse du LLM, avant l’exécution des Tools)- Exécution des Tools (si nécessaire)
- Reprise à l’étape 2 si des Tools ont été appelés
ProcessInputStepArgsLien direct vers processinputstepargs
messages:
messageList:
stepNumber:
steps:
systemMessages:
model:
toolChoice?:
activeTools?:
tools?:
providerOptions?:
modelSettings?:
structuredOutput?:
abort:
retry: true pour demander au LLM de réessayer l’étape avec un retour.retryCount:
ProcessorContext. Commence à 0 ; utilisez-le pour plafonner les nouvelles tentatives déclenchées par les Processors.tracingContext?:
requestContext?:
ProcessInputStepResultLien direct vers processinputstepresult
processInputStep peut renvoyer plusieurs formes :
- Objet
ProcessInputStepResult: remplacez n’importe quelle combinaison des propriétés ci-dessous pour cette étape (décrites ensuite). MessageList: renvoyez la même instance demessageListpour signaler que vous avez modifié les messages sur place.MastraDBMessage[]: renvoyez un tableau de messages transformés. Il remplace les messages de l’étape.voidouundefined: ne renvoyez rien pour laisser l’étape inchangée.
La forme objet peut renvoyer n’importe quelle combinaison des propriétés suivantes :
model?:
toolChoice?:
activeTools?:
tools?:
messages?:
messageList?:
systemMessages?:
providerOptions?:
modelSettings?:
structuredOutput?:
Chaînage des ProcessorsLien direct vers Chaînage des Processors
Lorsque plusieurs Processors implémentent processInputStep, ils s’exécutent dans l’ordre et les modifications se propagent dans la chaîne :
Processor 1: receives { model: 'gpt-5.4' } → returns { model: 'gpt-5.4-mini' }
Processor 2: receives { model: 'gpt-5.4-mini' } → returns { toolChoice: 'none' }
Final: model = 'gpt-5.4-mini', toolChoice = 'none'
Isolation des messages systèmeLien direct vers Isolation des messages système
Les messages système sont réinitialisés à leurs valeurs d’origine au début de chaque étape. Les modifications apportées dans processInputStep n’affectent que l’étape actuelle, pas les étapes suivantes.
Cas d’utilisationLien direct vers Cas d’utilisation
- Changer dynamiquement de modèle selon le numéro d’étape ou le contexte
- Désactiver les Tools après un certain nombre d’étapes
- Ajouter ou remplacer dynamiquement des Tools selon le contexte de la conversation
- Transformer les types de parties de message entre les Providers (par exemple,
reasoning→thinkingpour Anthropic) - Modifier les messages selon le numéro d’étape ou le contexte accumulé
- Ajouter des instructions système propres à une étape
- Ajuster les options du Provider à chaque étape (par exemple, le contrôle du cache)
- Modifier le schéma de sortie structurée selon le contexte de l’étape
processLLMRequestLien direct vers processllmrequest
Traite la requête LLM finale après la conversion de la MessageList en LanguageModelV2Prompt par Mastra et avant l’appel au Provider. Utilisez cette méthode pour les réécritures temporaires tenant compte du modèle, qui ne doivent affecter que la requête sortante actuelle.
Les modifications du prompt renvoyé sont transmises au modèle uniquement pour l’appel actuel. Elles ne sont pas repersistées dans MessageList, la Memory, l’historique de l’UI ni les appels ultérieurs au Provider.
processLLMRequest?(
args: ProcessLLMRequestArgs,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult;
ProcessLLMRequestArgsLien direct vers processllmrequestargs
prompt:
model:
stepNumber:
steps:
state:
abort:
tripwire.retryCount:
ProcessorContext. Commence à 0 ; utilisez-le pour plafonner les nouvelles tentatives déclenchées par les Processors.requestContext?:
tracingContext?:
writer?:
writer.custom() pour émettre un fragment data-*.abortSignal?:
Valeur de retourLien direct vers Valeur de retour
processLLMRequest renvoie ProcessLLMRequestResult, dont le type est { prompt?: LanguageModelV2Prompt } | undefined | void.
- Renvoyez
{ prompt }pour remplacer le prompt sortant de l’appel actuel au Provider. - Renvoyez
undefinedouvoidpour transmettre le prompt d’origine sans modification.
Cas d’utilisationLien direct vers Cas d’utilisation
- Supprimer ou restructurer les parties du prompt propres au Provider avant un appel au modèle
- Normaliser les rôles ou le contenu afin de respecter les exigences d’entrée d’un Provider
- Adapter les formats des résultats de Tools lors d’un changement de Provider au milieu de la boucle
processLLMResponseLien direct vers processllmresponse
Traite la réponse du LLM une fois l’étape terminée, ou après la relecture d’une réponse mise en cache, et après la collecte des fragments de réponse par les Processors de sortie. Ce hook est associé à processLLMRequest : utilisez processLLMRequest pour stocker un état, tel qu’une clé de cache, avant l’appel au Provider, puis processLLMResponse pour agir sur la réponse terminée, par exemple en l’écrivant dans un cache.
L’objet state est la même instance que celle transmise à processLLMRequest pour la même étape. Les Processors peuvent donc mettre en relation les traitements antérieurs et postérieurs à l’appel.
processLLMResponse?(
args: ProcessLLMResponseArgs,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult;
ProcessLLMResponseArgsLien direct vers processllmresponseargs
chunks:
{ type, payload }).model:
stepNumber:
steps:
state:
processLLMRequest pour la même étape. Utilisez-le pour transmettre des données entre les deux hooks, par exemple une clé de cache.fromCache:
true, la réponse a été relue depuis un cache par l’intermédiaire de processLLMRequest, qui a renvoyé { response }. Les Processors qui écrivent dans un cache doivent ignorer l’écriture lorsque cette valeur est true.warnings?:
request?:
rawResponse?:
abort:
retryCount:
0 ; utilisez-le pour plafonner les nouvelles tentatives déclenchées par les Processors.requestContext?:
tracingContext?:
writer?:
abortSignal?:
Valeur de retourLien direct vers Valeur de retour
processLLMResponse renvoie ProcessLLMResponseResult, dont le type est undefined | void. La valeur de retour est réservée à de futures extensions.
Cas d’utilisationLien direct vers Cas d’utilisation
- Écrire les réponses du LLM dans un cache après un appel en direct, en association avec la dérivation de la clé de cache dans
processLLMRequest - Journaliser ou enregistrer la réponse complète à des fins d’analyse
- Déclencher des effets secondaires selon la réponse terminée
processAPIErrorLien direct vers processapierror
Gère les erreurs de rejet de l’API du LLM avant qu’elles ne deviennent des erreurs finales. Cette méthode s’exécute lorsque l’appel à l’API échoue avec une erreur non réessayable, telle qu’un code d’état 400 ou 422. Contrairement à processOutputStep, qui s’exécute après les réponses réussies, elle s’exécute lorsque l’API rejette la requête.
Ajoutez les Processors qui implémentent processAPIError au tableau errorProcessors d’un Agent.
Les Processors peuvent examiner l’erreur et modifier la requête, par exemple en ajoutant des messages à la messageList. Renvoyez { retry: true } pour effectuer une nouvelle tentative avec l’état modifié.
processAPIError?(args: ProcessAPIErrorArgs): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void;
ProcessAPIErrorArgsLien direct vers processapierrorargs
error:
messages:
messageList:
stepNumber:
steps:
state:
retryCount:
abort:
writer?:
writer.custom() pour émettre un fragment data-*.requestContext?:
abortSignal?:
ProcessAPIErrorResultLien direct vers processapierrorresult
retry:
Cas d’utilisationLien direct vers Cas d’utilisation
- Gérer les rejets propres à l’API en modifiant la requête puis en réessayant
- Convertir les erreurs non réessayables en erreurs réessayables grâce à des modifications de la requête
- Mettre en œuvre des stratégies de récupération d’erreur propres au modèle
Exemple : récupération d’erreur personnaliséeLien direct vers Exemple : récupération d’erreur personnalisée
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'
export class ErrorRecoveryProcessor implements Processor {
id = 'error-recovery'
processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
// Only retry once
if (retryCount > 0) return
// Check for a specific API error
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
// Trim older messages to fit within context
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}
processOutputStreamLien direct vers processoutputstream
Traite les fragments de sortie diffusés avec une gestion intégrée de l’état. Permet aux Processors d’accumuler des fragments et de prendre des décisions à partir d’un contexte plus large.
processOutputStream?(args: ProcessOutputStreamArgs): Promise<ChunkType | null | undefined>;
ProcessOutputStreamArgsLien direct vers processoutputstreamargs
part:
streamParts:
state:
abort:
tripwire. Transmettez retry: true pour demander une nouvelle tentative au LLM au lieu de terminer.retryCount:
ProcessorContext. Commence à 0 ; utilisez-le pour plafonner les nouvelles tentatives déclenchées par les Processors.messageList?:
tracingContext?:
requestContext?:
writer?:
Valeur de retourLien direct vers Valeur de retour
processOutputStream renvoie Promise<ChunkType | null | undefined>.
- Renvoyez le
ChunkTypepour émettre le fragment. Renvoyez lepartd’origine pour l’émettre sans modification, ou un nouveauChunkTypepour émettre un fragment modifié. - Renvoyez
nullpour supprimer le fragment. Rien n’est envoyé au Processor suivant ni au client. - Renvoyez
undefined, y compris la valeurundefinedimplicite d’une instructionreturn;ou d’une méthode arrivant à sa fin, pour supprimer le fragment.nulletundefinedse comportent de la même manière.
La suppression d’un fragment n’affecte que ce fragment. Le flux continue et le fragment suivant est toujours traité. Pour arrêter entièrement le flux, appelez abort().
processOutputResultLien direct vers processoutputresult
Traite le résultat de sortie complet une fois la diffusion ou la génération terminée.
processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;
ProcessOutputResultArgsLien direct vers processoutputresultargs
messages:
messageList:
state:
result:
text (texte accumulé), usage (utilisation des tokens avec inputTokens, outputTokens et totalTokens), finishReason (raison de la fin de la génération) et steps (tous les résultats des étapes LLM, chacun avec toolCalls, toolResults, reasoning, sources, files, etc.).abort:
tripwire.retryCount:
ProcessorContext. Commence à 0 ; utilisez-le pour plafonner les nouvelles tentatives déclenchées par les Processors.tracingContext?:
requestContext?:
writer?:
processOutputStepLien direct vers processoutputstep
Traite la sortie après chaque réponse du LLM dans la boucle agentique, avant l’exécution des Tools. Contrairement à processOutputResult, qui ne s’exécute qu’une fois à la fin, cette méthode s’exécute à chaque étape. Elle est idéale pour mettre en œuvre des garde-fous capables de déclencher de nouvelles tentatives.
processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;
ProcessOutputStepArgsLien direct vers processoutputstepargs
messages:
messageList:
stepNumber:
finishReason?:
providerMetadata?:
steps est vide.toolCalls?:
text?:
usage:
inputTokens, outputTokens, totalTokens).systemMessages:
steps:
state:
abort:
retry: true pour demander au LLM de réessayer l’étape.retryCount:
tracingContext?:
requestContext?:
Cas d’utilisationLien direct vers Cas d’utilisation
- Mettre en œuvre des garde-fous de qualité capables de demander de nouvelles tentatives
- Valider la sortie du LLM avant l’exécution des Tools
- Ajouter une journalisation ou des métriques propres à chaque étape
- Mettre en œuvre une modération de la sortie avec possibilité de nouvelle tentative
Exemple : garde-fou de qualité avec nouvelle tentativeLien direct vers Exemple : garde-fou de qualité avec nouvelle tentative
import type { Processor } from '@mastra/core/processors'
export class QualityGuardrail implements Processor {
id = 'quality-guardrail'
async processOutputStep({ text, abort, retryCount }) {
const score = await evaluateResponseQuality(text)
if (score < 0.7) {
if (retryCount < 3) {
// Request retry with feedback for the LLM
abort('Response quality too low. Please provide more detail.', {
retry: true,
metadata: { qualityScore: score },
})
} else {
// Max retries reached, block the response
abort('Response quality too low after multiple attempts.')
}
}
return []
}
}
processToolResultLien direct vers processtoolresult
Traite le résultat d’un Tool après le retour de tool.execute() et avant que ce résultat soit ajouté à la liste des messages ou transmis à l’appel suivant au LLM. Cette méthode est symétrique à processOutputStep, qui se déclenche avant l’exécution du Tool. Utilisez-la pour rechercher une injection de prompt dans la sortie du Tool, masquer les champs sensibles ou interrompre l’exécution avec abort('reason', { retry: true }).
Pour remplacer le résultat du Tool, modifiez messageList sur place à l’aide de messageList.updateToolInvocation. Le runtime relit le résultat postérieur au Processor depuis la liste des messages et remplace le fragment de flux du résultat du Tool en aval avant sa mise en file d’attente. Les clients de diffusion voient ainsi la valeur traitée.
Cette méthode ne se déclenche pas lorsque tool.execute() lève une erreur ; elle est appelée uniquement lors des exécutions réussies d’un Tool pour lesquelles un résultat est disponible.
processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;
ProcessToolResultArgsLien direct vers processtoolresultargs
messages:
messageList:
updateToolInvocation pour remplacer le résultat du Tool par une valeur masquée ou transformée.stepNumber:
toolName:
toolCallId:
args:
result:
tool.execute() après son passage dans ensureSerializable. Pour les Tools exécutés par le Provider, par exemple web_search d’Anthropic, il s’agit du résultat brut provenant du flux du Provider, qui ne passe pas dans ensureSerializable.providerExecuted?:
systemMessages:
steps:
state:
abort:
retry: true pour demander au LLM de réessayer l’étape en utilisant le motif d’interruption comme retour.retryCount:
tracingContext?:
requestContext?:
Cas d’utilisationLien direct vers Cas d’utilisation
- Rechercher une injection de prompt dans la sortie du Tool avant que le LLM ne la voie.
- Masquer les champs sensibles dans les retours des Tools (PII, secrets, identifiants).
- Interrompre l’exécution lorsqu’un Tool renvoie du contenu qui enfreint la politique.
- Journaliser ou instrumenter les retours des Tools à des fins de conformité ou d’audit.
Exemple : masquer les champs sensiblesLien direct vers Exemple : masquer les champs sensibles
import type { Processor } from '@mastra/core/processors'
export class RedactToolResult implements Processor {
id = 'redact-tool-result'
async processToolResult({ toolName, toolCallId, args, result, messageList }) {
if (toolName !== 'lookup-customer') return
const redacted = {
...(result as Record<string, unknown>),
ssn: '[REDACTED]',
email: '[REDACTED]',
}
messageList.updateToolInvocation({
type: 'tool-invocation',
toolInvocation: {
state: 'result',
toolCallId,
toolName,
args,
result: redacted,
},
})
}
}
Exemple : bloquer l’injection de prompt dans la sortie d’un ToolLien direct vers Exemple : bloquer l’injection de prompt dans la sortie d’un Tool
import type { Processor } from '@mastra/core/processors'
export class ScanToolResult implements Processor {
id = 'scan-tool-result'
async processToolResult({ result, abort }) {
const text = typeof result === 'string' ? result : JSON.stringify(result)
if (containsPromptInjection(text)) {
abort('blocked by scan-tool-result: suspected prompt injection')
}
}
}
function containsPromptInjection(text: string): boolean {
return /ignore (all )?(previous|prior) instructions/i.test(text)
}
Types de ProcessorsLien direct vers Types de Processors
Mastra fournit des alias de types afin de garantir que les Processors implémentent les méthodes requises :
// Must implement processInput, processInputStep, processLLMRequest, or processLLMResponse (or any combination)
type InputProcessor = Processor &
(
| { processInput: required }
| { processInputStep: required }
| { processLLMRequest: required }
| { processLLMResponse: required }
)
// Must implement processOutputStream, processOutputStep, OR processOutputResult (or any combination)
type OutputProcessor = Processor &
(
| { processOutputStream: required }
| { processOutputStep: required }
| { processOutputResult: required }
)
// Must implement processAPIError
type ErrorProcessor = Processor & { processAPIError: required }
Configurez les Processors qui implémentent processAPIError dans errorProcessors :
const agent = new Agent({
id: 'agent',
errorProcessors: [new PrefillErrorHandler()],
})
Exemples d’utilisationLien direct vers Exemples d’utilisation
Processor d’entrée de baseLien direct vers Processor d’entrée de base
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class LowercaseProcessor implements Processor {
id = 'lowercase'
async processInput({ messages }): Promise<MastraDBMessage[]> {
return messages.map(msg => ({
...msg,
content: {
...msg.content,
parts: msg.content.parts?.map(part =>
part.type === 'text' ? { ...part, text: part.text.toLowerCase() } : part,
),
},
}))
}
}
Processor par étape avec processInputStepLien direct vers per-step-processor-with-processinputstep
import type {
Processor,
ProcessInputStepArgs,
ProcessInputStepResult,
} from '@mastra/core/processors'
export class DynamicModelProcessor implements Processor {
id = 'dynamic-model'
async processInputStep({
stepNumber,
steps,
toolChoice,
}: ProcessInputStepArgs): Promise<ProcessInputStepResult> {
// Use a fast model for initial response
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
// Switch to powerful model after tool calls
if (steps.length > 0 && steps[steps.length - 1].toolCalls?.length) {
return { model: 'openai/gpt-5.6-sol' }
}
// Disable tools after 5 steps to force completion
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
return {}
}
}
Transformation de messages avec processInputStepLien direct vers message-transformer-with-processinputstep
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class ReasoningTransformer implements Processor {
id = 'reasoning-transformer'
async processInputStep({ messages, messageList }) {
// Transform reasoning parts to thinking parts at each step
// This is useful when switching between model providers
for (const msg of messages) {
if (msg.role === 'assistant' && msg.content.parts) {
for (const part of msg.content.parts) {
if (part.type === 'reasoning') {
;(part as any).type = 'thinking'
}
}
}
}
return messageList
}
}
Processor hybride (entrée et sortie)Lien direct vers Processor hybride (entrée et sortie)
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
import type { ChunkType } from '@mastra/core/stream'
export class ContentFilter implements Processor {
id = 'content-filter'
private blockedWords: string[]
constructor(blockedWords: string[]) {
this.blockedWords = blockedWords
}
async processInput({ messages, abort }): Promise<MastraDBMessage[]> {
for (const msg of messages) {
const text = msg.content.parts
?.filter(p => p.type === 'text')
.map(p => p.text)
.join(' ')
if (this.blockedWords.some(word => text?.includes(word))) {
abort('Blocked content detected in input')
}
}
return messages
}
async processOutputStream({ part, abort }): Promise<ChunkType | null> {
if (part.type === 'text-delta') {
if (this.blockedWords.some(word => part.payload.text.includes(word))) {
abort('Blocked content detected in output')
}
}
return part
}
}
Accumulation du flux avec étatLien direct vers Accumulation du flux avec état
import type { Processor } from '@mastra/core/processors'
import type { ChunkType } from '@mastra/core/stream'
export class WordCounter implements Processor {
id = 'word-counter'
async processOutputStream({ part, state }): Promise<ChunkType> {
// Initialize state on first chunk
if (!state.wordCount) {
state.wordCount = 0
}
// Count words in text chunks
if (part.type === 'text-delta') {
const words = part.payload.text.split(/\s+/).filter(Boolean)
state.wordCount += words.length
}
// Log word count on finish
if (part.type === 'finish') {
console.log(`Total words: ${state.wordCount}`)
}
return part
}
}
Cycle de vie de l’étatLien direct vers Cycle de vie de l’état
Chaque Processor reçoit un objet state dans processLLMRequest, processLLMResponse, processOutputStream, processOutputStep, processOutputResult et processAPIError. L’état possède trois propriétés importantes :
- Propre au Processor : chaque Processor reçoit son propre objet
state, indexé par l’iddu Processor. Les Processors dont les ID sont différents ne peuvent ni lire ni remplacer l’état des autres. - Propre à la requête : un nouvel objet d’état est créé au début de chaque appel à
agent.generate()ouagent.stream(). L’état ne fuit ni entre les requêtes ni entre les utilisateurs. - Partagé entre les méthodes : au sein d’une requête, le même objet
stateest transmis àprocessLLMRequest(avant l’appel au Provider),processLLMResponse(une fois l’étape terminée),processOutputStream(pour chaque fragment),processOutputStep(après chaque étape LLM),processOutputResult(une fois à la fin) etprocessAPIError(lorsqu’un appel au LLM échoue). Par exemple,processLLMRequestpeut stocker une clé de cache, puisprocessLLMResponsepeut la relire pour écrire la réponse.
Initialisez les champs de manière défensive lors du premier accès, car state commence comme un objet vide :
import type { Processor } from '@mastra/core/processors'
export class WordCounter implements Processor {
id = 'word-counter'
async processOutputStream({ part, state }) {
state.wordCount ??= 0
if (part.type === 'text-delta') {
state.wordCount += part.payload.text.split(/\s+/).filter(Boolean).length
}
return part
}
}
Interruption et fragments tripwireLien direct vers Interruption et fragments tripwire
La fonction abort de chaque méthode lève une erreur TripWire qui arrête le traitement et émet un fragment tripwire dans le flux de sortie. Les clients peuvent détecter ce fragment afin de distinguer une réponse bloquée d’une fin normale.
abort('Blocked content detected', { retry: false, metadata: { category: 'pii' } })
reason: explication lisible. Apparaît danstripwire.payload.reason.retry: lorsque cette valeur esttrue, l’Agent réessaie la même étape en utilisantreasoncomme retour. Les nouvelles tentatives ne s’exécutent que simaxProcessorRetriesest défini sur l’Agent ou l’appel ; sinon, la requête est interrompue. LorsqueerrorProcessorsest configuré,maxProcessorRetriesprend par défaut la valeur10pour cet appel.metadata: données structurées facultatives jointes au fragmenttripwirepour les consommateurs en aval.
Le fragment tripwire émis présente la forme suivante :
type TripwireChunk = {
type: 'tripwire'
runId: string
from: 'AGENT'
payload: {
reason: string
retry?: boolean
metadata?: unknown
processorId: string
}
}
Dans les appels sans diffusion (agent.generate()), le résultat expose les mêmes informations dans result.tripwire et result.finishReason === 'other'.
Émettre des fragments de données personnalisésLien direct vers Émettre des fragments de données personnalisés
Les Processors qui ont accès à writer peuvent diffuser des fragments data-* personnalisés vers le client en appelant writer.custom(chunk). Les Tools peuvent faire de même par l’intermédiaire de leur propre writer. C’est le seul moyen pour un Processor d’émettre du contenu en dehors des fragments de texte et de Tool habituels.
await writer.custom({
type: 'data-moderation',
runId,
from: 'AGENT',
data: { level: 'warn', reason: 'Possibly unsafe' },
})
Lorsque la Memory est configurée, les fragments data-* personnalisés émis depuis processOutputStream ou processOutputResult sont enregistrés comme parties du message de l’assistant. Définissez transient: true sur l’objet fragment pour le diffuser sans l’enregistrer dans la Memory :
await writer.custom({
type: 'data-progress',
data: { status: 'Processing' },
transient: true,
})
Transmettez transient comme propriété du fragment, et non comme deuxième argument de writer.custom(). Le deuxième argument contient les options du writer, telles que messageId.
Par défaut, les Processors ne voient pas les fragments data-* dans processOutputStream, afin de ne pas traiter accidentellement la télémétrie des Tools ou leur propre sortie. Activez explicitement ce comportement en définissant processDataParts: true sur le Processor :
class ModerationCollector implements Processor {
id = 'moderation-collector'
processDataParts = true
async processOutputStream({ part, state }) {
if (part.type === 'data-moderation') {
state.warnings ??= []
state.warnings.push(part.data)
}
return part
}
}
Le type du fragment doit commencer par data- pour être traité comme un fragment de données personnalisé. Renvoyer null ou undefined depuis processOutputStream supprime toujours le fragment ; un Processor peut donc examiner, modifier ou filtrer des données personnalisées de la même manière que des fragments de texte.
Configurer des Processors sur un AgentLien direct vers Configurer des Processors sur un Agent
Les Processors sont associés à un Agent au moyen de trois tableaux :
import { Agent } from '@mastra/core/agent'
import { PrefillErrorHandler } from '@mastra/core/processors'
const agent = new Agent({
id: 'support-agent',
name: 'support-agent',
model: 'openai/gpt-5',
instructions: '...',
inputProcessors: [new ContentFilter(['secret'])],
outputProcessors: [new WordCounter()],
errorProcessors: [new PrefillErrorHandler()],
maxProcessorRetries: 3,
})
inputProcessors: s’exécute avant le LLM. Reçoit les messages d’entrée.outputProcessors: s’exécute pendant ou après la réponse du LLM. Reçoit les fragments ou les messages de sortie.errorProcessors: s’exécute lorsque l’appel à l’API du LLM lève une erreur. Reçoit l’erreur brute.
Chaque tableau accepte également une fonction afin que les Processors puissent être construits pour chaque requête à partir de RequestContext :
new Agent({
id: 'processor-interface-agent',
inputProcessors: ({ requestContext }) => {
const blockedWords = requestContext.get('blockedWords') ?? []
return [new ContentFilter(blockedWords)]
},
})
Remplacements propres à chaque appelLien direct vers Remplacements propres à chaque appel
agent.generate() et agent.stream() acceptent inputProcessors, outputProcessors, errorProcessors et maxProcessorRetries. Lorsqu’un tableau de Processors est défini sur l’appel, il remplace le tableau correspondant configuré sur l’Agent pour cette requête. Les Processors de Memory, Workspace, Skill, Channel et Browser ajoutés automatiquement par Mastra sont toujours conservés et s’exécutent autour de votre tableau.
await agent.stream('Summarize this', {
outputProcessors: [new StreamFilter()],
maxProcessorRetries: 5,
})
La valeur maxProcessorRetries transmise lors de l’appel remplace la valeur par défaut de l’Agent. Si aucune n’est définie, les nouvelles tentatives demandées par les Processors sont traitées comme des interruptions.
Liens associésLien direct vers Liens associés
- Vue d’ensemble des Processors : guide conceptuel des Processors
- Garde-fous : Processors de sécurité et de validation
- Processors de Memory : Processors propres à la Memory