メインコンテンツへ移動

Processor インターフェース

Processor インターフェースは、Mastra のすべての Processor に共通する規約を定義します。Processor は1つ以上のメソッドを実装し、Agent 実行パイプラインのさまざまな段階を処理できます。

Processor メソッドの実行タイミング
Processor メソッドの実行タイミングへの直接リンク

Processor メソッドは、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 │
│ │
└────────────────────────────────────────────────────────────────────┘
メソッド実行タイミングユースケース
processInput開始時に1回、Agent ループの前最初のユーザー入力を検証/変換し、コンテキストを追加
processInputStepAgent ループの各ステップ、各 LLM 呼び出しの前ステップ間のメッセージを変換し、Tool の結果を処理
processLLMRequestLLM リクエスト変換後、Provider 呼び出しの前変更を永続化せず、現在の呼び出しで送信する LanguageModelV2Prompt を書き換え
processAPIErrorLLM API 呼び出しが失敗したときAPI による拒否を調査し、必要に応じて状態/メッセージを変更して再試行を要求
processOutputStreamLLM レスポンス中の各ストリーミングチャンクストリーミングコンテンツをフィルタリング/変更し、パターンをリアルタイムで検出
processLLMResponseLLM ステップが完了し、ストリームチャンクを収集した後完全なレスポンスを取得またはキャッシュし、processLLMRequest と対になる呼び出し後の副作用を実行
processOutputStep各 LLM レスポンスの後、Tool 実行の前出力品質を検証し、再試行を伴う Guardrail を実装
processToolResultTool ごとに、tool.execute() が返った後、結果をメッセージリストへ追加する前Tool 出力のプロンプトインジェクションを検査し、機密フィールドを墨消しし、ポリシー違反時に中止
processOutputResult生成完了後に1回最終レスポンスを後処理し、結果をログに記録

インターフェース定義
インターフェース定義への直接リンク

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
}

プロパティ
プロパティへの直接リンク

id:

string
Processor の一意な識別子。トレースとデバッグに使用します。

name?:

string
Processor の任意の表示名。指定しない場合は id を使用します。

description?:

string
トレースと Studio に表示される、任意の人が読める説明。

processorIndex?:

number
結合された Processor リスト内での位置。Processor が Memory、Workspace、呼び出しごとのオーバーライドと統合されるとき、Mastra が実行時に設定します。自分で設定する必要はありません。

processDataParts?:

boolean
true の場合、processOutputStream メソッドは Tool が writer.custom() で送出した data-* チャンクも受け取ります。デフォルトは false です。

onViolation?:

(violation: ProcessorViolation) => void | Promise<void>
方法(block または warn)にかかわらず、Processor がポリシー違反を検出したときに呼び出される任意のコールバック。アラート、外部システムへのログ記録、ユーザーへのメール送信などの副作用に使用します。このコールバックがスローしたエラーは、Processor のロジックを妨げないよう暗黙的に捕捉されます。violation オブジェクトには processorId、message、Processor 固有の detail フィールドが含まれます。

メッセージ引数
メッセージ引数への直接リンク

ほとんどの Processor メソッドは messagesmessageList の両方を受け取ります。どちらも同じ基盤の会話を指しますが、公開方法が異なります。

messagesmessageList の違い
messages-vs-messagelistへの直接リンク

  • messages:現在の段階にスコープされた MastraDBMessage オブジェクトの通常の配列。processInputprocessInputStep では System メッセージを除きます。processOutputResultprocessOutputStep では最新の LLM レスポンスを含みます。この配列は messageList を基盤とするため、メッセージの content.parts をその場で編集すると、後続の Processor と永続化処理にも反映されます。
  • messageList:実行を支える稼働中の MessageList インスタンス。フィルタリング済みのビュー(input、response、remembered、all)、複数の出力形式(db、ui、core)、会話を変更するメソッドを公開します。

現在の段階のメッセージを読み取り、map 処理し、フィールドを軽く編集するだけなら messages を使用します。次の場合は messageList を使用します。

  • 出力処理中に入力メッセージを読む場合など、別の段階のメッセージを読み取る。
  • メッセージ全体を追加、削除、または置換する。
  • サードパーティ API 用の UI メッセージや core メッセージなど、別の形式に変換する。

messages は常に messageList から導出されるため、メッセージの追加、削除、並べ替えには messageList の変更が標準的な方法です。メッセージの内容をその場で編集する場合(content.parts の書き換えなど)は、messages を直接変更しても同じです。messages から新しい配列を返すと、Mastra は現在の段階の messageList と照合します。

