跳至主要內容

Processor 介面

Processor 介面定義 Mastra 所有 processor 必須遵守的合約。Processor 可實作一個或多個方法,以處理 Agent 執行 pipeline 的不同階段。

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在開始時執行一次,位於 agentic loop 之前驗證或轉換最初的使用者輸入,加入 context
processInputStep在 agentic loop 的每個步驟、每次 LLM 呼叫之前在步驟之間轉換訊息及處理 Tool 結果
processLLMRequest轉換 LLM 請求後、呼叫 Provider 前改寫目前呼叫送出的 LanguageModelV2Prompt,但不持久保存變更
processAPIError當 LLM API 呼叫失敗時檢查 API 拒絕、按需要修改 state 或訊息,並要求重試
processOutputStreamLLM 回應期間的每個串流 chunk篩選或修改串流內容,並即時偵測模式
processLLMResponseLLM 步驟完成並收集串流 chunk 後擷取或快取完整回應,並執行與 processLLMRequest 配對的呼叫後副作用
processOutputStep每次 LLM 回應後、Tool 執行前驗證輸出質素,實作可重試的 guardrail
processToolResult每個 Tool 各自執行;在 tool.execute() 傳回後、結果加入訊息清單前掃描 Tool 輸出中的 prompt injection、遮蓋敏感欄位,並在違反政策時中止
processOutputResult生成完成後執行一次後處理最終回應並記錄結果

介面定義
介面定義 的直接連結

interface Processor<TId extends string = string, TTripwireMetadata = unknown> {
readonly id: TId
readonly name?: string
readonly description?: string
/** Index of this processor in the workflow (set at runtime when combining processors). */
processorIndex?: number
/** When true, processOutputStream also receives `data-*` chunks. Default: false. */
processDataParts?: boolean

/** Callback invoked when this processor detects a violation, regardless of strategy. */
onViolation?: (violation: ProcessorViolation) => void | Promise<void>

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

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

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

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

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

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

processOutputStep?(args: ProcessOutputStepArgs<TTripwireMetadata>): ProcessorMessageResult

processToolResult?(args: ProcessToolResultArgs<TTripwireMetadata>): ProcessorMessageResult

processOutputResult?(args: ProcessOutputResultArgs<TTripwireMetadata>): ProcessorMessageResult
}

屬性
屬性 的直接連結

id:

string
Processor 的唯一識別碼,用於 tracing 及除錯。

name?:

string
Processor 的選填顯示名稱。如未提供,會使用 id。

description?:

string
顯示於 tracing 及 Studio 的選填易讀描述。

processorIndex?:

number
Processor 在合併後 processor 清單中的位置。Mastra 將 processor 與 memory、Workspace 及每次呼叫覆寫設定合併時,會在執行期間設定此值;你無須自行設定。

processDataParts?:

boolean
設為 true 時,processOutputStream 方法亦會接收 Tool 透過 writer.custom() 發出的 data-* chunk。預設為 false。

onViolation?:

(violation: ProcessorViolation) => void | Promise<void>
當 processor 偵測到違反政策時呼叫的選填 callback,無論策略為封鎖還是警告均會執行。可用於發出警示、記錄至外部系統或向使用者傳送電郵等副作用。此 callback 擲出的錯誤會被靜默擷取,以免干擾 processor 邏輯。violation 物件包含 processorId、message,以及 processor 專用的 detail 欄位。

訊息引數
訊息引數 的直接連結

大部分 processor 方法都會同時接收 messagesmessageList。兩者指向相同的底層對話,但呈現方式不同。

messagesmessageList
messages-vs-messagelist 的直接連結

  • messages: 以目前階段為範圍的普通 MastraDBMessage 物件陣列。 processInputprocessInputStep 不包括系統訊息。 processOutputResultprocessOutputStep 包括最新的 LLM 回應。 此陣列由 messageList 支援,因此就地編輯訊息的 content.parts,下游 processor 及持久保存程序都會看見變更。
  • messageList: 支援該次執行的即時 MessageList instance。它提供經篩選的檢視(input、response、remembered、all)、多種輸出格式(db、ui、core),以及修改對話的方法。

只需讀取、逐項處理或輕微編輯目前階段訊息的欄位時,使用 messages。以下情況則使用 messageList

  • 讀取另一階段的訊息,例如處理輸出時讀取輸入訊息。
  • 新增、移除或取代整則訊息。
  • 轉換成其他格式,例如供第三方 API 使用的 UI 或 core 訊息。

messages 一律衍生自 messageList,因此修改 messageList 是新增、移除或重新排序訊息的標準方式。若就地編輯訊息內容(例如改寫 content.parts),直接修改 messages 的效果相同。如從 messages 傳回新陣列,Mastra 會針對目前階段將其與 messageList 協調一致。

