Processor 介面
Processor 介面定義了 Mastra 中所有 processor 必須遵循的規範。Processor 可實作一或多個方法,以處理 Agent 執行管線的不同階段。
Processor 方法的執行時機「Processor 方法的執行時機」的直接連結
Processor 方法會在 Agent 執行生命週期的不同時間點執行:
┌────────────────────────────────────────────────────────────────────┐
│ Agent Execution Flow │
├────────────────────────────────────────────────────────────────────┤
│ │
│ User Input │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ processInput │ ← Runs ONCE at start │
│ └───────────┬────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────────────┐ │
│ │ Agentic Loop │ │
│ │ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processInputStep │ ← Runs at EACH step │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processLLMRequest │ ← Before provider call │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ LLM Execution ──── API Error? ───┐ │ │
│ │ │ │ │ │
│ │ │ ┌───────────┴──────────┐ │ │
│ │ │ │ processAPIError │ │ │
│ │ │ └──────────────────────┘ │ │
│ │ │ (retry loops back to LLM) │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processOutputStream │ ← Runs on EACH stream chunk │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processLLMResponse │ ← After stream completes │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processOutputStep │ ← Runs after EACH LLM step │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ Tool Execution (if needed) │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processToolResult │ ← Runs per tool, after each │ │
│ │ └───────────┬────────────┘ tool.execute() returns │ │
│ │ │ │ │
│ │ └──────── Loop back if tools called ────────────│ │
│ │ │ │
│ └──────────────────────────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ processOutputResult │ ← Runs ONCE after completion │
│ └────────────────────────┘ │
│ │ │
│ ▼ │
│ Final Response │
│ │
└────────────────────────────────────────────────────────────────────┘
| 方法 | 執行時機 | 使用情境 |
|---|---|---|
processInput | 開始時執行一次,位於 Agentic loop 之前 | 驗證或轉換使用者最初的輸入,以及新增內容脈絡 |
processInputStep | Agentic loop 的每個步驟,在每次 LLM 呼叫前 | 在步驟之間轉換訊息及處理 Tool 結果 |
processLLMRequest | LLM 請求完成轉換後、呼叫 Provider 前 | 改寫目前呼叫對外送出的 LanguageModelV2Prompt,且不保存變更 |
processAPIError | LLM API 呼叫失敗時 | 檢查 API 拒絕原因、視需要修改狀態或訊息,並要求重試 |
processOutputStream | LLM 回應期間的每個串流區塊 | 篩選或修改串流內容,並即時偵測模式 |
processLLMResponse | LLM 步驟完成並收集串流區塊後 | 擷取或快取完整回應,並執行與 processLLMRequest 配對的呼叫後副作用 |
processOutputStep | 每次 LLM 回應後、執行 Tool 前 | 驗證輸出品質,以及實作可重試的防護機制 |
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:
name?:
description?:
processorIndex?:
processDataParts?:
data-* 區塊。預設為 false。onViolation?:
訊息引數「訊息引數」的直接連結
大多數 processor 方法都會同時接收 messages 與 messageList。兩者指向相同的底層對話,但呈現方式不同。
messages 與 messageList 的比較「messages-vs-messagelist」的直接連結
messages:由MastraDBMessage物件組成的普通陣列,範圍限定於目前階段。對processInput與processInputStep而言,不包含系統訊息;對processOutputResult與processOutputStep而言,則包含最新的 LLM 回應。此陣列由messageList支援,因此直接編輯訊息的content.parts,下游 processor 與保存作業都會看到變更。messageList:支援此次執行的即時MessageList執行個體。它提供經篩選的檢視(輸入、回應、記憶及全部)、多種輸出格式(db、ui 及 core),以及修改對話的方法。
若只需讀取、映射或小幅編輯目前階段訊息中的欄位,請使用 messages。若有下列需求,請使用 messageList:
- 讀取其他階段的訊息,例如處理輸出時讀取輸入訊息。
- 新增、移除或替換整則訊息。
- 轉換成其他格式,例如供第三方 API 使用的 UI 或 core 訊息。
messages 一律衍生自 messageList,因此要新增、移除或重新排序訊息,修改 messageList 才是標準做法。若要直接編輯訊息內容(例如改寫 content.parts),直接修改 messages 的效果相同。若從 messages 傳回新陣列,Mastra 會依目前階段將它與 messageList 協調一致。
保存「保存」的直接連結
啟用記憶體時,只有在所有 processor 完成後最終留在 messageList 中的內容會保存至儲存空間。以下兩種傳回方式的保存結果相同:
- 直接修改
messageList(或傳回相同的MessageList執行個體)時,記錄的修改會就地套用,因此儲存的對話會反映變更。 - 傳回
MastraDBMessage[]或{ messages, systemMessages }時,Mastra 會依目前階段將傳回的陣列與messageList協調一致,移除缺少的訊息並替換系統訊息。
傳回不同的 MessageList 執行個體會導致錯誤。請一律修改傳給 processor 的執行個體。
從訊息讀取文字「從訊息讀取文字」的直接連結
MastraDBMessage.content 使用結構化物件,不支援字串。讀取使用者或助理文字的標準方式是使用 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是主要來源。單則訊息可包含多個部分,包括 Tool 呼叫、Tool 結果及檔案部分等非文字部分。讀取part.text前,請先以part.type === 'text'篩選。message.content.content是為了向後相容而保留的扁平化字串。僅在parts為空或不存在時用作備援。MastraDBMessage上的message.content本身絕不會是普通字串。舊版CoreMessage的資料結構可能是字串,但 processor 一律接收MastraDBMessage。
方法「方法」的直接連結
processInput「processinput」的直接連結
在輸入訊息傳送至 LLM 前進行處理。此方法在 Agent 開始執行時執行一次。
processInput?(args: ProcessInputArgs): Promise<ProcessInputResult> | ProcessInputResult;
ProcessInputArgs「processinputargs」的直接連結
messages:
systemMessages:
messageList:
abort:
retry: true 可要求 LLM 根據意見回饋重試此步驟。retryCount:
tracingContext?:
requestContext?:
ProcessInputResult「processinputresult」的直接連結
此方法可傳回以下三種類型之一:
MastraDBMessage[]:
MessageList:
{ messages, systemMessages }:
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 中的執行順序」的直接連結
processInput(開始時執行一次)- inputProcessors 的
processInputStep(每個步驟執行,位於 LLM 呼叫前) prepareStep回呼(作為 processInputStep 管線的一部分,在 inputProcessors 後執行)- inputProcessors 的
processLLMRequest(提示完成轉換後、呼叫 Provider 前) - 執行 LLM
- outputProcessors 的
processOutputStream(每個串流區塊) - inputProcessors 的
processLLMResponse(串流完成後,與processLLMRequest配對) - outputProcessors 的
processOutputStep(LLM 回應後、執行 Tool 前) - 執行 Tool(如有需要)
- 若呼叫了 Tool,則從步驟 2 重新執行
ProcessInputStepArgs「processinputstepargs」的直接連結
messages:
messageList:
stepNumber:
steps:
systemMessages:
model:
toolChoice?:
activeTools?:
tools?:
providerOptions?:
modelSettings?:
structuredOutput?:
abort:
retry: true 可要求 LLM 根據意見回饋重試此步驟。retryCount:
ProcessorContext 的目前重試次數。從 0 開始;可用來限制由 processor 觸發的重試。tracingContext?:
requestContext?:
ProcessInputStepResult「processinputstepresult」的直接連結
processInputStep 可傳回多種資料結構:
ProcessInputStepResult物件:為此步驟覆寫下列屬性的任意組合(後續會說明)。MessageList:傳回相同的messageList執行個體,表示你已就地修改訊息。MastraDBMessage[]:傳回轉換後的訊息陣列,以替換此步驟的訊息。void或undefined:不傳回任何內容,讓此步驟維持不變。
物件形式可傳回以下屬性的任意組合:
model?:
toolChoice?:
activeTools?:
tools?:
messages?:
messageList?:
systemMessages?:
providerOptions?:
modelSettings?:
structuredOutput?:
Processor 串接「Processor 串接」的直接連結
多個 processor 實作 processInputStep 時,會依序執行,且變更會串接傳遞:
Processor 1: receives { model: 'gpt-5.4' } → returns { model: 'gpt-5.4-mini' }
Processor 2: receives { model: 'gpt-5.4-mini' } → returns { toolChoice: 'none' }
Final: model = 'gpt-5.4-mini', toolChoice = 'none'
系統訊息隔離「系統訊息隔離」的直接連結
每個步驟開始時,系統訊息都會重設為原始值。在 processInputStep 中所做的修改只會影響目前步驟,不會影響後續步驟。
使用情境「使用情境」的直接連結
- 根據步驟編號或內容脈絡動態切換模型
- 執行一定數量的步驟後停用 Tool
- 根據對話內容脈絡動態新增或替換 Tool
- 在不同 Provider 之間轉換訊息部分的類型(例如針對 Anthropic 將
reasoning轉為thinking) - 根據步驟編號或累積的內容脈絡修改訊息
- 新增步驟專屬的系統指令
- 調整每個步驟的 Provider 選項(例如快取控制)
- 根據步驟內容脈絡修改結構化輸出 schema
processLLMRequest「processllmrequest」的直接連結
Mastra 將 MessageList 轉換成 LanguageModelV2Prompt 後、呼叫 Provider 前,此方法會處理最終的 LLM 請求。適合用於只應影響目前對外請求、依模型而定的暫時性改寫。
傳回的提示變更只會轉送給模型供目前呼叫使用,不會保存回 MessageList、記憶體、UI 歷程記錄或後續 Provider 呼叫。
processLLMRequest?(
args: ProcessLLMRequestArgs,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult;
ProcessLLMRequestArgs「processllmrequestargs」的直接連結
prompt:
model:
stepNumber:
steps:
state:
abort:
tripwire 區塊。retryCount:
ProcessorContext 的目前重試次數。從 0 開始;可用來限制由 processor 觸發的重試。requestContext?:
tracingContext?:
writer?:
writer.custom() 即可發出 data-* 區塊。abortSignal?:
傳回值「傳回值」的直接連結
processLLMRequest 會傳回 ProcessLLMRequestResult,其類型為 { prompt?: LanguageModelV2Prompt } | undefined | void。
- 傳回
{ prompt },即可替換目前 Provider 呼叫對外送出的提示。 - 傳回
undefined或void,則會原封不動地轉送原始提示。
使用情境「使用情境」的直接連結
- 在呼叫模型前移除 Provider 專屬的提示部分,或調整其資料結構
- 將角色或內容正規化,以符合 Provider 的輸入需求
- 在迴圈中途切換 Provider 時調整 Tool 結果格式
processLLMResponse「processllmresponse」的直接連結
在步驟完成(或重播快取回應)且輸出 processor 收集完回應區塊後,處理 LLM 回應。此 hook 與 processLLMRequest 配對:在呼叫 Provider 前使用 processLLMRequest 暫存狀態(例如快取鍵),並使用 processLLMResponse 對完成的回應執行作業(例如寫入快取)。
state 物件與同一步驟傳給 processLLMRequest 的執行個體相同,因此 processor 可將呼叫前後的作業建立關聯。
processLLMResponse?(
args: ProcessLLMResponseArgs,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult;
ProcessLLMResponseArgs「processllmresponseargs」的直接連結
chunks:
{ type, payload })。model:
stepNumber:
steps:
state:
processLLMRequest 共用的 processor 專屬狀態。可用來在兩個 hook 之間傳遞資料(例如快取鍵)。fromCache:
true,表示回應是透過 processLLMRequest 傳回 { response },從快取重播而來。寫入快取的 processor 應在此值為 true 時略過寫入。warnings?:
request?:
rawResponse?:
abort:
retryCount:
0 開始;可用來限制由 processor 觸發的重試。requestContext?:
tracingContext?:
writer?:
abortSignal?:
傳回值「傳回值」的直接連結
processLLMResponse 會傳回 ProcessLLMResponseResult,其類型為 undefined | void。傳回值保留供未來擴充使用。
使用情境「使用情境」的直接連結
- 即時呼叫後將 LLM 回應寫入快取(與
processLLMRequest中衍生快取鍵的作業配對) - 記錄完整回應,以供分析使用
- 根據完成的回應觸發副作用
processAPIError「processapierror」的直接連結
在 LLM API 拒絕錯誤成為最終錯誤前進行處理。當 API 呼叫因不可重試的錯誤(例如 400 或 422 狀態碼)而失敗時,此方法便會執行。processOutputStep 會在成功回應後執行,而此方法會在 API 拒絕請求時執行。
請將實作 processAPIError 的 processor 加入 Agent 的 errorProcessors 陣列。
Processor 可檢查錯誤並修改請求,例如將訊息附加到 messageList。傳回 { retry: true } 即可使用修改後的狀態重試。
processAPIError?(args: ProcessAPIErrorArgs): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void;
ProcessAPIErrorArgs「processapierrorargs」的直接連結
error:
messages:
messageList:
stepNumber:
steps:
state:
retryCount:
abort:
writer?:
writer.custom() 即可發出 data-* 區塊。requestContext?:
abortSignal?:
ProcessAPIErrorResult「processapierrorresult」的直接連結
retry:
使用情境「使用情境」的直接連結
- 修改請求並重試,以處理 API 專屬的拒絕錯誤
- 修改請求,將不可重試的錯誤轉為可重試錯誤
- 實作模型專屬的錯誤復原策略
範例:自訂錯誤復原「範例:自訂錯誤復原」的直接連結
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'
export class ErrorRecoveryProcessor implements Processor {
id = 'error-recovery'
processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
// Only retry once
if (retryCount > 0) return
// Check for a specific API error
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
// Trim older messages to fit within context
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}
processOutputStream「processoutputstream」的直接連結
使用內建狀態管理功能處理串流輸出區塊。Processor 可累積區塊,並根據較完整的內容脈絡做出決策。
processOutputStream?(args: ProcessOutputStreamArgs): Promise<ChunkType | null | undefined>;
ProcessOutputStreamArgs「processoutputstreamargs」的直接連結
part:
streamParts:
state:
abort:
tripwire 區塊。傳入 retry: true 可要求 LLM 再次嘗試,而非結束執行。retryCount:
ProcessorContext 的目前重試次數。從 0 開始;可用來限制由 processor 觸發的重試。messageList?:
tracingContext?:
requestContext?:
writer?:
傳回值「傳回值」的直接連結
processOutputStream 會傳回 Promise<ChunkType | null | undefined>。
- 傳回
ChunkType即可發出區塊。傳回原始part會原封不動地發出,傳回新的ChunkType則會發出修改後的區塊。 - 傳回
null會捨棄區塊,不會傳送任何內容給下一個 processor 或使用者端。 - 傳回
undefined(包括return;陳述式隱含的undefined,或方法執行至結尾而未傳回值)會捨棄區塊。null與undefined的行為相同。
捨棄區塊只會影響該單一區塊。串流會繼續,下一個區塊仍會受到處理。若要完全停止串流,請呼叫 abort()。
processOutputResult「processoutputresult」的直接連結
串流或內容產生完成後,處理完整的輸出結果。
processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;
ProcessOutputResultArgs「processoutputresultargs」的直接連結
messages:
messageList:
state:
result:
text(累積文字)、usage(token 用量,包括 inputTokens、outputTokens、totalTokens)、finishReason(內容產生結束的原因),以及 steps(所有 LLM 步驟結果,每個結果都包含 toolCalls、toolResults、reasoning、sources、files 等)。abort:
tripwire 區塊。retryCount:
ProcessorContext 的目前重試次數。從 0 開始;可用來限制由 processor 觸發的重試。tracingContext?:
requestContext?:
writer?:
processOutputStep「processoutputstep」的直接連結
在 Agentic loop 中每次 LLM 回應後、執行 Tool 前處理輸出。processOutputResult 只在結束時執行一次,而此方法會在每個步驟執行。這是實作可觸發重試之防護機制的理想方法。
processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;
ProcessOutputStepArgs「processoutputstepargs」的直接連結
messages:
messageList:
stepNumber:
finishReason?:
providerMetadata?:
steps 為空,也會提供。toolCalls?:
text?:
usage:
inputTokens、outputTokens、totalTokens)。systemMessages:
steps:
state:
abort:
retry: true 可要求 LLM 重試此步驟。retryCount:
tracingContext?:
requestContext?:
使用情境「使用情境」的直接連結
- 實作可要求重試的品質防護機制
- 執行 Tool 前驗證 LLM 輸出
- 為每個步驟新增記錄或指標
- 實作可重試的輸出審核
範例:可重試的品質防護機制「範例:可重試的品質防護機制」的直接連結
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() 傳回 Tool 結果後、該結果加入訊息清單或傳給下一次 LLM 呼叫前進行處理。此方法與在 Tool 執行前觸發的 processOutputStep 相互對應。可用來掃描 Tool 輸出是否有 prompt injection、遮蔽敏感欄位,或使用 abort('reason', { retry: true }) 中止執行。
若要替換 Tool 結果,請透過 messageList.updateToolInvocation 就地修改 messageList。執行階段會在 processor 執行後重新從訊息清單讀取結果,並在排入佇列前覆寫下游的 Tool 結果串流區塊,因此串流使用者端會看到處理後的值。
若 tool.execute() 擲出錯誤,此方法不會觸發;只有成功執行 Tool 且有可用結果時才會呼叫。
processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;
ProcessToolResultArgs「processtoolresultargs」的直接連結
messages:
messageList:
updateToolInvocation,即可使用遮蔽或轉換後的值替換 Tool 結果。stepNumber:
toolName:
toolCallId:
args:
result:
ensureSerializable 處理後的 tool.execute() 輸出。若 Tool 由 Provider 執行(例如 Anthropic web_search),此值是 Provider 串流的原始結果,不會經過 ensureSerializable 處理。providerExecuted?:
systemMessages:
steps:
state:
abort:
retry: true,可要求 LLM 以中止原因作為意見回饋來重試此步驟。retryCount:
tracingContext?:
requestContext?:
使用情境「使用情境」的直接連結
- 在 LLM 看到 Tool 輸出前,掃描其中是否有 prompt injection。
- 遮蔽 Tool 傳回內容中的敏感欄位(PII、密鑰、認證資訊)。
- Tool 傳回的內容違反政策時中止執行。
- 基於法規遵循或稽核需求,記錄或檢測 Tool 傳回內容。
範例:遮蔽敏感欄位「範例:遮蔽敏感欄位」的直接連結
import type { Processor } from '@mastra/core/processors'
export class RedactToolResult implements Processor {
id = 'redact-tool-result'
async processToolResult({ toolName, toolCallId, args, result, messageList }) {
if (toolName !== 'lookup-customer') return
const redacted = {
...(result as Record<string, unknown>),
ssn: '[REDACTED]',
email: '[REDACTED]',
}
messageList.updateToolInvocation({
type: 'tool-invocation',
toolInvocation: {
state: 'result',
toolCallId,
toolName,
args,
result: redacted,
},
})
}
}
範例:封鎖 Tool 輸出中的 prompt injection「範例:封鎖 Tool 輸出中的 prompt injection」的直接連結
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」的直接連結
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」的直接連結
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」的直接連結
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class ReasoningTransformer implements Processor {
id = 'reasoning-transformer'
async processInputStep({ messages, messageList }) {
// Transform reasoning parts to thinking parts at each step
// This is useful when switching between model providers
for (const msg of messages) {
if (msg.role === 'assistant' && msg.content.parts) {
for (const part of msg.content.parts) {
if (part.type === 'reasoning') {
;(part as any).type = 'thinking'
}
}
}
}
return messageList
}
}
混合 processor(輸入與輸出)「混合 processor(輸入與輸出)」的直接連結
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
import type { ChunkType } from '@mastra/core/stream'
export class ContentFilter implements Processor {
id = 'content-filter'
private blockedWords: string[]
constructor(blockedWords: string[]) {
this.blockedWords = blockedWords
}
async processInput({ messages, abort }): Promise<MastraDBMessage[]> {
for (const msg of messages) {
const text = msg.content.parts
?.filter(p => p.type === 'text')
.map(p => p.text)
.join(' ')
if (this.blockedWords.some(word => text?.includes(word))) {
abort('Blocked content detected in input')
}
}
return messages
}
async processOutputStream({ part, abort }): Promise<ChunkType | null> {
if (part.type === 'text-delta') {
if (this.blockedWords.some(word => part.payload.text.includes(word))) {
abort('Blocked content detected in output')
}
}
return part
}
}
使用狀態的串流累加 processor「使用狀態的串流累加 processor」的直接連結
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
}
}
狀態生命週期「狀態生命週期」的直接連結
每個 processor 都會在 processLLMRequest、processLLMResponse、processOutputStream、processOutputStep、processOutputResult 與 processAPIError 中接收 state 物件。狀態有三個重要特性:
- 各 processor 獨立:每個 processor 都會取得自己的
state物件,並以 processor 的id作為索引鍵。ID 不同的 processor 無法讀取或覆寫彼此的狀態。 - 各請求獨立:每次呼叫
agent.generate()或agent.stream()時,都會在開始時建立新的狀態物件。狀態不會在不同請求或不同使用者之間洩漏。 - 跨方法共用:在單一請求中,相同的
state物件會傳給processLLMRequest(呼叫 Provider 前)、processLLMResponse(步驟完成後)、processOutputStream(每個區塊)、processOutputStep(每個 LLM 步驟後)、processOutputResult(結束時執行一次)及processAPIError(LLM 呼叫失敗時)。例如,processLLMRequest可暫存快取鍵,processLLMResponse再讀取該鍵以寫入回應。
由於 state 一開始是空物件,請在第一次存取時以防禦性方式初始化欄位:
import type { Processor } from '@mastra/core/processors'
export class WordCounter implements Processor {
id = 'word-counter'
async processOutputStream({ part, state }) {
state.wordCount ??= 0
if (part.type === 'text-delta') {
state.wordCount += part.payload.text.split(/\s+/).filter(Boolean).length
}
return part
}
}
中止與 tripwire 區塊「中止與 tripwire 區塊」的直接連結
每個方法上的 abort 函式都會擲出 TripWire 錯誤,停止處理並在輸出串流中發出 tripwire 區塊。使用者端可偵測此區塊,以區分遭封鎖的回應與正常結束。
abort('Blocked content detected', { retry: false, metadata: { category: 'pii' } })
reason:方便閱讀的說明,會出現在tripwire.payload.reason。retry:若為true,Agent 會重試相同步驟,並將reason作為意見回饋。只有在 Agent 或呼叫上設定maxProcessorRetries時才會進行重試,否則請求會中止。若已設定errorProcessors,該次呼叫的maxProcessorRetries預設為10。metadata:附加至tripwire區塊的選用結構化資料,供下游使用者使用。
發出的 tripwire 區塊具有以下資料結構:
type TripwireChunk = {
type: 'tripwire'
runId: string
from: 'AGENT'
payload: {
reason: string
retry?: boolean
metadata?: unknown
processorId: string
}
}
在非串流呼叫(agent.generate())中,結果會透過 result.tripwire 及 result.finishReason === 'other' 提供相同資訊。
發出自訂資料區塊「發出自訂資料區塊」的直接連結
可存取 writer 的 processor 能呼叫 writer.custom(chunk),將自訂 data-* 區塊串流傳送至使用者端。Tool 也能透過自己的 writer 執行相同作業。這是 processor 在一般文字及 Tool 區塊以外發出內容的唯一方式。
await writer.custom({
type: 'data-moderation',
runId,
from: 'AGENT',
data: { level: 'warn', reason: 'Possibly unsafe' },
})
設定記憶體時,從 processOutputStream 或 processOutputResult 發出的自訂 data-* 區塊,會儲存為助理訊息的一部分。在區塊物件上設定 transient: true,即可串流傳送而不儲存至記憶體:
await writer.custom({
type: 'data-progress',
data: { status: 'Processing' },
transient: true,
})
請將 transient 作為區塊的屬性傳入,不要當作 writer.custom() 的第二個引數。第二個引數包含 messageId 等 writer 選項。
依預設,processor 在 processOutputStream 中不會看到 data-* 區塊,以免意外處理 Tool 遙測資訊或自己產生的輸出。請在 processor 上設定 processDataParts: true 以選擇加入:
class ModerationCollector implements Processor {
id = 'moderation-collector'
processDataParts = true
async processOutputStream({ part, state }) {
if (part.type === 'data-moderation') {
state.warnings ??= []
state.warnings.push(part.data)
}
return part
}
}
區塊 type 必須以 data- 開頭,才會視為自訂資料區塊。從 processOutputStream 傳回 null 或 undefined 仍會捨棄區塊,因此 processor 可使用篩選文字區塊的相同方式,檢查、修改或篩選自訂資料。
在 Agent 上設定 processor「在 Agent 上設定 processor」的直接連結
Processor 會透過三個陣列附加至 Agent:
import { Agent } from '@mastra/core/agent'
import { PrefillErrorHandler } from '@mastra/core/processors'
const agent = new Agent({
id: 'support-agent',
name: 'support-agent',
model: 'openai/gpt-5',
instructions: '...',
inputProcessors: [new ContentFilter(['secret'])],
outputProcessors: [new WordCounter()],
errorProcessors: [new PrefillErrorHandler()],
maxProcessorRetries: 3,
})
inputProcessors:在 LLM 前執行,接收輸入訊息。outputProcessors:在 LLM 回應期間或回應後執行,接收輸出區塊或訊息。errorProcessors:LLM API 呼叫擲出錯誤時執行,接收原始錯誤。
每個陣列也接受函式,因此能依每次請求使用 RequestContext 建立 processor:
new Agent({
id: 'processor-interface-agent',
inputProcessors: ({ requestContext }) => {
const blockedWords = requestContext.get('blockedWords') ?? []
return [new ContentFilter(blockedWords)]
},
})
每次呼叫的覆寫設定「每次呼叫的覆寫設定」的直接連結
agent.generate() 與 agent.stream() 接受 inputProcessors、outputProcessors、errorProcessors 及 maxProcessorRetries。若在呼叫上設定任何 processor 陣列,該次請求便會以它替換 Agent 上設定的對應陣列。Mastra 自動新增的記憶體、Workspace、Skill、channel 及 browser processor 一律會保留,並在你的陣列前後執行。
await agent.stream('Summarize this', {
outputProcessors: [new StreamFilter()],
maxProcessorRetries: 5,
})
呼叫中傳入的 maxProcessorRetries 會覆寫 Agent 的預設值。若兩處都未設定,processor 要求的重試會視為中止。
相關內容「相關內容」的直接連結
- Processor 概觀:Processor 的概念指南
- 防護機制:安全性與驗證 processor
- 記憶體 processor:記憶體專屬 processor