永続化
永続化への直接リンク

Memory が有効な場合、すべての Processor の完了後に messageList に残った内容だけがストレージへ永続化されます。永続化については、次の2つの戻り方は同等です。

  • messageList を直接変更する(または同じ MessageList インスタンスを返す)場合、記録された変更がその場で適用されるため、保存される会話に変更が反映されます。
  • MastraDBMessage[] または { messages, systemMessages } を返す場合、Mastra は返された配列を現在の段階の messageList と照合し、存在しないメッセージを削除して System メッセージを置換します。

別の MessageList インスタンスを返すとエラーになります。必ず Processor に渡されたインスタンスを変更してください。

メッセージからのテキスト読み取り
メッセージからのテキスト読み取りへの直接リンク

MastraDBMessage.content は構造化オブジェクトを使用します。文字列はサポートされません。ユーザーまたは Assistant のテキストを読み取る標準的な方法は 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
}

要点:

  • message.content.parts が主要な情報源です。1つのメッセージには、Tool 呼び出し、Tool の結果、ファイル part など、テキスト以外も含む複数の part を格納できます。part.text を読む前に part.type === 'text' でフィルタリングしてください。
  • message.content.content は後方互換性のために保持されるフラット化された文字列です。parts が空または存在しない場合のフォールバックとしてのみ使用します。
  • MastraDBMessagemessage.content 自体が通常の文字列になることはありません。従来の CoreMessage 形式では文字列の場合がありますが、Processor は常に MastraDBMessage を受け取ります。

メソッド
メソッドへの直接リンク

processInput
processinputへの直接リンク

入力メッセージを LLM に送信する前に処理します。Agent の実行開始時に1回実行されます。

processInput?(args: ProcessInputArgs): Promise<ProcessInputResult> | ProcessInputResult;

ProcessInputArgs
processinputargsへの直接リンク

messages:

MastraDBMessage[]
処理するユーザーと Assistant のメッセージ(System メッセージを除く)。

systemMessages:

CoreMessage[]
すべての System メッセージ(Agent の指示、Memory コンテキスト、ユーザー提供)。変更して返せます。

messageList:

MessageList
高度なメッセージ管理に使用する完全な MessageList インスタンス。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。実行を停止する TripWire エラーをスローします。フィードバックを付けて LLM にステップの再試行を要求するには retry: true を渡します。

retryCount:

number
この生成で Processor が再試行を発生させた回数。再試行回数の制限に使用します。Mastra から常に渡され、0から始まります。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

requestContext?:

RequestContext
threadId や resourceId などの実行メタデータを持つ、リクエストスコープのコンテキスト。

ProcessInputResult
processinputresultへの直接リンク

このメソッドは、3つの型のいずれかを返せます。

MastraDBMessage[]:

array
変換済みメッセージの配列。System メッセージは変更されません。

MessageList:

MessageList
渡されたものと同じ messageList インスタンス。直接変更したことを示します。

{ messages, systemMessages }:

object
変換済みメッセージと変更済み System メッセージの両方を持つオブジェクト。

processInputStep
processinputstepへの直接リンク

Agent ループの各ステップで、LLM に送信する前に入力メッセージを処理します。開始時に1回だけ実行される processInput と異なり、Tool 呼び出しの継続を含むすべてのステップで実行されます。

processInputStep?<TTripwireMetadata = unknown>(
args: ProcessInputStepArgs<TTripwireMetadata>,
):
| Promise<ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined>
| ProcessInputStepResult
| MessageList
| MastraDBMessage[]
| void
| undefined;

Agent ループ内の実行順序
Agent ループ内の実行順序への直接リンク

  1. processInput(開始時に 1 回)
  2. inputProcessors の processInputStep(各ステップの LLM 呼び出し前)
  3. prepareStep コールバック(processInputStep パイプラインの一部として、inputProcessors の後に実行)
  4. inputProcessors の processLLMRequest(プロンプト変換後、Provider 呼び出し前)
  5. LLM の実行
  6. outputProcessors の processOutputStream(ストリーミングの各チャンク)
  7. inputProcessors の processLLMResponse(ストリーム完了後、processLLMRequest と対になる処理)
  8. outputProcessors の processOutputStep(LLM レスポンス後、Tool 実行前)
  9. Tool の実行(必要な場合)
  10. Tool が呼び出された場合はステップ 2 から繰り返す

ProcessInputStepArgs
processinputstepargsへの直接リンク

messages:

MastraDBMessage[]
前のステップの Tool 呼び出しと結果を含むすべてのメッセージ(読み取り専用のスナップショット)。