持久保存
持久保存 的直接連結

啟用 memory 後,只有全部 processor 完成時最終保留在 messageList 的內容才會持久保存至儲存空間。以下兩種傳回方式的持久保存效果相同:

  • 直接修改 messageList(或傳回同一個 MessageList instance)時,已記錄的修改會就地套用,因此儲存的對話會反映你的變更。
  • 傳回 MastraDBMessage[]{ messages, systemMessages } 時,Mastra 會針對目前階段將傳回的陣列與 messageList 協調一致,移除缺少的訊息並取代系統訊息。

傳回另一個 MessageList instance 會造成錯誤。請一律修改傳入 processor 的 instance。

從訊息讀取文字
從訊息讀取文字 的直接連結

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 是主要來源。單一訊息可包含多個 part,包括 Tool 呼叫、Tool 結果及檔案 part 等非文字 part。 讀取 part.text 前,請先按 part.type === 'text' 篩選。
  • message.content.content 是為向後兼容而保留的扁平字串。只有在 parts 為空或缺少時才用作後備。
  • message.content 本身在 MastraDBMessage 上絕不會是普通字串。舊有 CoreMessage 結構可能是字串,但 processor 一律接收 MastraDBMessage

方法
方法 的直接連結

processInput
processinput 的直接連結

在輸入訊息傳送至 LLM 前處理訊息。Agent 開始執行時只會運行一次。

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

ProcessInputArgs
processinputargs 的直接連結

messages:

MastraDBMessage[]
要處理的使用者及 assistant 訊息(不包括系統訊息)。

systemMessages:

CoreMessage[]
所有系統訊息(Agent 指示、memory context、使用者提供的內容)。可修改後傳回。

messageList:

MessageList
用於進階訊息管理的完整 MessageList instance。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。它會擲出 TripWire 錯誤以停止執行。傳入 retry: true 可要求 LLM 使用回饋重試該步驟。

retryCount:

number
Processor 針對這次生成觸發重試的次數。用此值限制重試次數。Mastra 一律會傳入此值,初始值為 0。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

requestContext?:

RequestContext
以請求為範圍的 context,包含 threadId 及 resourceId 等執行 metadata。

ProcessInputResult
processinputresult 的直接連結

此方法可傳回以下三種類型之一:

MastraDBMessage[]:

array
轉換後的訊息陣列。系統訊息維持不變。

MessageList:

MessageList
傳入的同一個 messageList instance,表示你已直接修改它。

{ messages, systemMessages }:

object
同時包含轉換後訊息及修改後系統訊息的物件。

processInputStep
processinputstep 的直接連結

在 agentic loop 的每個步驟、輸入訊息傳送至 LLM 前處理訊息。processInput 只在開始時執行一次,此方法則會在每個步驟執行,包括延續 Tool 呼叫的步驟。

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

Agentic loop 中的執行次序
Agentic loop 中的執行次序 的直接連結

  1. processInput (開始時執行一次)
  2. processInputStep(來自 inputProcessors) (每個步驟、LLM 呼叫之前)
  3. prepareStep callback (作為 processInputStep pipeline 一部分執行,位於 inputProcessors 之後)
  4. processLLMRequest(來自 inputProcessors) (prompt 轉換後、Provider 呼叫前)
  5. LLM 執行
  6. processOutputStream(來自 outputProcessors) (每個串流 chunk)
  7. processLLMResponse(來自 inputProcessors) (串流完成後,與 processLLMRequest 配對)
  8. processOutputStep(來自 outputProcessors) (LLM 回應後、Tool 執行前)
  9. Tool 執行(如需要)
  10. 如有呼叫 Tool,從步驟 2 重複

ProcessInputStepArgs
processinputstepargs 的直接連結

messages:

MastraDBMessage[]
所有訊息,包括先前步驟的 Tool 呼叫及結果(唯讀快照)。

messageList:

MessageList
用於管理訊息的 MessageList instance。可直接修改,亦可在結果中傳回。

stepNumber:

number
目前步驟編號(從 0 開始)。步驟 0 是最初的 LLM 呼叫。

steps:

StepResult[]
先前步驟的結果,包括 text、toolCalls 及 toolResults。

systemMessages:

CoreMessage[]
所有系統訊息(唯讀快照)。在結果中傳回即可取代。

model:

MastraLanguageModelV2
目前使用的 model。在結果中傳回另一個 model 即可切換。

toolChoice?:

ToolChoice
目前的 Tool 選擇設定('auto'、'none'、'required' 或指定 Tool)。

