Processor
Processor は、Agent を通過するメッセージを変換、検証、制御します。Agent の実行パイプライン内の特定の時点で実行されるため、言語モデルに届く前の入力や、ユーザーに返る前の出力を変更できます。
Processor は次のように設定します。
inputProcessors:メッセージが言語モデルに届く前に実行されます。outputProcessors:言語モデルがレスポンスを生成した後、ユーザーに返る前に実行されます。
個別の Processor オブジェクトを使用するか、Mastra の Workflow プリミティブで Workflow に構成できます。Workflow を使うと、Processor の実行順序、並列処理、条件ロジックを高度に制御できます。
入力と出力の両方のロジックを実装する Processor もあり、変換する位置に応じてどちらの配列でも使用できます。
一部の組み込み Processor は、非表示の system reminder signal も送信します。この signal は生の Memory 履歴に保存され、次のモデル呼び出し前に <system-reminder>...</system-reminder> コンテキストへ変換されます。ただし、明示的に有効にしない限り、標準の UI 向けメッセージ変換とデフォルトの Memory recall では非表示になります。signal を保持せず現在の呼び出しだけに渡すには、transient: true を指定します。
Processor を使用する場面Processor を使用する場面への直接リンク
Processor は次の用途に使用します。
- ユーザー入力の正規化や検証
- Agent への Guardrail の追加
- プロンプトインジェクションや jailbreak の検出と防止
- 安全性やコンプライアンスのためのコンテンツモデレーション
- メッセージの変換(言語の翻訳、Tool 呼び出しのフィルタリングなど)
- トークン使用量やメッセージ履歴の長さの制限
- 機密情報(PII)の編集
- メッセージへのカスタムビジネスロジックの適用
Mastra には一般的なユースケース向けの Processor が複数用意されています。アプリケーション固有の要件に合わせてカスタム Processor も作成できます。
クイックスタートクイックスタートへの直接リンク
Processor をインポートしてインスタンス化し、Agent の inputProcessors または outputProcessors 配列に渡します。
import { Agent } from '@mastra/core/agent'
import { ModerationProcessor } from '@mastra/core/processors'
export const moderatedAgent = new Agent({
id: 'moderated-agent',
name: 'moderated-agent',
instructions: 'You are a helpful assistant',
model: 'openai/gpt-5-mini',
inputProcessors: [
new ModerationProcessor({
model: 'openai/gpt-5-mini',
categories: ['hate', 'harassment', 'violence'],
threshold: 0.7,
strategy: 'block',
}),
],
})
実行順序実行順序への直接リンク
Processor は配列に記述された順に実行されます。
inputProcessors: [new UnicodeNormalizer(), new PromptInjectionDetector(), new ModerationProcessor()]
出力 Processor では、この順序によってモデルのレスポンスに適用する変換の順番が決まります。
Memory が有効な場合Memory が有効な場合への直接リンク
Agent で Memory が有効な場合、Memory Processor がパイプラインに自動で追加されます。
入力 Processor:
[Memory Processors] → [Your inputProcessors]
最初に Memory がメッセージ履歴を読み込み、その後に独自の Processor が実行されます。
出力 Processor:
[Your outputProcessors] → [Memory Processors]
最初に独自の Processor が実行され、その後に Memory がメッセージを永続化します。
この順序では、abort() を呼ぶ出力 Guardrail は Memory Processor をスキップし、メッセージの保存を防ぎます。詳しくは Memory Processorを参照してください。
Agent に Processor を追加するAgent に Processor を追加するへの直接リンク
Processor は3つの配列で Agent に設定します。
import { Agent } from '@mastra/core/agent'
import { PrefillErrorHandler, TokenLimiter, ModerationProcessor } from '@mastra/core/processors'
const agent = new Agent({
id: 'support-agent',
name: 'support-agent',
model: 'openai/gpt-5',
instructions: '...',
inputProcessors: [
new TokenLimiter(4000),
new ModerationProcessor({ model: 'openai/gpt-5-nano' }),
],
outputProcessors: [new ModerationProcessor({ model: 'openai/gpt-5-nano' })],
errorProcessors: [new PrefillErrorHandler()],
})
inputProcessorsは LLM の前に実行されます。outputProcessorsは LLM のレスポンス中とレスポンス後に実行されます。errorProcessorsは LLM API 呼び出しが例外をスローしたときに実行され、Provider エラーから復旧できます。
各配列には、配列を返す関数も指定できます。これにより、RequestContext からリクエストごとに Processor を構築できます。
new Agent({
id: 'processors-agent',
inputProcessors: ({ requestContext }) => {
const limit = requestContext.get('tokenLimit') ?? 4000
return [new TokenLimiter(limit)]
},
})
呼び出しごとに Processor を上書きする呼び出しごとに Processor を上書きするへの直接リンク
agent.generate() と agent.stream() も同じ3つの配列を受け取ります。配列を渡すと、その呼び出しに限り、Agent に設定された対応する配列を置き換えます。Memory、Workspace、その他のフレームワーク管理 Processor は、引き続き指定した配列の前後で実行されます。
await agent.stream('Summarize this', {
inputProcessors: [new TokenLimiter(2000)],
maxProcessorRetries: 5,
})
カスタム Processor を作成するカスタム Processor を作成するへの直接リンク
カスタム Processor は Processor interface を実装します。
Processor のメソッドは、会話にアクセスするための2つの引数を受け取ります。
messages:現在の段階にあるMastraDBMessageオブジェクトのスナップショット配列。messageList:稼働中のMessageListinstance。ほかの段階の読み取りや、メッセージの追加、削除、置換に使用します。
テキストは message.content 自体ではなく message.content.parts にあります。ユーザーまたは assistant のテキストを読むには、parts を反復し、part.type === 'text' で絞り込みます。従来との互換性のため、平坦化された message.content.content 文字列もフォールバックとして使用できます。詳しくは Processor reference のメッセージ引数を参照してください。
入力メッセージを変換する入力メッセージを変換するへの直接リンク
import type { Processor, ProcessInputArgs } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class CustomInputProcessor implements Processor {
id = 'custom-input'
async processInput({ messages }: ProcessInputArgs): Promise<MastraDBMessage[]> {
// Transform messages before they reach the LLM.
// Text lives in content.parts — iterate parts and rewrite text parts only.
return messages.map(msg => ({
...msg,
content: {
...msg.content,
parts: msg.content.parts?.map(part =>
part.type === 'text' ? { ...part, text: part.text.toLowerCase() } : part,
),
},
}))
}
}
processInput() メソッドは messages、systemMessages、abort() 関数を受け取ります。メッセージを置き換えるには MastraDBMessage[] を返し、system message も変更するには { messages, systemMessages } を返します。
使用できるすべての引数と戻り値の型は Processor referenceを参照してください。
各ステップを制御する各ステップを制御するへの直接リンク
processInput() は Agent 実行の開始時に一度だけ実行されますが、processInputStep() は Agent ループの各ステップ(Tool 呼び出しの継続を含む)で実行されます。実行時のモデル切り替えや Tool choice の変更など、ステップごとの設定変更が可能です。
import type {
Processor,
ProcessInputStepArgs,
ProcessInputStepResult,
} from '@mastra/core/processors'
export class DynamicModelProcessor implements Processor {
id = 'dynamic-model'
async processInputStep({
stepNumber,
model,
toolChoice,
messageList,
}: ProcessInputStepArgs): Promise<ProcessInputStepResult> {
// Use a fast model for initial response
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
// Disable tools after 5 steps to force completion
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
// No changes for other steps
return {}
}
}
このメソッドは現在の stepNumber、model、tools、toolChoice、messages などを受け取ります。そのステップで上書きするプロパティを含むオブジェクト(例:{ model, toolChoice, tools, systemMessages })を返します。
使用できるすべての引数と戻り値の型は Processor referenceを参照してください。
Provider 呼び出し前に LLM リクエストを書き換えるProvider 呼び出し前に LLM リクエストを書き換えるへの直接リンク
Mastra がモデルに送る最終プロンプトを書き換えるには、processLLMRequest() を使用します。このフックは、Mastra が MessageList を Provider 向けのプロンプト形式(LanguageModelV2Prompt)に変換した後、Provider 呼び出しの直前に実行されます。
会話を変更する場合は、メッセージベースのフックを使用します。
processInput():Agent ループの開始前に会話を一度変更します。processInputStep():各 LLM 呼び出し前にメッセージやステップ設定を変更します。processLLMRequest():現在の Provider 呼び出しに対する送信プロンプトだけを変更します。
processLLMRequest() から返した変更は一時的です。MessageList、Memory、UI 履歴、将来の Provider 呼び出しには保存されません。このため、保存された会話履歴を変えるべきでない Provider 互換性のための書き換え、role/content の正規化、その他のモデル固有のプロンプト変更に適しています。
このメソッドは prompt、model、stepNumber、steps、state、共有 Processor コンテキストを受け取ります。processLLMRequest() から abort() を呼ぶと通常の tripwire レスポンスが送出され、呼び出しが停止します。
使用できるすべての引数と戻り値の型は Processor referenceを参照してください。
Provider 呼び出し後に LLM レスポンスを処理するProvider 呼び出し後に LLM レスポンスを処理するへの直接リンク
ステップが完了し、ストリームチャンクが収集された後の LLM レスポンスを処理するには、processLLMResponse() を使用します。このフックは processLLMRequest() と対になります。リクエストフックで状態(キャッシュキーなど)を保存し、レスポンスフックで読み戻してキャッシュへの書き込みなどの副作用を実行します。
state オブジェクトは、同じステップで processLLMRequest() に渡されるものと同じ instance です。fromCache が true の場合、レスポンスは実際のモデル呼び出しではなくキャッシュから再生されています。この場合、キャッシュへ書き込む Processor は書き込みをスキップしてください。
このメソッドは chunks、model、stepNumber、steps、state、fromCache、共有 Processor コンテキストを受け取ります。
使用できるすべての引数と戻り値の型は Processor referenceを参照してください。
prepareStep() callback を使用するuse-the-preparestep-callbackへの直接リンク
generate() または stream() の prepareStep() callback は processInputStep() の短縮形です。内部では、Mastra がこの関数を各ステップで呼ぶ Processor にラップします。processInputStep() と同じ引数と戻り値の型を受け取りますが、クラスを作成する必要はありません。
await agent.generate('Complex task', {
prepareStep: async ({ stepNumber, model }) => {
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
},
})
出力メッセージを変換する出力メッセージを変換するへの直接リンク
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class CustomOutputProcessor implements Processor {
id = 'custom-output'
async processOutputResult({ messages }): Promise<MastraDBMessage[]> {
// Transform messages after the LLM generates them
return messages.filter(msg => msg.role !== 'system')
}
}
このメソッドは、生成データ全体を含む result オブジェクト、text、usage(トークン数)、finishReason、steps(それぞれに toolCalls、toolResults などを含む)も受け取ります。使用量の追跡や Tool 呼び出しの検査に使用できます。
import type { Processor } from '@mastra/core/processors'
export class UsageTracker implements Processor {
id = 'usage-tracker'
async processOutputResult({ messages, result }) {
console.log(`Tokens: ${result.usage.inputTokens} in, ${result.usage.outputTokens} out`)
console.log(`Finish reason: ${result.finishReason}`)
return messages
}
}
ストリーミング出力をフィルタリングするストリーミング出力をフィルタリングするへの直接リンク
processOutputStream() メソッドは、クライアントに届く前のストリーミングチャンクを変換またはフィルタリングします。
import type { Processor } from '@mastra/core/processors'
import type { ChunkType } from '@mastra/core/stream'
export class StreamFilter implements Processor {
id = 'stream-filter'
async processOutputStream({ part }): Promise<ChunkType | null> {
// Drop text-delta chunks that contain the word "secret"
if (part.type === 'text-delta' && part.payload.text.includes('secret')) {
return null
}
// Return the (possibly modified) chunk to emit it
return part
}
}
戻り値:
ChunkTypeはそのチャンクを送出します。元のpartを返すと、変更せずそのまま通します。nullまたはundefinedはチャンクを破棄します。どちらも同じ動作のため、何も返さないメソッドもチャンクを破棄します。- 破棄の影響は1つのチャンクだけです。ストリーム全体を停止するには
abort()を呼びます。
Tool が writer.custom() で送出したカスタム data-* チャンクも受信するには、Processor に processDataParts = true を設定します。クライアントに届く前に、Tool が送出したデータチャンクを検査、変更、ブロックできます。
各レスポンスを検証する各レスポンスを検証するへの直接リンク
processOutputStep() メソッドは各 LLM ステップ後に実行され、レスポンスの検証と任意の再試行要求が可能です。
import type { Processor } from '@mastra/core/processors'
export class ResponseValidator implements Processor {
id = 'response-validator'
async processOutputStep({ text, abort, retryCount }) {
const isValid = await validateResponse(text)
if (!isValid && retryCount < 3) {
abort('Response did not meet requirements. Try again.', { retry: true })
}
return []
}
}
再試行の動作については、高度なパターンの再試行メカニズムを参照してください。
チャンクとステップ間でデータを保持するチャンクとステップ間でデータを保持するへの直接リンク
出力メソッドは、1つのリクエストの間保持される state オブジェクトを受け取ります。状態は Processor の id で分けられるため、各 Processor からは自身のデータだけが見えます。また、processOutputStream、processOutputStep、processOutputResult の間で共有されます。agent.generate() または agent.stream() を新しく呼び出すたびに、新しい状態オブジェクトが作成されます。
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
}
async processOutputResult({ messages, state }) {
console.log(`Total words: ${state.wordCount}`)
return messages
}
}
組み込み Utility Processor組み込み Utility Processorへの直接リンク
Mastra は一般的なタスク向けの Utility Processor を提供します。
セキュリティと検証の Processorについては、入力・出力 Guardrail とモデレーション Processor を説明する Guardrails ページを参照してください。 Memory 固有の Processorについては、メッセージ履歴、semantic recall、working memory を処理する Processor を説明する Memory Processor ページを参照してください。
TokenLimitertokenlimiterへの直接リンク
合計トークン数が指定した上限を超えたときに古いメッセージを削除し、コンテキストウィンドウの超過を防ぎます。最近のメッセージを優先し、system message を保持します。
import { Agent } from '@mastra/core/agent'
import { TokenLimiter } from '@mastra/core/processors'
const agent = new Agent({
id: 'my-agent',
name: 'my-agent',
model: 'openai/gpt-5.6-sol',
inputProcessors: [new TokenLimiter(127000)],
})
カスタム encoding、strategy、count mode のオプションは TokenLimiterProcessor referenceを参照してください。
ToolCallFiltertoolcallfilterへの直接リンク
LLM に送るメッセージから Tool 呼び出しと結果を削除し、冗長な Tool interaction によるトークンを節約します。特定の Tool だけを除外することもできます。このフィルターが影響するのは LLM 入力だけで、フィルタリングされたメッセージも Memory には保存されます。
デフォルトでは、ToolCallFilter は Agent ループ開始前の初期入力をフィルタリングします。最近の Tool 生成ステップを保持しながら各ループステップでもフィルタリングするには、filterAfterToolSteps を使用します。
new ToolCallFilter({
filterAfterToolSteps: 2,
})
フィルタリング済みの完了した Tool 結果について、簡潔な toModelOutput 履歴を保持するには preserveModelOutput: true を設定します。モデル向けの出力だけを保持し、生の Tool 引数と結果は削除します。
new ToolCallFilter({
preserveModelOutput: true,
})
設定オプションは ToolCallFilter referenceを、Memory 前のフィルタリングは Memory Processorページを参照してください。
ToolSearchProcessortoolsearchprocessorへの直接リンク
多数の Tool ライブラリを持つ Agent に、実行時の Tool 検索を提供します。すべての Tool を最初から渡す代わりに、search_tools と load_tool のメタ Tool を Agent に渡します。Agent は必要に応じてキーワードで Tool を検索して読み込めるため、コンテキストのトークン使用量を削減できます。
設定オプションと使用例は ToolSearchProcessor referenceを参照してください。
ProviderHistoryCompatproviderhistorycompatへの直接リンク
Agent が複数のモデル Provider 間でメッセージを再利用する際に、Provider 固有の履歴の非互換性を処理します。Provider 呼び出し前に送信 LLM リクエストを書き換えるか、既知の Provider API エラーから復旧して再試行できます。
Provider の履歴互換性ルール、リアクティブな API エラー復旧、カスタム互換性ルール、予測可能な Processor 順序が必要な場合は、ProviderHistoryCompat を明示的に追加します。
設定、組み込みルール、カスタムルールのオプションは ProviderHistoryCompat referenceを参照してください。
レスポンスキャッシュレスポンスキャッシュへの直接リンク
この機能はベータ版です。API が安定するまで、メジャーバージョンを上げずに破壊的変更が行われる可能性があります。
レスポンスキャッシュは、Agent が同一のリクエストを受け取ったときに LLM 呼び出しを省略し、以前キャッシュしたレスポンスを再生します。レイテンシーを短縮し、繰り返し呼び出すコストを回避できます。
キャッシュは ResponseCache 入力 Processor として実装されています。Mastra は Agent レベルのオプションを提供していません。有効にするには Processor を明示的に登録します。これにより、Mastra がフィードバックを収集している間、API surface を小さく保てます。呼び出しごとの上書きは RequestContext を介して渡されます。
レスポンスキャッシュを使用する場面レスポンスキャッシュを使用する場面への直接リンク
プロンプトテンプレート、推奨プロンプトボタン、Agent 検索の再質問、同じ入力を繰り返し分類する Guardrail LLM など、同じ形式のリクエストがユーザー間やセッション間で繰り返される場合に使用します。Tool による外部の副作用が発生する呼び出しでは、キャッシュヒット時に Tool 呼び出しが再実行されず再生されるため、使用しないでください。
クイックスタートクイックスタートへの直接リンク
Agent の inputProcessors に ResponseCache を追加し、バックエンドとして任意の MastraServerCache を渡します。開発時は InMemoryServerCache をそのまま使用できます。
import { Agent } from '@mastra/core/agent'
import { InMemoryServerCache } from '@mastra/core/cache'
import { ResponseCache } from '@mastra/core/processors'
const cache = new InMemoryServerCache()
export const searchAgent = new Agent({
id: 'search-agent',
name: 'Search Agent',
instructions: 'You answer questions concisely.',
model: 'openai/gpt-5',
inputProcessors: [new ResponseCache({ cache, ttl: 600 })], // 10 minutes
})
最初の呼び出しでは通常どおり LLM を実行し、レスポンスをキャッシュに書き込みます。解決済みプロンプトが同一の以降の呼び出しでは、LLM を呼び出さずキャッシュ済みレスポンスを返します。
RequestContext で呼び出しごとに上書きするRequestContext で呼び出しごとに上書きするへの直接リンク
呼び出しごとの設定は RequestContext を介して渡されます。新しいコンテキストを構築するには ResponseCache.context()、既存のコンテキストにマージするには ResponseCache.applyContext() を使用します。
import { ResponseCache } from '@mastra/core/processors'
import { RequestContext } from '@mastra/core/request-context'
// Fresh context with the override
await agent.stream('hello', {
requestContext: ResponseCache.context({ key: 'custom-key', bust: true }),
})
// Or merge into an existing context
const ctx = new RequestContext()
ctx.set('caller-meta', { userId: 'u-123' })
ResponseCache.applyContext(ctx, { bust: true })
await agent.stream('hello', { requestContext: ctx })
呼び出しごとに上書きできるフィールドは次のとおりです。
key:文字列または関数。このリクエストだけ、自動生成されるキャッシュキーを上書きします。scope:文字列またはnull。このリクエストだけ、tenant/user scope を上書きします。nullは scope を無効にします。bust:boolean。キャッシュの読み取りをスキップしますが、完了時には書き込みます(「強制更新」ボタンなどに便利です)。
cache、ttl、agentId は constructor に残ります。instance レベルの項目であり、呼び出しごとの変更は安全ではありません。
Tenant scopeTenant scopeへの直接リンク
デフォルトでは、ResponseCache はリクエストコンテキストの MASTRA_RESOURCE_ID_KEY を検索し、キャッシュの scope として使用します。つまり、resource ID がすでに設定されている Agent(Memory 経由など)では、ユーザー単位の分離が自動的に行われます。ユーザー間でキャッシュされたレスポンスが共有されることはありません。
別の scope が必要な場合は明示的に上書きします。
new Agent({
id: 'processors-agent',
inputProcessors: [
new ResponseCache({
cache,
scope: 'org-123', // explicit tenant scope
}),
],
})
すべての呼び出し元で意図的に entry を共有するには scope: null を渡します。公開済みでパーソナライズされていないことが明らかなコンテンツにだけ使用してください。
カスタムキャッシュバックエンドカスタムキャッシュバックエンドへの直接リンク
ResponseCache は任意の MastraServerCache を受け取ります。本番環境では @mastra/redis の RedisCache を使用します。
import { Agent } from '@mastra/core/agent'
import { ResponseCache } from '@mastra/core/processors'
import { RedisCache } from '@mastra/redis'
const cache = new RedisCache({ url: process.env.REDIS_URL })
export const agent = new Agent({
id: 'cached-agent',
name: 'Cached Agent',
instructions: '...',
model: 'openai/gpt-5',
inputProcessors: [new ResponseCache({ cache })],
})
カスタムバックエンドでは、MastraServerCache を継承して抽象メソッドを実装します(Processor が呼ぶのは get と set だけです)。
キャッシュの実装キャッシュの実装への直接リンク
ResponseCache は processLLMRequest(キャッシュ検索、ヒット時の早期終了)と processLLMResponse(完了時のキャッシュ書き込み)にフックします。どちらも Memory の読み込みと先行する入力 Processor によるプロンプト変換の_後_、Agent ループ内で実行されます。
そのため、キャッシュキーは Mastra がモデルに送ろうとしている解決済みの LanguageModelV2Prompt から生成されます。Memory が読み込まれ、先行する入力 Processor が実行された_後_にキーが作成され、Agent の Tool ループ内の各ステップが個別にキャッシュされます。
キャッシュキーの内容キャッシュキーの内容への直接リンク
key を指定しない場合、Processor はこのステップで LLM のレスポンスを変える入力から決定論的にキーを生成します。対象は agentId、stepNumber(Tool ループの各ステップに固有のキャッシュ entry を割り当てるため)、scope、モデルの識別情報(provider、modelId、spec version)、解決済みの prompt(Memory と Processor の処理後)です。これらの入力が変わると、キャッシュは自動的に無効になります。
マルチモーダルプロンプトも含まれます。画像とファイルの part は値としてキーに入ります。URL は完全な href、インラインのバイナリデータ(Uint8Array、ArrayBuffer)はバイト列の digest を提供します。そのため、参照する画像だけが異なる2つのリクエストには別々のキャッシュ entry が割り当てられます。
キャッシュキーをカスタマイズするキャッシュキーをカスタマイズするへの直接リンク
constructor または呼び出しごとに key を関数として渡すと、これらの入力の任意の subset から独自のキャッシュキーを生成できます。関数は決定論的ハッシュが使用するものと同じ入力を受け取り、文字列(または Promise<string>)を返します。
import { ResponseCache, buildResponseCacheKey } from '@mastra/core/processors'
await agent.stream(input, {
requestContext: ResponseCache.context({
// Cache only on the model id and the resolved prompt tail — ignore
// step number, scope, etc.
key: ({ model, prompt }) => `qa:${model.modelId}:${JSON.stringify(prompt).slice(-200)}`,
}),
})
// Or reuse the deterministic helper while overriding individual fields:
await agent.stream(input, {
requestContext: ResponseCache.context({
key: inputs => buildResponseCacheKey({ ...inputs, scope: 'global' }),
}),
})
関数が例外をスローした場合、呼び出しで引き続きキャッシュを利用できるよう、Processor はデフォルトのキー生成にフォールバックします。
キャッシュヒットの仕組みキャッシュヒットの仕組みへの直接リンク
キャッシュがヒットすると、Processor は processLLMRequest からキャッシュ済みチャンクを返して LLM 呼び出しを早期終了します。Agent ループはモデルを呼び出す代わりに、そのチャンクからストリームを合成します。agent.generate() はチャンクを FullOutput に収集します。agent.stream() はキャッシュ済みバッファからのチャンクを持つ MastraModelOutput を返すため、fullStream を反復するか text、usage、finishReason を await する consumer にはキャッシュ済みの値が表示されます。
キャッシュへの書き込みはレスポンス完了後に行われます。失敗した実行(エラー、tripwire の発動)はキャッシュされないため、次回の呼び出しで正常に再試行されます。
高度なパターン高度なパターンへの直接リンク
maxSteps で最終レスポンスを保証するensure-a-final-response-with-maxstepsへの直接リンク
maxSteps で Agent の実行を制限すると、最終ステップで Tool 呼び出しを試みた場合に空のレスポンスを返すことがあります。processInputStep() と sendSignal を使用して、最終ステップにリアクティブな reminder を注入します。system message を変更せず signal を追加するため、プロンプトキャッシュを維持できます。
import type { Processor, ProcessInputStepArgs } from '@mastra/core/processors'
export class EnsureFinalResponseProcessor implements Processor {
readonly id = 'ensure-final-response'
private maxSteps: number
constructor(maxSteps: number) {
this.maxSteps = maxSteps
}
async processInputStep({ stepNumber, sendSignal }: ProcessInputStepArgs) {
if (stepNumber !== this.maxSteps - 1) {
return
}
await sendSignal?.({
type: 'reactive',
contents:
`This is your final step (step ${stepNumber + 1} of ${this.maxSteps}). ` +
`Do not call any more tools. Summarize what you have found and give the user a complete final answer now.`,
attributes: { reason: 'max-steps-reached', step: stepNumber + 1 },
})
}
}
signal は、モデルがインラインで認識する <system-reminder> user message として渡されます。
<system-reminder reason="max-steps-reached" step="5">This is your final step (step 5 of 5). Do not call any more tools. Summarize what you have found and give the user a complete final answer now.</system-reminder>
Processor を inputProcessors に追加し、signal tag を説明する system prompt を含め、同じ maxSteps 値を generate() または stream() に渡します。
import { Agent } from '@mastra/core/agent'
import { EnsureFinalResponseProcessor } from '../processors/ensure-final-response'
const MAX_STEPS = 5
const agent = new Agent({
id: 'agent',
instructions: `You are a helpful assistant.
Some messages you receive may contain <system-reminder>...</system-reminder> tags.
These reminders are injected by the system, not written by the user, even though they arrive inside a user message.
Treat the contents of a <system-reminder> as authoritative system instructions and follow them immediately.
Do not mention the reminder to the user or quote the tags back to them.`,
inputProcessors: [new EnsureFinalResponseProcessor(MAX_STEPS)],
// ...
})
await agent.generate('Your prompt', { maxSteps: MAX_STEPS })
Reactive signal のデフォルトは tagName: 'system-reminder' です。Processor が送出する signal について詳しくは Signalを参照してください。
保持せず reminder を渡す保持せず reminder を渡すへの直接リンク
デフォルトでは、Processor から送信された signal は会話の一部になります。ストレージに書き込まれ、後のターンで再びプロンプトに入ります。毎ターン再注入する指示では、コピーが蓄積し、モデルが自身の過去の reminder を模倣すべき以前のコンテキストとして扱い始めるため望ましくありません。signal を保持せず現在の呼び出しだけでモデルに渡すには、transient: true を設定します。
使用する場面: 会話が長くなっても、短い誘導指示をモデルの直近のコンテキストに残したい場合です。たとえば「現在のタスクに集中する」「回答を3文以内にする」、またはアプリケーションの現在の状態に応じたターン単位の制約です。最新のメッセージ付近に保つため、各ターンで再注入します。
import type { Processor, ProcessInputStepArgs } from '@mastra/core/processors'
export class SteeringReminderProcessor implements Processor {
readonly id = 'steering-reminder'
async processInputStep({ sendSignal }: ProcessInputStepArgs) {
await sendSignal?.({
type: 'reactive',
contents: 'Stay on the current task and keep answers under three sentences.',
transient: true,
})
}
}
一時的な signal も現在の呼び出しのプロンプトには現れるため、モデルは最新のターン付近で認識します。保持されないので、各ターンで再送しても蓄積した履歴ではなく新しいコピーが1つだけコンテキストに入り、保存された thread 履歴には一切現れません。何も書き込まれないため、ターンをまたいで安定したプロンプトキャッシュの prefix も維持できます。
カスタムストリームイベントを送出するカスタムストリームイベントを送出するへの直接リンク
出力 Processor は writer オブジェクトを受け取り、ストリーミング中にカスタムデータチャンクをクライアントへ返せます。元のストリームをブロックせずにモデレーション結果をストリーミングしたり、UI 更新 signal を送信したりする用途に便利です。
import type { Processor } from '@mastra/core/processors'
export class ModerationProcessor implements Processor {
id = 'moderation'
async processOutputResult({ messages, writer }) {
// Run moderation on the final output
const text = messages
.filter(m => m.role === 'assistant')
.flatMap(m => m.content.parts?.filter(p => p.type === 'text'))
.map(p => p.text)
.join(' ')
const result = await runModeration(text)
if (result.requiresChange) {
// Emit a custom event to the client with the moderated text
await writer?.custom({
type: 'data-moderation-update',
data: {
originalText: text,
moderatedText: result.moderatedText,
reason: result.reason,
},
})
}
return messages
}
}
クライアントでは、ストリーム内のカスタムチャンク型をリッスンします。
const stream = await agent.stream('Hello')
for await (const chunk of stream.fullStream) {
if (chunk.type === 'data-moderation-update') {
// Update the UI with moderated text
updateDisplayedMessage(chunk.data.moderatedText)
}
}
カスタムチャンク型には data- prefix を使用する必要があります(例:data-moderation-update、data-status)。
デフォルトでは、processOutputStream() は data-* チャンクをスキップし、Tool telemetry や別の Processor の出力を誤って処理しないようにします。Processor でこれらのチャンクを検査、変更、ブロックするには、その Processor に processDataParts = true を設定します。
class ModerationCollector implements Processor {
id = 'moderation-collector'
processDataParts = true
async processOutputStream({ part, state }) {
if (part.type === 'data-moderation-update') {
state.warnings ??= []
state.warnings.push(part.data)
}
return part
}
}
メッセージにメタデータを追加するメッセージにメタデータを追加するへの直接リンク
processOutputResult でメッセージにカスタムメタデータを追加できます。このメタデータにはレスポンスオブジェクトからアクセスできます。
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class MetadataProcessor implements Processor {
id = 'metadata-processor'
async processOutputResult({
messages,
}: {
messages: MastraDBMessage[]
}): Promise<MastraDBMessage[]> {
return messages.map(msg => {
if (msg.role === 'assistant') {
return {
...msg,
content: {
...msg.content,
metadata: {
...msg.content.metadata,
processedAt: new Date().toISOString(),
customData: 'your data here',
},
},
}
}
return msg
})
}
}
generate() でメタデータにアクセスします。
const result = await agent.generate('Hello')
// The response includes uiMessages with processor-added metadata
const assistantMessage = result.response?.uiMessages?.find(m => m.role === 'assistant')
console.log(assistantMessage?.metadata?.customData)
ストリーミングでは、finish チャンクの payload または stream.response promise からメタデータにアクセスします。
Workflow を Processor として使用するWorkflow を Processor として使用するへの直接リンク
Mastra Workflow を Processor として使用すると、並列実行、条件分岐、エラー処理を備えた複雑な処理パイプラインを作成できます。
import { createWorkflow, createStep } from '@mastra/core/workflows'
import {
ProcessorStepSchema,
PromptInjectionDetector,
PIIDetector,
ModerationProcessor,
} from '@mastra/core/processors'
import { Agent } from '@mastra/core/agent'
// Create a workflow that runs multiple checks in parallel
const moderationWorkflow = createWorkflow({
id: 'moderation-pipeline',
inputSchema: ProcessorStepSchema,
outputSchema: ProcessorStepSchema,
})
.parallel([
createStep(
new PIIDetector({
strategy: 'redact',
}),
),
createStep(
new PromptInjectionDetector({
strategy: 'block',
}),
),
createStep(
new ModerationProcessor({
strategy: 'block',
}),
),
])
.map(async ({ inputData }) => {
return inputData['processor:pii-detector']
})
.commit()
// Use the workflow as an input processor
const agent = new Agent({
id: 'moderated-agent',
name: 'Moderated Agent',
model: 'openai/gpt-5.6-sol',
inputProcessors: [moderationWorkflow],
})
.parallel() ステップの後、各分岐の結果には Processor ID(例:processor:pii-detector)をキーとして割り当てます。次のステップが受け取る出力の分岐を選ぶには .map() を使用します。
分岐が redact などの変更を伴う strategy を使用する場合、変換済みメッセージを後続へ渡すため、その分岐に map します。すべての分岐が block だけなら、どれを選んでも構いません。どの分岐もメッセージを変更しないため、任意の1つを選びます。
Agent を Mastra に登録すると、Processor Workflow も Workflow として自動的に登録され、Studio で表示、デバッグできます。
再試行メカニズム再試行メカニズムへの直接リンク
Processor は、フィードバックを添えて LLM にレスポンスの再試行を要求できます。品質チェック、出力検証、反復的な改善の実装に便利です。
import type { Processor } from '@mastra/core/processors'
export class QualityChecker implements Processor {
id = 'quality-checker'
async processOutputStep({ text, abort, retryCount }) {
const qualityScore = await evaluateQuality(text)
if (qualityScore < 0.7 && retryCount < 3) {
// Request a retry with feedback for the LLM
abort('Response quality score too low. Please provide a more detailed answer.', {
retry: true,
metadata: { score: qualityScore },
})
}
return []
}
}
const agent = new Agent({
id: 'quality-agent',
name: 'Quality Agent',
model: 'openai/gpt-5.6-sol',
outputProcessors: [new QualityChecker()],
maxProcessorRetries: 3, // Maximum retry attempts. If unset, retries are disabled (unless errorProcessors are configured, in which case it defaults to 10).
})
再試行メカニズムは次のように動作します。
processOutputStep()とprocessInputStep()メソッドで動作します。- 中止理由を LLM のコンテキストに追加してステップを再実行します。
retryCountparameter で再試行回数を追跡します。- Agent または呼び出しに明示的な
maxProcessorRetries上限が必要です。
Error Processor の再試行上限Error Processor の再試行上限への直接リンク
processAPIError() には別のデフォルトがあります。errorProcessors が設定され、maxProcessorRetries が省略されている場合、ランタイムは最大 10 回の再試行を許可します。再試行回数を制限する必要がある場合は、上限を明示的に設定します。
StreamErrorRetryProcessor では、その maxRetries も同じ値に設定します。デフォルトは 1 であり、Agent の上限より低くなることがあります。あるリクエストで Processor だけを再試行メカニズムとして使用する場合、モデルの再試行は 0 にします。
違反 callback違反 callbackへの直接リンク
すべての Processor は onViolation property を公開します。これは、abort() が呼ばれた場合(block strategy)と Processor が警告を出した場合(warn strategy)のどちらでも、ポリシー違反が検出されるたびに実行されます。Processor の主要ロジックに影響を与えず、アラート、ログ、副作用に使用できます。
import { ModerationProcessor, CostGuardProcessor } from '@mastra/core/processors'
const moderation = new ModerationProcessor({
model: 'openai/gpt-5-nano',
strategy: 'block',
})
moderation.onViolation = ({ processorId, message, detail }) => {
// Log to external monitoring, send alerts, update dashboards
monitor.track('processor_violation', { processorId, message, detail })
}
const costGuard = new CostGuardProcessor({
maxCost: 10.0,
scope: 'resource',
window: '30d',
})
costGuard.onViolation = ({ processorId, message, detail }) => {
alertSystem.notify(`[${processorId}] ${message}`)
}
callback は次の内容を持つ ProcessorViolation オブジェクトを受け取ります。
processorId:違反を検出した Processor の IDmessage:違反内容を人が読める形式で表した説明detail:Processor 固有のメタデータ(コスト使用量、検出した PII の種類、モデレーションカテゴリなど)
onViolation は基底 Processor interfaceの一部なので、カスタム Processor でも使用できます。いずれかの Processor が abort() を呼ぶと、runner が自動的に実行します。Processor パイプラインへの干渉を防ぐため、callback 内でスローされたエラーは暗黙に捕捉されます。
Abort と tripwire チャンクAbort と tripwire チャンクへの直接リンク
abort(reason, options) を呼ぶと TripWire エラーがスローされ、処理が終了します。ストリームでは、Mastra がクライアントで検出可能な tripwire チャンクを送出します。
for await (const chunk of stream.fullStream) {
if (chunk.type === 'tripwire') {
console.log('Blocked by', chunk.payload.processorId, '-', chunk.payload.reason)
break
}
}
agent.generate() では、result.finishReason === 'other' のとき、同じ情報が result.tripwire として公開されます。
abort は2つ目の options 引数を受け取ります。
retry: trueは、終了せず再試行するよう Agent に要求します。入力・出力 Processor の再試行には、Agent または呼び出しでmaxProcessorRetriesを設定する必要があります。metadataは構造化データをtripwireチャンクに追加し、下流の consumer がpii、quality、moderationなどのカテゴリで分岐できるようにします。
API エラー処理API エラー処理への直接リンク
processAPIError メソッドは LLM API による拒否を処理します。これはネットワークやサーバー障害ではなく、API がリクエストを拒否するエラー(400 や 422 のステータスコードなど)です。API がメッセージ形式を拒否した場合に、リクエストを変更して再試行できます。
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'
export class ContextLengthHandler implements Processor {
id = 'context-length-handler'
processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
if (retryCount > 0) return
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}
Mastra には、Anthropic の「assistant message prefill」エラーを自動処理する組み込みの PrefillErrorHandler があります。この Processor は自動的に注入され、設定は不要です。
関連ドキュメント関連ドキュメントへの直接リンク
- Guardrails:セキュリティと検証の Processor
- Memory Processor:Memory 固有の Processor と自動連携
- Processor interface:Processor の完全な API reference
- ToolSearchProcessor reference:実行時 Tool 検索の API reference