messageList:

MessageList
メッセージ管理用の MessageList インスタンス。直接変更するか、結果で返せます。

stepNumber:

number
現在のステップ番号(0始まり)。ステップ0は最初の LLM 呼び出しです。

steps:

StepResult[]
text、toolCalls、toolResults を含む、前のステップの結果。

systemMessages:

CoreMessage[]
すべての System メッセージ(読み取り専用のスナップショット)。置換するには結果で返します。

model:

MastraLanguageModelV2
現在使用中のモデル。切り替えるには結果で別のモデルを返します。

toolChoice?:

ToolChoice
現在の Tool 選択設定('auto'、'none'、'required'、または特定の Tool)。

activeTools?:

string[]
現在有効な Tool 名。Tool を制限するにはフィルタリング済みの配列を返します。

tools?:

ToolSet
このステップで現在利用できる Tool。Tool を追加/置換するには結果で返します。

providerOptions?:

SharedV2ProviderOptions
Provider 固有のオプション(Anthropic の cacheControl、OpenAI の reasoningEffort など)。

modelSettings?:

CallSettings
temperature、maxTokens、topP などのモデル設定。

structuredOutput?:

StructuredOutputOptions
構造化出力の設定(schema、output mode)。変更するには結果で返します。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。実行を停止する TripWire エラーをスローします。フィードバックを付けて LLM にステップの再試行を要求するには retry: true を渡します。

retryCount:

number
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

requestContext?:

RequestContext
実行メタデータを持つ、リクエストスコープのコンテキスト。

ProcessInputStepResult
processinputstepresultへの直接リンク

processInputStep は複数の形式を返せます。

  • ProcessInputStepResult オブジェクト:このステップについて、以下で説明するプロパティを任意の組み合わせでオーバーライドします。
  • MessageList:同じ messageList インスタンスを返し、メッセージをその場で変更したことを示します。
  • MastraDBMessage[]:変換済みのメッセージ配列を返し、ステップのメッセージを置換します。
  • void または undefined:何も返さず、ステップを変更しません。

オブジェクト形式では、次のプロパティを任意の組み合わせで返せます。

model?:

LanguageModelV2 | string
このステップのモデルを変更します。モデルインスタンス、または 'openai/gpt-5.5' のような router ID を指定できます。

toolChoice?:

ToolChoice
このステップの Tool 選択動作を変更します。

activeTools?:

string[]
このステップで利用できる Tool をフィルタリングします。

tools?:

ToolSet
このステップの Tool を置換または変更します。統合には spread を使用します:{ tools: { ...tools, newTool } }。

messages?:

MastraDBMessage[]
すべてのメッセージを置換します。messageList と併用できません。

messageList?:

MessageList
同じ messageList インスタンスを返します(変更したことを示します)。messages と併用できません。

systemMessages?:

CoreMessage[]
このステップについてのみ、すべての System メッセージを置換します。

providerOptions?:

SharedV2ProviderOptions
このステップの Provider 固有オプションを変更します。

modelSettings?:

CallSettings
このステップのモデル設定を変更します。

structuredOutput?:

StructuredOutputOptions
このステップの構造化出力設定を変更します。

Processor のチェーン
Processor のチェーンへの直接リンク

複数の Processor が processInputStep を実装している場合は順番に実行され、変更が後続へ引き継がれます。

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'

System メッセージの分離
System メッセージの分離への直接リンク

System メッセージは各ステップの開始時に元の値へリセットされます。processInputStep で加えた変更は現在のステップだけに影響し、後続のステップには影響しません。

ユースケース
ユースケースへの直接リンク

  • ステップ番号またはコンテキストに基づく動的なモデル切り替え
  • 一定のステップ数を超えた後の Tool 無効化
  • 会話コンテキストに基づく Tool の動的な追加または置換
  • Provider 間でのメッセージ part 型の変換(Anthropic 向けの reasoningthinking など)
  • ステップ番号または蓄積されたコンテキストに基づくメッセージ変更
  • ステップ固有の System 指示の追加
  • ステップごとの Provider オプションの調整(キャッシュ制御など)
  • ステップコンテキストに基づく構造化出力スキーマの変更

processLLMRequest
processllmrequestへの直接リンク

Mastra が MessageListLanguageModelV2Prompt に変換した後、Provider を呼び出す前に、最終的な LLM リクエストを処理します。現在の送信リクエストだけに影響させる、一時的かつモデルを考慮した書き換えに使用します。

返されたプロンプトの変更は、現在の呼び出しについてのみモデルへ転送されます。MessageList、Memory、UI 履歴、後続の Provider 呼び出しには永続化されません。