activeTools?:

string[]
目前啟用的 Tool 名稱。傳回已篩選的陣列以限制 Tool。

tools?:

ToolSet
此步驟目前可用的 Tool。在結果中傳回以新增或取代 Tool。

providerOptions?:

SharedV2ProviderOptions
Provider 專用選項(例如 Anthropic cacheControl、OpenAI reasoningEffort)。

modelSettings?:

CallSettings
model 設定,例如 temperature、maxTokens、topP。

structuredOutput?:

StructuredOutputOptions
結構化輸出設定(schema、輸出模式)。在結果中傳回即可修改。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。它會擲出 TripWire 錯誤以停止執行。傳入 retry: true 可要求 LLM 使用回饋重試該步驟。

retryCount:

number
來自 ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

ProcessInputStepResult
processinputstepresult 的直接連結

processInputStep 可傳回多種結構:

  • ProcessInputStepResult 物件: 覆寫此步驟下列屬性的任何組合(詳見下文)。
  • MessageList: 傳回同一個 messageList instance,表示你已就地修改訊息。
  • MastraDBMessage[]: 傳回轉換後的訊息陣列,取代該步驟的訊息。
  • voidundefined: 不傳回任何內容,讓該步驟維持不變。

物件形式可傳回以下屬性的任何組合:

model?:

LanguageModelV2 | string
變更此步驟的 model。可以是 model instance 或類似 '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 instance(表示你已修改它)。不可與 messages 一併使用。

systemMessages?:

CoreMessage[]
只取代此步驟的所有系統訊息。

providerOptions?:

SharedV2ProviderOptions
變更此步驟的 Provider 專用選項。

modelSettings?:

CallSettings
變更此步驟的 model 設定。

structuredOutput?:

StructuredOutputOptions
變更此步驟的結構化輸出設定。

Processor 串接
Processor 串接 的直接連結

多個 processor 實作 processInputStep 時,會依次執行並串接變更:

Processor 1: receives { model: 'gpt-5.4' } → returns { model: 'gpt-5.4-mini' }
Processor 2: receives { model: 'gpt-5.4-mini' } → returns { toolChoice: 'none' }
Final: model = 'gpt-5.4-mini', toolChoice = 'none'

系統訊息隔離
系統訊息隔離 的直接連結

每個步驟開始時,系統訊息都會重設為原始值。在 processInputStep 作出的修改只影響目前步驟,不影響後續步驟。

