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 ループの前 | 最初のユーザー入力を検証/変換し、コンテキストを追加 |
processInputStep | Agent ループの各ステップ、各 LLM 呼び出しの前 | ステップ間のメッセージを変換し、Tool の結果を処理 |
processLLMRequest | LLM リクエスト変換後、Provider 呼び出しの前 | 変更を永続化せず、現在の呼び出しで送信する LanguageModelV2Prompt を書き換え |
processAPIError | LLM API 呼び出しが失敗したとき | API による拒否を調査し、必要に応じて状態/メッセージを変更して再試行を要求 |
processOutputStream | LLM レスポンス中の各ストリーミングチャンク | ストリーミングコンテンツをフィルタリング/変更し、パターンをリアルタイムで検出 |
processLLMResponse | LLM ステップが完了し、ストリームチャンクを収集した後 | 完全なレスポンスを取得またはキャッシュし、processLLMRequest と対になる呼び出し後の副作用を実行 |
processOutputStep | 各 LLM レスポンスの後、Tool 実行の前 | 出力品質を検証し、再試行を伴う Guardrail を実装 |
processToolResult | Tool ごとに、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:
name?:
description?:
processorIndex?:
processDataParts?:
data-* チャンクも受け取ります。デフォルトは false です。onViolation?:
メッセージ引数メッセージ引数への直接リンク
ほとんどの Processor メソッドは messages と messageList の両方を受け取ります。どちらも同じ基盤の会話を指しますが、公開方法が異なります。
messages と messageList の違いmessages-vs-messagelistへの直接リンク
messages:現在の段階にスコープされたMastraDBMessageオブジェクトの通常の配列。processInputとprocessInputStepでは System メッセージを除きます。processOutputResultとprocessOutputStepでは最新の 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が空または存在しない場合のフォールバックとしてのみ使用します。MastraDBMessageのmessage.content自体が通常の文字列になることはありません。従来のCoreMessage形式では文字列の場合がありますが、Processor は常にMastraDBMessageを受け取ります。
メソッドメソッドへの直接リンク
processInputprocessinputへの直接リンク
入力メッセージを LLM に送信する前に処理します。Agent の実行開始時に1回実行されます。
processInput?(args: ProcessInputArgs): Promise<ProcessInputResult> | ProcessInputResult;
ProcessInputArgsprocessinputargsへの直接リンク
messages:
systemMessages:
messageList:
abort:
retry: true を渡します。retryCount:
tracingContext?:
requestContext?:
ProcessInputResultprocessinputresultへの直接リンク
このメソッドは、3つの型のいずれかを返せます。
MastraDBMessage[]:
MessageList:
{ messages, systemMessages }:
processInputStepprocessinputstepへの直接リンク
Agent ループの各ステップで、LLM に送信する前に入力メッセージを処理します。開始時に1回だけ実行される processInput と異なり、Tool 呼び出しの継続を含むすべてのステップで実行されます。
processInputStep?<TTripwireMetadata = unknown>(
args: ProcessInputStepArgs<TTripwireMetadata>,
):
| Promise<ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined>
| ProcessInputStepResult
| MessageList
| MastraDBMessage[]
| void
| undefined;
Agent ループ内の実行順序Agent ループ内の実行順序への直接リンク
processInput(開始時に 1 回)- inputProcessors の
processInputStep(各ステップの LLM 呼び出し前) prepareStepコールバック(processInputStep パイプラインの一部として、inputProcessors の後に実行)- inputProcessors の
processLLMRequest(プロンプト変換後、Provider 呼び出し前) - LLM の実行
- outputProcessors の
processOutputStream(ストリーミングの各チャンク) - inputProcessors の
processLLMResponse(ストリーム完了後、processLLMRequestと対になる処理) - outputProcessors の
processOutputStep(LLM レスポンス後、Tool 実行前) - Tool の実行(必要な場合)
- Tool が呼び出された場合はステップ 2 から繰り返す
ProcessInputStepArgsprocessinputstepargsへの直接リンク
messages:
messageList:
stepNumber:
steps:
systemMessages:
model:
toolChoice?:
activeTools?:
tools?:
providerOptions?:
modelSettings?:
structuredOutput?:
abort:
retry: true を渡します。retryCount:
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。tracingContext?:
requestContext?:
ProcessInputStepResultprocessinputstepresultへの直接リンク
processInputStep は複数の形式を返せます。
ProcessInputStepResultオブジェクト:このステップについて、以下で説明するプロパティを任意の組み合わせでオーバーライドします。MessageList:同じmessageListインスタンスを返し、メッセージをその場で変更したことを示します。MastraDBMessage[]:変換済みのメッセージ配列を返し、ステップのメッセージを置換します。voidまたはundefined:何も返さず、ステップを変更しません。
オブジェクト形式では、次のプロパティを任意の組み合わせで返せます。
model?:
toolChoice?:
activeTools?:
tools?:
messages?:
messageList?:
systemMessages?:
providerOptions?:
modelSettings?:
structuredOutput?:
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 向けの
reasoning→thinkingなど) - ステップ番号または蓄積されたコンテキストに基づくメッセージ変更
- ステップ固有の System 指示の追加
- ステップごとの Provider オプションの調整(キャッシュ制御など)
- ステップコンテキストに基づく構造化出力スキーマの変更
processLLMRequestprocessllmrequestへの直接リンク
Mastra が MessageList を LanguageModelV2Prompt に変換した後、Provider を呼び出す前に、最終的な LLM リクエストを処理します。現在の送信リクエストだけに影響させる、一時的かつモデルを考慮した書き換えに使用します。
返されたプロンプトの変更は、現在の呼び出しについてのみモデルへ転送されます。MessageList、Memory、UI 履歴、後続の Provider 呼び出しには永続化されません。
processLLMRequest?(
args: ProcessLLMRequestArgs,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult;
ProcessLLMRequestArgsprocessllmrequestargsへの直接リンク
prompt:
model:
stepNumber:
steps:
state:
abort:
tripwire チャンクを送出します。retryCount:
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。requestContext?:
tracingContext?:
writer?:
data-* チャンクを送出するには writer.custom() を呼び出します。abortSignal?:
戻り値戻り値への直接リンク
processLLMRequest は、{ prompt?: LanguageModelV2Prompt } | undefined | void である ProcessLLMRequestResult を返します。
- 現在の Provider 呼び出しで送信するプロンプトを置換するには、
{ prompt }を返します。 - 元のプロンプトを変更せず転送するには、
undefinedまたはvoidを返します。
ユースケースユースケースへの直接リンク
- モデル呼び出し前に Provider 固有のプロンプト part を削除または整形する
- Provider の入力要件に合わせてロールまたはコンテンツを正規化する
- ループ中に Provider を切り替えるとき、Tool の結果形式を適応させる
processLLMResponseprocessllmresponseへの直接リンク
ステップの完了後(またはキャッシュ済みレスポンスの再生後)、output processor がレスポンスチャンクを収集した後に、LLM レスポンスを処理します。このフックは processLLMRequest と対になります。Provider 呼び出し前に processLLMRequest でキャッシュキーなどの state を保存し、完了したレスポンスへの処理(キャッシュへの書き込みなど)を processLLMResponse で行います。
state オブジェクトは、同じステップの processLLMRequest に渡されたものと同じインスタンスです。そのため、Processor は呼び出し前後の処理を関連付けられます。
processLLMResponse?(
args: ProcessLLMResponseArgs,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult;
ProcessLLMResponseArgsprocessllmresponseargsへの直接リンク
chunks:
{ type, payload })です。model:
stepNumber:
steps:
state:
processLLMRequest と共有する、Processor ごとの state。2つのフック間でキャッシュキーなどのデータを渡すために使用します。fromCache:
true の場合、processLLMRequest が { response } を返し、レスポンスがキャッシュから再生されたことを示します。キャッシュへ書き込む Processor は、これが true のとき書き込みをスキップしてください。warnings?:
request?:
rawResponse?:
abort:
retryCount:
0 から始まり、Processor が発生させる再試行の上限設定に使用します。requestContext?:
tracingContext?:
writer?:
abortSignal?:
戻り値戻り値への直接リンク
processLLMResponse は、undefined | void である ProcessLLMResponseResult を返します。戻り値は将来の拡張用に予約されています。
ユースケースユースケースへの直接リンク
- ライブ呼び出し後に LLM レスポンスをキャッシュへ書き込む(
processLLMRequestでのキャッシュキー導出と組み合わせる) - 分析用に完全なレスポンスをログまたは記録する
- 完了したレスポンスに基づいて副作用を発生させる
processAPIErrorprocessapierrorへの直接リンク
LLM API による拒否エラーが最終エラーとして表面化する前に処理します。API 呼び出しが再試行不能なエラー(ステータスコード 400 や 422 など)で失敗したときに実行されます。成功したレスポンスの後に実行される processOutputStep と異なり、API がリクエストを拒否したときに実行されます。
processAPIError を実装する Processor は、Agent の errorProcessors 配列に追加します。
Processor はエラーを調査し、messageList へのメッセージ追加などによってリクエストを変更できます。変更後の状態で再試行するには { retry: true } を返します。
processAPIError?(args: ProcessAPIErrorArgs): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void;
ProcessAPIErrorArgsprocessapierrorargsへの直接リンク
error:
messages:
messageList:
stepNumber:
steps:
state:
retryCount:
abort:
writer?:
data-* チャンクを送出するには writer.custom() を呼び出します。requestContext?:
abortSignal?:
ProcessAPIErrorResultprocessapierrorresultへの直接リンク
retry:
ユースケースユースケースへの直接リンク
- リクエストを変更して再試行し、API 固有の拒否を処理する
- リクエストを変更し、再試行不能なエラーを再試行可能にする
- モデル固有のエラー復旧方法を実装する
例:カスタムエラー復旧例:カスタムエラー復旧への直接リンク
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 }
}
}
}
}
processOutputStreamprocessoutputstreamへの直接リンク
組み込みの state 管理を使用して、ストリーミング出力チャンクを処理します。Processor はチャンクを蓄積し、より広いコンテキストに基づいて判断できます。
processOutputStream?(args: ProcessOutputStreamArgs): Promise<ChunkType | null | undefined>;
ProcessOutputStreamArgsprocessoutputstreamargsへの直接リンク
part:
streamParts:
state:
abort:
tripwire チャンクを送出する TripWire エラーをスローします。終了せず LLM の再試行を要求するには retry: true を渡します。retryCount:
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。messageList?:
tracingContext?:
requestContext?:
writer?:
戻り値戻り値への直接リンク
processOutputStream は Promise<ChunkType | null | undefined> を返します。
- チャンクを送出するには
ChunkTypeを返します。変更せず送出するには元のpartを、変更したチャンクを送出するには新しいChunkTypeを返します。 - チャンクを破棄するには
nullを返します。後続の Processor やクライアントには何も送信されません。 - チャンクを破棄するには
undefinedも返せます(return;文やメソッド末尾への到達による暗黙的なundefinedを含む)。nullとundefinedの動作は同じです。
チャンクの破棄は、その1つのチャンクだけに影響します。ストリームは継続し、次のチャンクも処理されます。ストリーム全体を停止するには abort() を呼び出します。
processOutputResultprocessoutputresultへの直接リンク
ストリーミングまたは生成の完了後に、完全な出力結果を処理します。
processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;
ProcessOutputResultArgsprocessoutputresultargsへの直接リンク
messages:
messageList:
state:
result:
text(蓄積テキスト)、usage(inputTokens、outputTokens、totalTokens を含むトークン使用量)、finishReason(生成が終了した理由)、steps(toolCalls、toolResults、reasoning、sources、files などを含む各 LLM ステップの全結果)を持つ、解決済みの生成結果。abort:
tripwire チャンクを送出します。retryCount:
ProcessorContext の現在の再試行回数。0 から始まり、Processor が発生させる再試行の上限設定に使用します。tracingContext?:
requestContext?:
writer?:
processOutputStepprocessoutputstepへの直接リンク
Agent ループの各 LLM レスポンスの後、Tool を実行する前に出力を処理します。最後に1回実行される processOutputResult と異なり、すべてのステップで実行されます。再試行を発生させられる Guardrail の実装に最適なメソッドです。
processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;
ProcessOutputStepArgsprocessoutputstepargsへの直接リンク
messages:
messageList:
stepNumber:
finishReason?:
providerMetadata?:
steps が空になるコンテンツフィルターによるブロックを含め、モデルステップが Provider メタデータを生成した場合に存在します。toolCalls?:
text?:
usage:
inputTokens、outputTokens、totalTokens)。systemMessages:
steps:
state:
abort:
retry: true を渡します。retryCount:
tracingContext?:
requestContext?:
ユースケースユースケースへの直接リンク
- 再試行を要求できる品質 Guardrail の実装
- Tool 実行前の LLM 出力の検証
- ステップごとのログまたはメトリクスの追加
- 再試行機能を持つ出力モデレーションの実装
例:再試行を伴う品質 Guardrail例:再試行を伴う品質 Guardrailへの直接リンク
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 []
}
}
processToolResultprocesstoolresultへの直接リンク
tool.execute() が返った後、その結果をメッセージリストへ追加する前、または次の LLM 呼び出しへ渡す前に、Tool の結果を処理します。Tool 実行前に発生する processOutputStep と対称です。Tool 出力のプロンプトインジェクション検査、機密フィールドの墨消し、abort('reason', { retry: true }) による実行中止に使用します。
Tool の結果を置換するには、messageList.updateToolInvocation で messageList をその場で変更します。ランタイムは Processor 処理後の結果をメッセージリストから再び読み取り、後続の Tool 結果ストリームチャンクをキューへ追加する前に上書きします。そのため、ストリーミングクライアントには処理済みの値が表示されます。
tool.execute() が例外をスローした場合、このメソッドは発生しません。結果を利用できる、成功した Tool 実行についてのみ呼び出されます。
processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;
ProcessToolResultArgsprocesstoolresultargsへの直接リンク
messages:
messageList:
updateToolInvocation を呼び出します。stepNumber:
toolName:
toolCallId:
args:
result:
ensureSerializable を通過した後の tool.execute() の出力です。Provider 実行型 Tool(Anthropic の web_search など)では、ensureSerializable を通らない Provider ストリームの生の結果です。providerExecuted?:
systemMessages:
steps:
state:
abort:
retry: true を渡します。retryCount:
tracingContext?:
requestContext?:
ユースケースユースケースへの直接リンク
- LLM が見る前に、Tool 出力のプロンプトインジェクションを検査する。
- Tool の戻り値から機密フィールド(PII、シークレット、認証情報)を墨消しする。
- Tool がポリシー違反のコンテンツを返したときに実行を中止する。
- コンプライアンスまたは監査用に、Tool の戻り値をログ記録または計測する。
例:機密フィールドの墨消し例:機密フィールドの墨消しへの直接リンク
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 出力内のプロンプトインジェクションをブロックへの直接リンク
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への直接リンク
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 を使用するステップ単位の Processorper-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 {}
}
}
processInputStep を使用するメッセージ変換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(入力と出力)ハイブリッド Processor(入力と出力)への直接リンク
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 を使用するストリーム蓄積への直接リンク
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 は、processLLMRequest、processLLMResponse、processOutputStream、processOutputStep、processOutputResult、processAPIError で state オブジェクトを受け取ります。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として表示されます。retry:trueの場合、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.tripwire と result.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,
})
transient は writer.custom() の第2引数ではなく、チャンクのプロパティとして渡します。第2引数には messageId などの writer オプションを指定します。
デフォルトでは、Tool のテレメトリーや自身の出力を誤って処理しないよう、Processor は processOutputStream で data-* チャンクを受け取りません。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
}
}
カスタムデータチャンクとして扱うには、チャンクの type が data- で始まる必要があります。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() は、inputProcessors、outputProcessors、errorProcessors、maxProcessorRetries を受け取ります。呼び出しにいずれかの Processor 配列を設定すると、そのリクエストについて Agent に設定された対応する配列を置き換えます。Mastra が自動的に追加する Memory、Workspace、Skill、Channel、Browser の Processor は常に保持され、指定した配列の前後で実行されます。
await agent.stream('Summarize this', {
outputProcessors: [new StreamFilter()],
maxProcessorRetries: 5,
})
呼び出しで渡した maxProcessorRetries は Agent のデフォルト値をオーバーライドします。どちらにも設定されていない場合、Processor が要求した再試行は中止として扱われます。
関連情報関連情報への直接リンク
- Processors の概要:Processor の概念ガイド
- Guardrails:セキュリティと検証の Processor
- Memory Processors:Memory 固有の Processor