processLLMRequest?(
args: ProcessLLMRequestArgs,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult;

ProcessLLMRequestArgs
processllmrequestargsへの直接リンク

prompt:

LanguageModelV2Prompt
この呼び出しで Provider に送信する LLM リクエストプロンプト。

model:

MastraLanguageModel
プロンプトを受け取る解決済みモデル。Provider 固有の書き換えを限定するために使用します。

stepNumber:

number
現在のステップ番号(0始まり)。ステップ0は最初の LLM 呼び出しです。

steps:

StepResult[]
text、toolCalls、toolResults を含む、前のステップの結果。

state:

Record<string, unknown>
このリクエスト内のすべてのメソッド呼び出しで維持される、Processor ごとの state。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。実行を停止する TripWire エラーをスローし、tripwire チャンクを送出します。

retryCount:

number
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。

requestContext?:

RequestContext
実行メタデータを持つ、リクエストスコープのコンテキスト。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

writer?:

ProcessorStreamWriter
ストリーミング中にカスタムデータチャンクを送出する Stream writer。data-* チャンクを送出するには writer.custom() を呼び出します。

abortSignal?:

AbortSignal
操作をキャンセルするための Signal。

戻り値
戻り値への直接リンク

processLLMRequest は、{ prompt?: LanguageModelV2Prompt } | undefined | void である ProcessLLMRequestResult を返します。

  • 現在の Provider 呼び出しで送信するプロンプトを置換するには、{ prompt } を返します。
  • 元のプロンプトを変更せず転送するには、undefined または void を返します。

ユースケース
ユースケースへの直接リンク

  • モデル呼び出し前に Provider 固有のプロンプト part を削除または整形する
  • Provider の入力要件に合わせてロールまたはコンテンツを正規化する
  • ループ中に Provider を切り替えるとき、Tool の結果形式を適応させる

processLLMResponse
processllmresponseへの直接リンク

ステップの完了後(またはキャッシュ済みレスポンスの再生後)、output processor がレスポンスチャンクを収集した後に、LLM レスポンスを処理します。このフックは processLLMRequest と対になります。Provider 呼び出し前に processLLMRequest でキャッシュキーなどの state を保存し、完了したレスポンスへの処理(キャッシュへの書き込みなど)を processLLMResponse で行います。

state オブジェクトは、同じステップの processLLMRequest に渡されたものと同じインスタンスです。そのため、Processor は呼び出し前後の処理を関連付けられます。

processLLMResponse?(
args: ProcessLLMResponseArgs,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult;

ProcessLLMResponseArgs
processllmresponseargsへの直接リンク

chunks:

CachedLLMStepChunk[]
このステップで LLM 呼び出しが生成した(またはキャッシュから再生した)チャンク。簡略形式({ type, payload })です。

model:

MastraLanguageModel
レスポンスを生成した(または生成するはずだった)モデル。

stepNumber:

number
現在のステップ番号(0始まり)。

steps:

StepResult[]
このステップを含む、これまでに完了したすべてのステップ。

state:

Record<string, unknown>
同じステップの processLLMRequest と共有する、Processor ごとの state。2つのフック間でキャッシュキーなどのデータを渡すために使用します。

fromCache:

boolean
true の場合、processLLMRequest{ response } を返し、レスポンスがキャッシュから再生されたことを示します。キャッシュへ書き込む Processor は、これが true のとき書き込みをスキップしてください。

warnings?:

LanguageModelV2CallWarning[]
言語モデル呼び出しから報告された警告(未対応の設定など)。

request?:

unknown
利用できる場合の Provider リクエスト本文。トレースに役立ちます。

rawResponse?:

unknown
利用できる場合の生の Provider レスポンス。トレースに役立ちます。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。実行を停止する TripWire エラーをスローします。

retryCount:

number
現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。

requestContext?:

RequestContext
実行メタデータを含むリクエストスコープのコンテキスト。

tracingContext?:

TracingContext
Observability のための Trace コンテキスト。

writer?:

ProcessorStreamWriter
カスタムデータチャンクを送出する Stream writer。

abortSignal?:

AbortSignal
操作をキャンセルするための Signal。

戻り値
戻り値への直接リンク

processLLMResponse は、undefined | void である ProcessLLMResponseResult を返します。戻り値は将来の拡張用に予約されています。

ユースケース
ユースケースへの直接リンク

  • ライブ呼び出し後に LLM レスポンスをキャッシュへ書き込む(processLLMRequest でのキャッシュキー導出と組み合わせる)
  • 分析用に完全なレスポンスをログまたは記録する
  • 完了したレスポンスに基づいて副作用を発生させる

processAPIError
processapierrorへの直接リンク

LLM API による拒否エラーが最終エラーとして表面化する前に処理します。API 呼び出しが再試行不能なエラー(ステータスコード 400 や 422 など)で失敗したときに実行されます。成功したレスポンスの後に実行される processOutputStep と異なり、API がリクエストを拒否したときに実行されます。

processAPIError を実装する Processor は、Agent の errorProcessors 配列に追加します。

Processor はエラーを調査し、messageList へのメッセージ追加などによってリクエストを変更できます。変更後の状態で再試行するには { retry: true } を返します。

processAPIError?(args: ProcessAPIErrorArgs): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void;

ProcessAPIErrorArgs
processapierrorargsへの直接リンク

error:

unknown
LLM API 呼び出し中に発生したエラー。

messages:

MastraDBMessage[]
エラー発生時点のすべてのメッセージ。

messageList:

MessageList
メッセージ管理用の MessageList インスタンス。再試行前にリクエストを変えるには、これを変更します。

stepNumber:

number
現在のステップ番号(0始まり)。

steps:

StepResult[]
これまでに完了したすべてのステップ。

state:

Record<string, unknown>
このリクエスト内のすべてのメソッド呼び出しで維持される、Processor ごとの state。

retryCount:

number
エラーハンドラーの現在の再試行回数。再試行回数の制限に使用します。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。

writer?:

ProcessorStreamWriter
ストリーミング中にカスタムデータチャンクを送出する Stream writer。data-* チャンクを送出するには writer.custom() を呼び出します。

requestContext?:

RequestContext
Agent 呼び出しから渡された Request コンテキスト。

abortSignal?:

AbortSignal
操作をキャンセルするための Signal。

ProcessAPIErrorResult
processapierrorresultへの直接リンク

retry:

boolean
変更の適用後に LLM 呼び出しを再試行するかどうか。

ユースケース
ユースケースへの直接リンク

  • リクエストを変更して再試行し、API 固有の拒否を処理する
  • リクエストを変更し、再試行不能なエラーを再試行可能にする
  • モデル固有のエラー復旧方法を実装する

例:カスタムエラー復旧
例:カスタムエラー復旧への直接リンク

src/mastra/processors/error-recovery.ts
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 }
}
}
}
}