使用情境
使用情境 的直接連結

  • 按步驟編號或 context 動態切換 model
  • 在指定步數後停用 Tool
  • 按對話 context 動態新增或取代 Tool
  • 在 Provider 之間轉換訊息 part 類型(例如為 Anthropic 將 reasoning 轉為 thinking
  • 按步驟編號或累積的 context 修改訊息
  • 新增步驟專用的系統指示
  • 按步驟調整 Provider 選項(例如快取控制)
  • 按步驟 context 修改結構化輸出 schema

processLLMRequest
processllmrequest 的直接連結

Mastra 將 MessageList 轉換為 LanguageModelV2Prompt 後、呼叫 Provider 前,此方法會處理最終 LLM 請求。可用於只應影響目前送出請求、可識別 model 的暫時改寫。

傳回的 prompt 變更只會轉送至目前呼叫的 model,不會持久保存回 MessageList、memory、UI 記錄或後續 Provider 呼叫。

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

ProcessLLMRequestArgs
processllmrequestargs 的直接連結

prompt:

LanguageModelV2Prompt
這次呼叫將傳送至 Provider 的 LLM 請求 prompt。

model:

MastraLanguageModel
將接收 prompt 的已解析 model。用此值限定 Provider 專用改寫的範圍。

stepNumber:

number
目前步驟編號(從 0 開始)。步驟 0 是最初的 LLM 呼叫。

steps:

StepResult[]
先前步驟的結果,包括 text、toolCalls 及 toolResults。

state:

Record<string, unknown>
每個 processor 獨立的 state,會在此請求內所有方法呼叫之間持續保留。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。它會擲出 TripWire 錯誤以停止執行,並發出 tripwire chunk。

retryCount:

number
來自 ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

writer?:

ProcessorStreamWriter
在串流期間發出自訂資料 chunk 的 stream writer。呼叫 writer.custom() 可發出 data-* chunk。

abortSignal?:

AbortSignal
用於取消操作的 signal。

傳回值
傳回值 的直接連結

processLLMRequest 傳回 ProcessLLMRequestResult,即 { prompt?: LanguageModelV2Prompt } | undefined | void

  • 傳回 { prompt } 以取代目前 Provider 呼叫送出的 prompt。
  • 傳回 undefinedvoid,以原樣轉送原始 prompt。

使用情境
使用情境 的直接連結

  • 在呼叫 model 前移除或重新建構 Provider 專用的 prompt part
  • 標準化角色或內容,以符合 Provider 的輸入要求
  • 在 loop 中途切換 Provider 時調整 Tool 結果格式

processLLMResponse
processllmresponse 的直接連結

此方法會在步驟完成(或重播快取回應)且 output processor 收集回應 chunk 後處理 LLM 回應。此 hook 與 processLLMRequest 配對:在 Provider 呼叫前用 processLLMRequest 暫存 state(例如 cache key),再用 processLLMResponse 處理已完成的回應(例如寫入快取)。

state 物件與同一步驟傳入 processLLMRequest 的 instance 相同,因此 processor 可關聯呼叫前後的工作。

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

ProcessLLMResponseArgs
processllmresponseargs 的直接連結

chunks:

CachedLLMStepChunk[]
此步驟由 LLM 呼叫產生(或從快取重播)的 chunk,採精簡格式({ type, payload })。

model:

MastraLanguageModel
產生(或原應產生)回應的 model。

stepNumber:

number
目前步驟編號(從 0 開始)。

steps:

StepResult[]
目前為止所有已完成的步驟,包括此步驟。

state:

Record<string, unknown>
與同一步驟的 processLLMRequest 共用、每個 processor 獨立的 state。可用它在兩個 hook 之間傳遞資料(例如 cache key)。

fromCache:

boolean
設為 true 時,表示回應是透過 processLLMRequest 傳回 { response } 從快取重播。寫入快取的 processor 應在此值為 true 時略過寫入。

warnings?:

LanguageModelV2CallWarning[]
language model 呼叫回報的警告(例如不支援的設定)。

request?:

unknown
Provider 請求 body(如有),可用於 tracing。

rawResponse?:

unknown
原始 Provider 回應(如有),可用於 tracing。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。 它會擲出 TripWire 錯誤以停止執行。

retryCount:

number
目前重試次數。初始值為 0;用於限制 processor 觸發的重試。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

writer?:

ProcessorStreamWriter
用於發出自訂資料 chunk 的 stream writer。

abortSignal?:

AbortSignal
用於取消操作的 signal。

傳回值
傳回值 的直接連結

processLLMResponse 傳回 ProcessLLMResponseResult,即 undefined | void。傳回值保留供日後擴充。

使用情境
使用情境 的直接連結

  • 即時呼叫後將 LLM 回應寫入快取(與 processLLMRequest 中衍生 cache key 的操作配對)
  • 記錄完整回應以供分析
  • 按已完成的回應觸發副作用

processAPIError
processapierror 的直接連結

在 LLM API 拒絕錯誤成為最終錯誤前處理它。當 API 呼叫因不可重試的錯誤(例如 400 或 422 狀態碼)失敗時執行。processOutputStep 在成功回應後執行,而此方法會在 API 拒絕請求時執行。

請將實作 processAPIError 的 processor 加入 Agent 的 errorProcessors 陣列。

Processor 可檢查錯誤並修改請求,例如在 messageList 加入訊息。傳回 { retry: true } 可使用修改後的 state 重試。

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

ProcessAPIErrorArgs
processapierrorargs 的直接連結

error:

unknown
LLM API 呼叫期間發生的錯誤。

messages:

MastraDBMessage[]
發生錯誤時的所有訊息。

messageList:

MessageList
用於管理訊息的 MessageList instance。修改它可在重試前變更請求。

stepNumber:

number
目前步驟編號(從 0 開始)。

steps:

StepResult[]
目前為止所有已完成的步驟。

state:

Record<string, unknown>
每個 processor 獨立的 state,會在此請求內所有方法呼叫之間持續保留。

retryCount:

number
error handler 目前的重試次數。用此值限制重試次數。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。

writer?:

ProcessorStreamWriter
在串流期間發出自訂資料 chunk 的 stream writer。呼叫 writer.custom() 可發出 data-* chunk。

requestContext?:

RequestContext
從 Agent 呼叫傳入的 request context。

abortSignal?:

AbortSignal
用於取消操作的 signal。

ProcessAPIErrorResult
processapierrorresult 的直接連結

retry:

boolean
套用修改後是否重試 LLM 呼叫。

使用情境
使用情境 的直接連結

  • 修改請求並重試,以處理 API 專用的拒絕
  • 透過修改請求,將不可重試的錯誤轉為可重試
  • 實作 model 專用的錯誤復原策略

範例: 自訂錯誤復原
範例: 自訂錯誤復原 的直接連結

src/mastra/processors/error-recovery.ts
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'

export class ErrorRecoveryProcessor implements Processor {
id = 'error-recovery'

processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
// Only retry once
if (retryCount > 0) return

// Check for a specific API error
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
// Trim older messages to fit within context
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}

processOutputStream
processoutputstream 的直接連結

使用內置 state 管理來處理串流輸出 chunk,讓 processor 可累積 chunk,並按更完整的 context 作出決定。

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

ProcessOutputStreamArgs
processoutputstreamargs 的直接連結

part:

ChunkType
目前正在處理的串流 chunk。

streamParts:

ChunkType[]
目前為止在串流中看見的所有 chunk。

state:

Record<string, unknown>
可修改且每個 processor 獨立的 state,在單一請求的每個 chunk 及每次方法呼叫之間持續保留。每次新的 generate 或 stream 呼叫都會建立全新的 state 物件。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止串流的函數。它會擲出 TripWire 錯誤以結束串流,並發出 tripwire chunk。傳入 retry: true 可要求 LLM 再嘗試,而非結束。

retryCount:

number
來自 ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。

messageList?:

MessageList
用於存取對話記錄的 MessageList instance。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

writer?:

ProcessorStreamWriter
用於向 client 發出自訂資料 chunk 的 stream writer。呼叫 writer.custom() 可發出 data-* 類型的 chunk;串流期間可用。

傳回值
傳回值 的直接連結

processOutputStream 傳回 Promise<ChunkType | null | undefined>

  • 傳回 ChunkType 以發出 chunk。傳回原始 part 可原樣發出;傳回新的 ChunkType 則可發出修改後的 chunk。
  • 傳回 null 以捨棄 chunk。不會向下一個 processor 或 client 傳送任何內容。
  • 傳回 undefined(包括 return; 陳述式或方法執行至結尾時隱含的 undefined)以捨棄 chunk。nullundefined 的行為相同。

捨棄 chunk 只影響該個 chunk。串流會繼續,下一個 chunk 仍會接受處理。要完全停止串流,請呼叫 abort()


processOutputResult
processoutputresult 的直接連結

在串流或生成完成後處理完整的輸出結果。

processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;

ProcessOutputResultArgs
processoutputresultargs 的直接連結

messages:

MastraDBMessage[]
已生成的回應訊息。

messageList:

MessageList
用於管理訊息的 MessageList instance。

state:

Record<string, unknown>
每個 processor 獨立的 state,會在此請求內所有方法呼叫之間持續保留。 與 processOutputStream 及其他方法共用。

result:

OutputResult
已解析的生成結果,包含 text(累積文字)、usage(包含 inputTokens、outputTokens、totalTokens 的 token 用量)、finishReason(生成結束原因)及 steps(所有 LLM 步驟結果,每項包含 toolCalls、toolResults、reasoning、sources、files 等)。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。它會擲出 TripWire 錯誤以停止執行,並發出 tripwire chunk。

retryCount:

number
來自 ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

writer?:

ProcessorStreamWriter
用於向 client 發出自訂資料 chunk 的 stream writer。呼叫 writer.custom() 可發出 data-* 類型的 chunk;串流期間可用。

processOutputStep
processoutputstep 的直接連結

在 agentic loop 中每次 LLM 回應後、Tool 執行前處理輸出。processOutputResult 只在最後執行一次,此方法則會在每個步驟執行,最適合實作可觸發重試的 guardrail。

processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;

ProcessOutputStepArgs
processoutputstepargs 的直接連結

messages:

MastraDBMessage[]
所有訊息,包括最新的 LLM 回應。

messageList:

MessageList
用於管理訊息的 MessageList instance。

stepNumber:

number
目前步驟編號(從 0 開始)。

finishReason?:

string
LLM 的結束原因(stop、tool-use、length 等)。

providerMetadata?:

ProviderMetadata
結束步驟的 Provider 專用 metadata(例如 AWS Bedrock guardrail trace)。當 model 步驟產生 Provider metadata 時便會提供,包括 steps 為空的 content-filter 封鎖情況。

toolCalls?:

ToolCallInfo[]
此步驟進行的 Tool 呼叫(如有)。

text?:

string
此步驟生成的文字。

usage:

LanguageModelUsage
目前步驟的 token 用量(inputTokensoutputTokenstotalTokens)。

systemMessages:

CoreMessage[]
所有系統訊息,可供讀取及修改。

steps:

StepResult[]
目前為止所有已完成的步驟,包括目前步驟。

state:

Record<string, unknown>
每個 processor 獨立的 state,會在此請求內所有方法呼叫之間持續保留。 與 processOutputStream 及 processOutputResult 共用。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止處理的函數。 傳入 retry: true 可要求 LLM 重試該步驟。

retryCount:

number
Processor 觸發重試的次數。用此值限制重試次數。Mastra 一律會傳入此值,初始值為 0。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

使用情境
使用情境 的直接連結

  • 實作可要求重試的質素 guardrail
  • 在 Tool 執行前驗證 LLM 輸出
  • 新增逐步記錄或指標
  • 實作可重試的輸出內容審核

範例: 可重試的質素 guardrail
範例: 可重試的質素 guardrail 的直接連結

src/mastra/processors/quality-guardrail.ts
import type { Processor } from '@mastra/core/processors'

export class QualityGuardrail implements Processor {
id = 'quality-guardrail'

async processOutputStep({ text, abort, retryCount }) {
const score = await evaluateResponseQuality(text)

if (score < 0.7) {
if (retryCount < 3) {
// Request retry with feedback for the LLM
abort('Response quality too low. Please provide more detail.', {
retry: true,
metadata: { qualityScore: score },
})
} else {
// Max retries reached, block the response
abort('Response quality too low after multiple attempts.')
}
}

return []
}
}

processToolResult
processtoolresult 的直接連結

tool.execute() 傳回後、結果加入訊息清單或送入下一次 LLM 呼叫前處理 Tool 結果。此方法與 Tool 執行前觸發的 processOutputStep 對稱。可用此方法掃描 Tool 輸出中的 prompt injection、遮蓋敏感欄位,或用 abort('reason', { retry: true }) 中止執行。

要取代 Tool 結果,請透過 messageList.updateToolInvocation 就地修改 messageList。runtime 會從訊息清單重新讀取經 processor 處理的結果,並在下游 Tool 結果串流 chunk 排入佇列前覆寫它,讓串流 client 看見處理後的值。

tool.execute() 擲出錯誤時不會觸發此方法;只有 Tool 成功執行並有結果可用時才會呼叫。

processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;

ProcessToolResultArgs
processtoolresultargs 的直接連結

messages:

MastraDBMessage[]
所有訊息,包括含有 Tool 呼叫的目前 assistant 訊息。

messageList:

MessageList
用於管理訊息的 MessageList instance。 呼叫 updateToolInvocation,以遮蓋或轉換後的值取代 Tool 結果。

stepNumber:

number
目前步驟編號(從 0 開始)。

toolName:

string
已執行 Tool 的名稱。

toolCallId:

string
這次特定 Tool 呼叫的唯一識別碼。

args:

unknown
LLM 傳入 Tool 的引數。

result:

unknown
Tool 傳回的值。對於 client 執行的 Tool,這是 tool.execute() 經過 ensureSerializable 後的輸出。對於 Provider 執行的 Tool(例如 Anthropic web_search),則是 Provider 串流的原始結果,不會經過 ensureSerializable

providerExecuted?:

boolean
此結果是否來自 Anthropic web_search 等由 Provider 執行的 Tool。對 client 執行的 Tool,預設為 undefined。

systemMessages:

CoreMessage[]
所有系統訊息,可供讀取。

steps:

StepResult[]
目前為止所有已完成的步驟。

state:

Record<string, unknown>
每個 processor 獨立的 state,會在此請求內所有方法呼叫之間持續保留。 與同一 processor 的其他方法共用。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
中止執行的函數。傳入 retry: true,可要求 LLM 以中止原因作為回饋來重試該步驟。

retryCount:

number
Processor 觸發重試的次數,初始值為 0。

tracingContext?:

TracingContext
用於可觀測性的 tracing context。

requestContext?:

RequestContext
以請求為範圍並包含執行 metadata 的 context。

使用情境
使用情境 的直接連結

  • 在 LLM 看見 Tool 輸出前掃描 prompt injection。
  • 遮蓋 Tool 傳回值中的敏感欄位(PII、秘密、憑證)。
  • 當 Tool 傳回違反政策的內容時中止執行。
  • 記錄或監測 Tool 傳回值,以供合規或審計。

範例: 遮蓋敏感欄位
範例: 遮蓋敏感欄位 的直接連結

src/mastra/processors/redact-tool-result.ts
import type { Processor } from '@mastra/core/processors'

export class RedactToolResult implements Processor {
id = 'redact-tool-result'

async processToolResult({ toolName, toolCallId, args, result, messageList }) {
if (toolName !== 'lookup-customer') return

const redacted = {
...(result as Record<string, unknown>),
ssn: '[REDACTED]',
email: '[REDACTED]',
}

messageList.updateToolInvocation({
type: 'tool-invocation',
toolInvocation: {
state: 'result',
toolCallId,
toolName,
args,
result: redacted,
},
})
}
}

範例: 封鎖 Tool 輸出中的 prompt injection
範例: 封鎖 Tool 輸出中的 prompt injection 的直接連結

src/mastra/processors/scan-tool-result.ts
import type { Processor } from '@mastra/core/processors'

export class ScanToolResult implements Processor {
id = 'scan-tool-result'

async processToolResult({ result, abort }) {
const text = typeof result === 'string' ? result : JSON.stringify(result)
if (containsPromptInjection(text)) {
abort('blocked by scan-tool-result: suspected prompt injection')
}
}
}

function containsPromptInjection(text: string): boolean {
return /ignore (all )?(previous|prior) instructions/i.test(text)
}

Processor 類型
Processor 類型 的直接連結

Mastra 提供類型別名,確保 processor 實作所需方法:

// Must implement processInput, processInputStep, processLLMRequest, or processLLMResponse (or any combination)
type InputProcessor = Processor &
(
| { processInput: required }
| { processInputStep: required }
| { processLLMRequest: required }
| { processLLMResponse: required }
)

// Must implement processOutputStream, processOutputStep, OR processOutputResult (or any combination)
type OutputProcessor = Processor &
(
| { processOutputStream: required }
| { processOutputStep: required }
| { processOutputResult: required }
)

// Must implement processAPIError
type ErrorProcessor = Processor & { processAPIError: required }

errorProcessors 設定實作 processAPIError 的 processor:

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

使用範例
使用範例 的直接連結

基本輸入 processor
基本輸入 processor 的直接連結

src/mastra/processors/lowercase.ts
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'

export class LowercaseProcessor implements Processor {
id = 'lowercase'

async processInput({ messages }): Promise<MastraDBMessage[]> {
return messages.map(msg => ({
...msg,
content: {
...msg.content,
parts: msg.content.parts?.map(part =>
part.type === 'text' ? { ...part, text: part.text.toLowerCase() } : part,
),
},
}))
}
}

使用 processInputStep 的逐步 processor
per-step-processor-with-processinputstep 的直接連結

src/mastra/processors/dynamic-model.ts
import type {
Processor,
ProcessInputStepArgs,
ProcessInputStepResult,
} from '@mastra/core/processors'

export class DynamicModelProcessor implements Processor {
id = 'dynamic-model'

async processInputStep({
stepNumber,
steps,
toolChoice,
}: ProcessInputStepArgs): Promise<ProcessInputStepResult> {
// Use a fast model for initial response
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}

// Switch to powerful model after tool calls
if (steps.length > 0 && steps[steps.length - 1].toolCalls?.length) {
return { model: 'openai/gpt-5.6-sol' }
}

// Disable tools after 5 steps to force completion
if (stepNumber > 5) {
return { toolChoice: 'none' }
}

return {}
}
}

