> Discover all available pages from the documentation index: https://mastra.zisheng.pro/ja/llms.txt # Processor インターフェース `Processor` インターフェースは、Mastra のすべての Processor に共通する規約を定義します。Processor は1つ以上のメソッドを実装し、Agent 実行パイプラインのさまざまな段階を処理できます。 ## Processor メソッドの実行タイミング Processor メソッドは、Agent 実行ライフサイクルの異なる時点で実行されます。 ```text ┌────────────────────────────────────────────────────────────────────┐ │ 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回 | 最終レスポンスを後処理し、結果をログに記録 | ## インターフェース定義 ```typescript interface Processor { 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 processInput?( args: ProcessInputArgs, ): Promise | ProcessInputResult processInputStep?( args: ProcessInputStepArgs, ): | Promise | ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined processLLMRequest?( args: ProcessLLMRequestArgs, ): Promise | ProcessLLMRequestResult processLLMResponse?( args: ProcessLLMResponseArgs, ): Promise | ProcessLLMResponseResult processAPIError?( args: ProcessAPIErrorArgs, ): Promise | ProcessAPIErrorResult | void processOutputStream?( args: ProcessOutputStreamArgs, ): Promise processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult processOutputResult?(args: ProcessOutputResultArgs): 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`): 方法(block または warn)にかかわらず、Processor がポリシー違反を検出したときに呼び出される任意のコールバック。アラート、外部システムへのログ記録、ユーザーへのメール送信などの副作用に使用します。このコールバックがスローしたエラーは、Processor のロジックを妨げないよう暗黙的に捕捉されます。violation オブジェクトには processorId、message、Processor 固有の detail フィールドが含まれます。 ## メッセージ引数 ほとんどの Processor メソッドは `messages` と `messageList` の両方を受け取ります。どちらも同じ基盤の会話を指しますが、公開方法が異なります。 ### `messages` と `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` です。 ```typescript 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` を受け取ります。 ## メソッド ### `processInput` 入力メッセージを LLM に送信する前に処理します。Agent の実行開始時に1回実行されます。 ```typescript processInput?(args: ProcessInputArgs): Promise | ProcessInputResult; ``` #### `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` このメソッドは、3つの型のいずれかを返せます。 **MastraDBMessage\[]** (`array`): 変換済みメッセージの配列。System メッセージは変更されません。 **MessageList** (`MessageList`): 渡されたものと同じ messageList インスタンス。直接変更したことを示します。 **{ messages, systemMessages }** (`object`): 変換済みメッセージと変更済み System メッセージの両方を持つオブジェクト。 *** ### `processInputStep` Agent ループの各ステップで、LLM に送信する前に入力メッセージを処理します。開始時に1回だけ実行される `processInput` と異なり、Tool 呼び出しの継続を含むすべてのステップで実行されます。 ```typescript processInputStep?( args: ProcessInputStepArgs, ): | Promise | ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined; ``` #### 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` **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` `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 が `processInputStep` を実装している場合は順番に実行され、変更が後続へ引き継がれます。 ```text 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 メッセージは各ステップの開始時に**元の値へリセット**されます。`processInputStep` で加えた変更は現在のステップだけに影響し、後続のステップには影響しません。 #### ユースケース - ステップ番号またはコンテキストに基づく動的なモデル切り替え - 一定のステップ数を超えた後の Tool 無効化 - 会話コンテキストに基づく Tool の動的な追加または置換 - Provider 間でのメッセージ part 型の変換(Anthropic 向けの `reasoning` → `thinking` など) - ステップ番号または蓄積されたコンテキストに基づくメッセージ変更 - ステップ固有の System 指示の追加 - ステップごとの Provider オプションの調整(キャッシュ制御など) - ステップコンテキストに基づく構造化出力スキーマの変更 *** ### `processLLMRequest` Mastra が `MessageList` を `LanguageModelV2Prompt` に変換した後、Provider を呼び出す前に、最終的な LLM リクエストを処理します。現在の送信リクエストだけに影響させる、一時的かつモデルを考慮した書き換えに使用します。 返されたプロンプトの変更は、現在の呼び出しについてのみモデルへ転送されます。`MessageList`、Memory、UI 履歴、後続の Provider 呼び出しには永続化されません。 ```typescript processLLMRequest?( args: ProcessLLMRequestArgs, ): Promise | ProcessLLMRequestResult; ``` #### `ProcessLLMRequestArgs` **prompt** (`LanguageModelV2Prompt`): この呼び出しで Provider に送信する LLM リクエストプロンプト。 **model** (`MastraLanguageModel`): プロンプトを受け取る解決済みモデル。Provider 固有の書き換えを限定するために使用します。 **stepNumber** (`number`): 現在のステップ番号(0始まり)。ステップ0は最初の LLM 呼び出しです。 **steps** (`StepResult[]`): text、toolCalls、toolResults を含む、前のステップの結果。 **state** (`Record`): このリクエスト内のすべてのメソッド呼び出しで維持される、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` ステップの完了後(またはキャッシュ済みレスポンスの再生後)、output processor がレスポンスチャンクを収集した後に、LLM レスポンスを処理します。このフックは `processLLMRequest` と対になります。Provider 呼び出し前に `processLLMRequest` でキャッシュキーなどの state を保存し、完了したレスポンスへの処理(キャッシュへの書き込みなど)を `processLLMResponse` で行います。 `state` オブジェクトは、同じステップの `processLLMRequest` に渡されたものと同じインスタンスです。そのため、Processor は呼び出し前後の処理を関連付けられます。 ```typescript processLLMResponse?( args: ProcessLLMResponseArgs, ): Promise | ProcessLLMResponseResult; ``` #### `ProcessLLMResponseArgs` **chunks** (`CachedLLMStepChunk[]`): このステップで LLM 呼び出しが生成した(またはキャッシュから再生した)チャンク。簡略形式({ type, payload })です。 **model** (`MastraLanguageModel`): レスポンスを生成した(または生成するはずだった)モデル。 **stepNumber** (`number`): 現在のステップ番号(0始まり)。 **steps** (`StepResult[]`): このステップを含む、これまでに完了したすべてのステップ。 **state** (`Record`): 同じステップの 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` LLM API による拒否エラーが最終エラーとして表面化する前に処理します。API 呼び出しが再試行不能なエラー(ステータスコード 400 や 422 など)で失敗したときに実行されます。成功したレスポンスの後に実行される `processOutputStep` と異なり、API がリクエストを拒否したときに実行されます。 `processAPIError` を実装する Processor は、Agent の `errorProcessors` 配列に追加します。 Processor はエラーを調査し、`messageList` へのメッセージ追加などによってリクエストを変更できます。変更後の状態で再試行するには `{ retry: true }` を返します。 ```typescript processAPIError?(args: ProcessAPIErrorArgs): Promise | ProcessAPIErrorResult | void; ``` #### `ProcessAPIErrorArgs` **error** (`unknown`): LLM API 呼び出し中に発生したエラー。 **messages** (`MastraDBMessage[]`): エラー発生時点のすべてのメッセージ。 **messageList** (`MessageList`): メッセージ管理用の MessageList インスタンス。再試行前にリクエストを変えるには、これを変更します。 **stepNumber** (`number`): 現在のステップ番号(0始まり)。 **steps** (`StepResult[]`): これまでに完了したすべてのステップ。 **state** (`Record`): このリクエスト内のすべてのメソッド呼び出しで維持される、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` **retry** (`boolean`): 変更の適用後に LLM 呼び出しを再試行するかどうか。 #### ユースケース - リクエストを変更して再試行し、API 固有の拒否を処理する - リクエストを変更し、再試行不能なエラーを再試行可能にする - モデル固有のエラー復旧方法を実装する #### 例:カスタムエラー復旧 ```typescript 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` 組み込みの state 管理を使用して、ストリーミング出力チャンクを処理します。Processor はチャンクを蓄積し、より広いコンテキストに基づいて判断できます。 ```typescript processOutputStream?(args: ProcessOutputStreamArgs): Promise; ``` #### `ProcessOutputStreamArgs` **part** (`ChunkType`): 現在処理中のストリームチャンク。 **streamParts** (`ChunkType[]`): ストリームでこれまでに確認されたすべてのチャンク。 **state** (`Record`): 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() を呼び出します。ストリーミング中に利用できます。 #### 戻り値 `processOutputStream` は `Promise` を返します。 - チャンクを送出するには `ChunkType` を返します。変更せず送出するには元の `part` を、変更したチャンクを送出するには新しい `ChunkType` を返します。 - チャンクを破棄するには `null` を返します。後続の Processor やクライアントには何も送信されません。 - チャンクを破棄するには `undefined` も返せます(`return;` 文やメソッド末尾への到達による暗黙的な `undefined` を含む)。`null` と `undefined` の動作は同じです。 チャンクの破棄は、その1つのチャンクだけに影響します。ストリームは継続し、次のチャンクも処理されます。ストリーム全体を停止するには `abort()` を呼び出します。 *** ### `processOutputResult` ストリーミングまたは生成の完了後に、完全な出力結果を処理します。 ```typescript processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult; ``` #### `ProcessOutputResultArgs` **messages** (`MastraDBMessage[]`): 生成されたレスポンスメッセージ。 **messageList** (`MessageList`): メッセージ管理用の MessageList インスタンス。 **state** (`Record`): このリクエスト内のすべてのメソッド呼び出しで維持される、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` Agent ループの各 LLM レスポンスの後、Tool を実行する前に出力を処理します。最後に1回実行される `processOutputResult` と異なり、すべてのステップで実行されます。再試行を発生させられる Guardrail の実装に最適なメソッドです。 ```typescript processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult; ``` #### `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`): 現在のステップのトークン使用量(inputTokens、outputTokens、totalTokens)。 **systemMessages** (`CoreMessage[]`): 読み取り/変更用のすべての System メッセージ。 **steps** (`StepResult[]`): 現在のステップを含む、これまでに完了したすべてのステップ。 **state** (`Record`): このリクエスト内のすべてのメソッド呼び出しで維持される、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 ```typescript 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` `tool.execute()` が返った後、その結果をメッセージリストへ追加する前、または次の LLM 呼び出しへ渡す前に、Tool の結果を処理します。Tool 実行前に発生する `processOutputStep` と対称です。Tool 出力のプロンプトインジェクション検査、機密フィールドの墨消し、`abort('reason', { retry: true })` による実行中止に使用します。 Tool の結果を置換するには、`messageList.updateToolInvocation` で `messageList` をその場で変更します。ランタイムは Processor 処理後の結果をメッセージリストから再び読み取り、後続の Tool 結果ストリームチャンクをキューへ追加する前に上書きします。そのため、ストリーミングクライアントには処理済みの値が表示されます。 `tool.execute()` が例外をスローした場合、このメソッドは発生しません。結果を利用できる、成功した Tool 実行についてのみ呼び出されます。 ```typescript processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult; ``` #### `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`): このリクエスト内のすべてのメソッド呼び出しで維持される、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 の戻り値をログ記録または計測する。 #### 例:機密フィールドの墨消し ```typescript 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), ssn: '[REDACTED]', email: '[REDACTED]', } messageList.updateToolInvocation({ type: 'tool-invocation', toolInvocation: { state: 'result', toolCallId, toolName, args, result: redacted, }, }) } } ``` #### 例:Tool 出力内のプロンプトインジェクションをブロック ```typescript 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 の型 Mastra は、Processor が必須メソッドを実装することを保証する型エイリアスを提供します。 ```typescript // 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` に設定します。 ```typescript const agent = new Agent({ id: 'agent', errorProcessors: [new PrefillErrorHandler()], }) ``` ## 使用例 ### 基本的な input processor ```typescript 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 { 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 ```typescript import type { Processor, ProcessInputStepArgs, ProcessInputStepResult, } from '@mastra/core/processors' export class DynamicModelProcessor implements Processor { id = 'dynamic-model' async processInputStep({ stepNumber, steps, toolChoice, }: ProcessInputStepArgs): Promise { // 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` を使用するメッセージ変換 ```typescript 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(入力と出力) ```typescript 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 { 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 { if (part.type === 'text-delta') { if (this.blockedWords.some(word => part.payload.text.includes(word))) { abort('Blocked content detected in output') } } return part } } ``` ### state を使用するストリーム蓄積 ```typescript 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 { // 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 のライフサイクル すべての 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` は空のオブジェクトとして始まるため、最初にアクセスするときはフィールドを防御的に初期化してください。 ```typescript 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 チャンク 各メソッドの `abort` 関数は、処理を停止する `TripWire` エラーをスローし、出力ストリームに `tripwire` チャンクを送出します。クライアントはこのチャンクを検出し、ブロックされたレスポンスと通常の完了を区別できます。 ```typescript 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` チャンクの形式は次のとおりです。 ```typescript 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 チャンク以外のコンテンツを送出する方法は、これだけです。 ```typescript await writer.custom({ type: 'data-moderation', runId, from: 'AGENT', data: { level: 'warn', reason: 'Possibly unsafe' }, }) ``` Memory が設定されている場合、`processOutputStream` または `processOutputResult` から送出されたカスタム `data-*` チャンクは、Assistant メッセージの part として保存されます。Memory に保存せずストリーミングするには、チャンクオブジェクトに `transient: true` を設定します。 ```typescript 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` を設定すると受け取れます。 ```typescript 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 の設定 Processor は3つの配列を通じて Agent に設定します。 ```typescript 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 を構築できます。 ```typescript 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 は常に保持され、指定した配列の前後で実行されます。 ```typescript await agent.stream('Summarize this', { outputProcessors: [new StreamFilter()], maxProcessorRetries: 5, }) ``` 呼び出しで渡した `maxProcessorRetries` は Agent のデフォルト値をオーバーライドします。どちらにも設定されていない場合、Processor が要求した再試行は中止として扱われます。 ## 関連情報 - [Processors の概要](https://mastra.zisheng.pro/ja/docs/agents/processors):Processor の概念ガイド - [Guardrails](https://mastra.zisheng.pro/ja/docs/agents/guardrails):セキュリティと検証の Processor - [Memory Processors](https://mastra.zisheng.pro/ja/docs/memory/memory-processors):Memory 固有の Processor