processOutputStream
processoutputstreamへの直接リンク

組み込みの state 管理を使用して、ストリーミング出力チャンクを処理します。Processor はチャンクを蓄積し、より広いコンテキストに基づいて判断できます。

processOutputStream?(args: ProcessOutputStreamArgs): Promise<ChunkType | null | undefined>;

ProcessOutputStreamArgs
processoutputstreamargsへの直接リンク

part:

ChunkType
現在処理中のストリームチャンク。

streamParts:

ChunkType[]
ストリームでこれまでに確認されたすべてのチャンク。

state:

Record<string, unknown>
1つのリクエスト内のすべてのチャンクとメソッド呼び出しで維持される、変更可能な Processor ごとの state。generate または stream の新しい呼び出しごとに新しい state オブジェクトが作成されます。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
ストリームを中止する関数。ストリームを終了して tripwire チャンクを送出する TripWire エラーをスローします。終了せず LLM の再試行を要求するには retry: true を渡します。

retryCount:

number
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。

messageList?:

MessageList
会話履歴へアクセスするための MessageList インスタンス。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

requestContext?:

RequestContext
実行メタデータを持つ、リクエストスコープのコンテキスト。

writer?:

ProcessorStreamWriter
カスタムデータチャンクをクライアントへ送出する Stream writer。data-* 型のチャンクを送出するには writer.custom() を呼び出します。ストリーミング中に利用できます。

戻り値
戻り値への直接リンク

processOutputStreamPromise<ChunkType | null | undefined> を返します。

  • チャンクを送出するには ChunkType を返します。変更せず送出するには元の part を、変更したチャンクを送出するには新しい ChunkType を返します。
  • チャンクを破棄するには null を返します。後続の Processor やクライアントには何も送信されません。
  • チャンクを破棄するには undefined も返せます(return; 文やメソッド末尾への到達による暗黙的な undefined を含む)。nullundefined の動作は同じです。

チャンクの破棄は、その1つのチャンクだけに影響します。ストリームは継続し、次のチャンクも処理されます。ストリーム全体を停止するには abort() を呼び出します。


processOutputResult
processoutputresultへの直接リンク