使用 processInputStep 的訊息轉換 processor
message-transformer-with-processinputstep 的直接連結

src/mastra/processors/reasoning-transformer.ts
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'

export class ReasoningTransformer implements Processor {
id = 'reasoning-transformer'

async processInputStep({ messages, messageList }) {
// Transform reasoning parts to thinking parts at each step
// This is useful when switching between model providers
for (const msg of messages) {
if (msg.role === 'assistant' && msg.content.parts) {
for (const part of msg.content.parts) {
if (part.type === 'reasoning') {
;(part as any).type = 'thinking'
}
}
}
}
return messageList
}
}

混合 processor(輸入及輸出)
混合 processor(輸入及輸出) 的直接連結

src/mastra/processors/content-filter.ts
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
import type { ChunkType } from '@mastra/core/stream'

export class ContentFilter implements Processor {
id = 'content-filter'
private blockedWords: string[]

constructor(blockedWords: string[]) {
this.blockedWords = blockedWords
}

async processInput({ messages, abort }): Promise<MastraDBMessage[]> {
for (const msg of messages) {
const text = msg.content.parts
?.filter(p => p.type === 'text')
.map(p => p.text)
.join(' ')

if (this.blockedWords.some(word => text?.includes(word))) {
abort('Blocked content detected in input')
}
}
return messages
}

async processOutputStream({ part, abort }): Promise<ChunkType | null> {
if (part.type === 'text-delta') {
if (this.blockedWords.some(word => part.payload.text.includes(word))) {
abort('Blocked content detected in output')
}
}
return part
}
}

使用 state 的串流累加器
使用 state 的串流累加器 的直接連結

src/mastra/processors/word-counter.ts
import type { Processor } from '@mastra/core/processors'
import type { ChunkType } from '@mastra/core/stream'

export class WordCounter implements Processor {
id = 'word-counter'

async processOutputStream({ part, state }): Promise<ChunkType> {
// Initialize state on first chunk
if (!state.wordCount) {
state.wordCount = 0
}

// Count words in text chunks
if (part.type === 'text-delta') {
const words = part.payload.text.split(/\s+/).filter(Boolean)
state.wordCount += words.length
}

// Log word count on finish
if (part.type === 'finish') {
console.log(`Total words: ${state.wordCount}`)
}

return part
}
}

State 生命週期
State 生命週期 的直接連結

每個 processor 都會在 processLLMRequestprocessLLMResponseprocessOutputStreamprocessOutputStepprocessOutputResultprocessAPIError 中接收 state 物件。State 有三項重要特性:

  • 每個 processor 獨立: 每個 processor 都有自己的 state 物件,以 processor 的 id 作為 key。不同 id 的 processor 無法讀取或覆寫彼此的 state。
  • 每個請求獨立: 每次呼叫 agent.generate()agent.stream() 時都會建立新的 state 物件。State 不會在請求或使用者之間洩漏。
  • 跨方法共用: 在同一個請求內,同一個 state 物件會傳入 processLLMRequest(Provider 呼叫前)、processLLMResponse(步驟完成後)、processOutputStream(每個 chunk)、processOutputStep(每個 LLM 步驟後)、processOutputResult(最後一次)及 processAPIError(LLM 呼叫失敗時)。例如,processLLMRequest 可暫存 cache key,而 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 chunk