ストリーミングまたは生成の完了後に、完全な出力結果を処理します。

processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;

ProcessOutputResultArgs
processoutputresultargsへの直接リンク

messages:

MastraDBMessage[]
生成されたレスポンスメッセージ。

messageList:

MessageList
メッセージ管理用の MessageList インスタンス。

state:

Record<string, unknown>
このリクエスト内のすべてのメソッド呼び出しで維持される、Processor ごとの state。processOutputStream や他のメソッドと共有されます。

result:

OutputResult
text(蓄積テキスト)、usage(inputTokens、outputTokens、totalTokens を含むトークン使用量)、finishReason(生成が終了した理由)、steps(toolCalls、toolResults、reasoning、sources、files などを含む各 LLM ステップの全結果)を持つ、解決済みの生成結果。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。実行を停止する TripWire エラーをスローし、tripwire チャンクを送出します。

retryCount:

number
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

requestContext?:

RequestContext
実行メタデータを持つ、リクエストスコープのコンテキスト。

writer?:

ProcessorStreamWriter
カスタムデータチャンクをクライアントへ送出する Stream writer。data-* 型のチャンクを送出するには writer.custom() を呼び出します。ストリーミング中に利用できます。

processOutputStep
processoutputstepへの直接リンク

Agent ループの各 LLM レスポンスの後、Tool を実行する前に出力を処理します。最後に1回実行される processOutputResult と異なり、すべてのステップで実行されます。再試行を発生させられる Guardrail の実装に最適なメソッドです。

processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;

ProcessOutputStepArgs
processoutputstepargsへの直接リンク

messages:

MastraDBMessage[]
最新の LLM レスポンスを含むすべてのメッセージ。

messageList:

MessageList
メッセージ管理用の MessageList インスタンス。

stepNumber:

number
現在のステップ番号(0始まり)。

finishReason?:

string
LLM から返された終了理由(stop、tool-use、length など)。

providerMetadata?:

ProviderMetadata
終了ステップの Provider 固有メタデータ(AWS Bedrock の Guardrail トレースなど)。steps が空になるコンテンツフィルターによるブロックを含め、モデルステップが Provider メタデータを生成した場合に存在します。

toolCalls?:

ToolCallInfo[]
このステップで行われた Tool 呼び出し(存在する場合)。

text?:

string
このステップで生成されたテキスト。

usage:

LanguageModelUsage
現在のステップのトークン使用量(inputTokensoutputTokenstotalTokens)。

systemMessages:

CoreMessage[]
読み取り/変更用のすべての System メッセージ。

steps:

StepResult[]
現在のステップを含む、これまでに完了したすべてのステップ。

state:

Record<string, unknown>
このリクエスト内のすべてのメソッド呼び出しで維持される、Processor ごとの state。processOutputStream および processOutputResult と共有されます。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
処理を中止する関数。LLM にステップの再試行を要求するには retry: true を渡します。

retryCount:

number
Processor が再試行を発生させた回数。再試行回数の制限に使用します。Mastra から常に渡され、0から始まります。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

requestContext?:

RequestContext
実行メタデータを持つ、リクエストスコープのコンテキスト。

ユースケース
ユースケースへの直接リンク

  • 再試行を要求できる品質 Guardrail の実装
  • Tool 実行前の LLM 出力の検証
  • ステップごとのログまたはメトリクスの追加
  • 再試行機能を持つ出力モデレーションの実装

例:再試行を伴う品質 Guardrail
例:再試行を伴う品質 Guardrailへの直接リンク

src/mastra/processors/quality-guardrail.ts
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 []
}
}

processToolResult
processtoolresultへの直接リンク

tool.execute() が返った後、その結果をメッセージリストへ追加する前、または次の LLM 呼び出しへ渡す前に、Tool の結果を処理します。Tool 実行前に発生する processOutputStep と対称です。Tool 出力のプロンプトインジェクション検査、機密フィールドの墨消し、abort('reason', { retry: true }) による実行中止に使用します。

Tool の結果を置換するには、messageList.updateToolInvocationmessageList をその場で変更します。ランタイムは Processor 処理後の結果をメッセージリストから再び読み取り、後続の Tool 結果ストリームチャンクをキューへ追加する前に上書きします。そのため、ストリーミングクライアントには処理済みの値が表示されます。

tool.execute() が例外をスローした場合、このメソッドは発生しません。結果を利用できる、成功した Tool 実行についてのみ呼び出されます。

processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;

ProcessToolResultArgs
processtoolresultargsへの直接リンク

messages:

MastraDBMessage[]
Tool 呼び出しを持つ現在の Assistant メッセージを含む、すべてのメッセージ。

messageList:

MessageList
メッセージ管理用の MessageList インスタンス。Tool の結果を墨消しまたは変換した値に置換するには updateToolInvocation を呼び出します。

stepNumber:

number
現在のステップ番号(0始まり)。

toolName:

string
実行された Tool の名前。

toolCallId:

string
この特定の Tool 呼び出しの一意な識別子。

args:

unknown
LLM が Tool に渡した引数。

result:

unknown
Tool が返した値。クライアント実行型 Tool では、ensureSerializable を通過した後の tool.execute() の出力です。Provider 実行型 Tool(Anthropic の web_search など)では、ensureSerializable を通らない Provider ストリームの生の結果です。

providerExecuted?:

boolean
この結果が Anthropic web_search などの Provider 実行型 Tool から取得されたかどうか。クライアント実行型 Tool のデフォルトは undefined です。

systemMessages:

CoreMessage[]
読み取り用のすべての System メッセージ。

steps:

StepResult[]
これまでに完了したすべてのステップ。

state:

Record<string, unknown>
このリクエスト内のすべてのメソッド呼び出しで維持される、Processor ごとの state。同じ Processor の他の Processor メソッドと共有されます。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
実行を中止する関数。中止理由をフィードバックとして LLM にステップの再試行を要求するには retry: true を渡します。

retryCount:

number
Processor が再試行を発生させた回数。0から始まります。

tracingContext?:

TracingContext
可観測性のためのトレースコンテキスト。

requestContext?:

RequestContext
実行メタデータを持つ、リクエストスコープのコンテキスト。

ユースケース
ユースケースへの直接リンク

  • LLM が見る前に、Tool 出力のプロンプトインジェクションを検査する。
  • Tool の戻り値から機密フィールド(PII、シークレット、認証情報)を墨消しする。
  • Tool がポリシー違反のコンテンツを返したときに実行を中止する。
  • コンプライアンスまたは監査用に、Tool の戻り値をログ記録または計測する。

例:機密フィールドの墨消し
例:機密フィールドの墨消しへの直接リンク

src/mastra/processors/redact-tool-result.ts
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,
},
})
}
}

例:Tool 出力内のプロンプトインジェクションをブロック
例:Tool 出力内のプロンプトインジェクションをブロックへの直接リンク

src/mastra/processors/scan-tool-result.ts
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)
}

Processor の型
Processor の型への直接リンク

Mastra は、Processor が必須メソッドを実装することを保証する型エイリアスを提供します。

// 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 }

processAPIError を実装する Processor は errorProcessors に設定します。

const agent = new Agent({
id: 'agent',
errorProcessors: [new PrefillErrorHandler()],
})

使用例
使用例への直接リンク

基本的な input processor
基本的な input processorへの直接リンク

src/mastra/processors/lowercase.ts
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,
),
},
}))
}
}

processInputStep を使用するステップ単位の Processor
per-step-processor-with-processinputstepへの直接リンク

src/mastra/processors/dynamic-model.ts
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 {}
}
}

processInputStep を使用するメッセージ変換
message-transformer-with-processinputstepへの直接リンク

src/mastra/processors/reasoning-transformer.ts
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(入力と出力)
ハイブリッド Processor(入力と出力)への直接リンク

src/mastra/processors/content-filter.ts
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
}
}

state を使用するストリーム蓄積
state を使用するストリーム蓄積への直接リンク

src/mastra/processors/word-counter.ts
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
}
}

state のライフサイクル
state のライフサイクルへの直接リンク

すべての Processor は、processLLMRequestprocessLLMResponseprocessOutputStreamprocessOutputStepprocessOutputResultprocessAPIErrorstate オブジェクトを受け取ります。state には3つの重要な特性があります。

  • Processor ごと:各 Processor は、その id をキーとする独自の state オブジェクトを受け取ります。ID が異なる Processor 同士は、互いの state を読み取ったり上書きしたりできません。
  • リクエストごとagent.generate() または agent.stream() の各呼び出し開始時に、新しい state オブジェクトが作成されます。state がリクエスト間やユーザー間で漏れることはありません。
  • メソッド間で共有:1つのリクエスト内では、同じ state オブジェクトが processLLMRequest(Provider 呼び出し前)、processLLMResponse(ステップ完了後)、processOutputStream(各チャンク)、processOutputStep(各 LLM ステップ後)、processOutputResult(最後に1回)、processAPIError(LLM 呼び出し失敗時)へ渡されます。たとえば、processLLMRequest でキャッシュキーを保存し、processLLMResponse で読み戻してレスポンスを書き込めます。