中止及 tripwire chunk 的直接連結

每個方法的 abort 函數都會擲出 TripWire 錯誤,以停止處理並在輸出串流發出 tripwire chunk。Client 可偵測該 chunk,以區分被封鎖的回應與正常完成。

abort('Blocked content detected', { retry: false, metadata: { category: 'pii' } })
  • reason: 易於理解的說明,會顯示為 tripwire.payload.reason
  • retry: 設為 true 時,Agent 會重試相同步驟,並將 reason 作為回饋。只有在 Agent 或呼叫上設定 maxProcessorRetries 時才會重試,否則請求會中止。設定 errorProcessors 後,該次呼叫的 maxProcessorRetries 預設為 10
  • metadata: 附加至 tripwire chunk、供下游使用者使用的選填結構化資料。

發出的 tripwire chunk 結構如下:

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

在非串流呼叫(agent.generate())中,結果會透過 result.tripwireresult.finishReason === 'other' 提供相同資料。

發出自訂資料 chunk
發出自訂資料 chunk 的直接連結

可存取 writer 的 processor 能呼叫 writer.custom(chunk),向 client 串流傳送自訂 data-* chunk。Tool 也可透過自己的 writer 執行相同操作。這是 processor 在一般文字及 Tool chunk 以外發出內容的唯一方式。

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

設定 memory 後,從 processOutputStreamprocessOutputResult 發出的自訂 data-* chunk 會儲存為 assistant 訊息的 part。在 chunk 物件設定 transient: true,即可串流傳送而不儲存至 memory:

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

請將 transient 作為 chunk 的屬性傳入,而非 writer.custom() 的第二個引數。第二個引數包含 messageId 等 writer 選項。

預設情況下,processor 在 processOutputStream不會看見 data-* chunk,以免意外處理 Tool telemetry 或自己的輸出。在 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
}
}

Chunk 的 type 必須以 data- 開頭,才會視為自訂資料 chunk。從 processOutputStream 傳回 nullundefined 仍會捨棄 chunk,因此 processor 可像篩選文字 chunk 一樣檢查、修改或篩選自訂資料。

在 Agent 設定 processor
在 Agent 設定 processor 的直接連結

Processor 透過三個陣列附加至 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 回應期間或之後執行,接收輸出 chunk 或訊息。
  • errorProcessors: 在 LLM API 呼叫擲出錯誤時執行,接收原始錯誤。

每個陣列亦接受函數,讓 processor 可按請求從 RequestContext 建立:

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

每次呼叫的覆寫設定
每次呼叫的覆寫設定 的直接連結

agent.generate()agent.stream() 接受 inputProcessorsoutputProcessorserrorProcessorsmaxProcessorRetries。呼叫時設定任何 processor 陣列,都會為該請求取代 Agent 上相應的陣列。Mastra 自動加入的 memory、Workspace、Skill、channel 及 browser processor 一律會保留,並在你的陣列前後執行。

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

呼叫時傳入的 maxProcessorRetries 會覆寫 Agent 預設值。如兩者均未設定,processor 要求的重試會視為中止。