state は空のオブジェクトとして始まるため、最初にアクセスするときはフィールドを防御的に初期化してください。

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

中止と tripwire チャンク
中止と tripwire チャンクへの直接リンク

各メソッドの abort 関数は、処理を停止する TripWire エラーをスローし、出力ストリームに tripwire チャンクを送出します。クライアントはこのチャンクを検出し、ブロックされたレスポンスと通常の完了を区別できます。

abort('Blocked content detected', { retry: false, metadata: { category: 'pii' } })
  • reason:人が読める説明。tripwire.payload.reason として表示されます。
  • retrytrue の場合、Agent は reason をフィードバックとして渡し、同じステップを再試行します。再試行は Agent または呼び出しに maxProcessorRetries が設定されている場合に限り実行されます。それ以外の場合、リクエストは中止されます。errorProcessors が設定されている場合、その呼び出しの maxProcessorRetries はデフォルトで 10 です。
  • metadata:後続のコンシューマー向けに tripwire チャンクへ添付する任意の構造化データ。

送出される tripwire チャンクの形式は次のとおりです。

type TripwireChunk = {
type: 'tripwire'
runId: string
from: 'AGENT'
payload: {
reason: string
retry?: boolean
metadata?: unknown
processorId: string
}
}

非ストリーミング呼び出し(agent.generate())では、結果の result.tripwireresult.finishReason === 'other' から同じ情報を取得できます。

カスタムデータチャンクの送出
カスタムデータチャンクの送出への直接リンク

writer にアクセスできる Processor は、writer.custom(chunk) を呼び出して、カスタム data-* チャンクをクライアントへストリーミングできます。Tool も独自の writer で同じことができます。Processor が通常のテキストチャンクと Tool チャンク以外のコンテンツを送出する方法は、これだけです。

await writer.custom({
type: 'data-moderation',
runId,
from: 'AGENT',
data: { level: 'warn', reason: 'Possibly unsafe' },
})

Memory が設定されている場合、processOutputStream または processOutputResult から送出されたカスタム data-* チャンクは、Assistant メッセージの part として保存されます。Memory に保存せずストリーミングするには、チャンクオブジェクトに transient: true を設定します。

await writer.custom({
type: 'data-progress',
data: { status: 'Processing' },
transient: true,
})

transientwriter.custom() の第2引数ではなく、チャンクのプロパティとして渡します。第2引数には messageId などの writer オプションを指定します。

デフォルトでは、Tool のテレメトリーや自身の出力を誤って処理しないよう、Processor は processOutputStreamdata-* チャンクを受け取りません。Processor に processDataParts: true を設定すると受け取れます。

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

カスタムデータチャンクとして扱うには、チャンクの typedata- で始まる必要があります。processOutputStream から null または undefined を返すと、引き続きそのチャンクは破棄されます。そのため、Processor はテキストチャンクと同様にカスタムデータを検査、変更、フィルタリングできます。

Agent への Processor の設定
Agent への Processor の設定への直接リンク

Processor は3つの配列を通じて Agent に設定します。

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:LLM の前に実行され、入力メッセージを受け取ります。
  • outputProcessors:LLM レスポンスの処理中または処理後に実行され、出力チャンクまたはメッセージを受け取ります。
  • errorProcessors:LLM API 呼び出しが例外をスローしたときに実行され、生のエラーを受け取ります。

各配列は関数も受け取るため、RequestContext からリクエストごとに Processor を構築できます。

new Agent({
id: 'processor-interface-agent',
inputProcessors: ({ requestContext }) => {
const blockedWords = requestContext.get('blockedWords') ?? []
return [new ContentFilter(blockedWords)]
},
})

呼び出しごとのオーバーライド
呼び出しごとのオーバーライドへの直接リンク

agent.generate()agent.stream() は、inputProcessorsoutputProcessorserrorProcessorsmaxProcessorRetries を受け取ります。呼び出しにいずれかの Processor 配列を設定すると、そのリクエストについて Agent に設定された対応する配列を置き換えます。Mastra が自動的に追加する Memory、Workspace、Skill、Channel、Browser の Processor は常に保持され、指定した配列の前後で実行されます。

await agent.stream('Summarize this', {
outputProcessors: [new StreamFilter()],
maxProcessorRetries: 5,
})

呼び出しで渡した maxProcessorRetries は Agent のデフォルト値をオーバーライドします。どちらにも設定されていない場合、Processor が要求した再試行は中止として扱われます。