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 或訊息,並要求重試 |
processOutputStream | LLM 回應期間的每個串流 chunk | 篩選或修改串流內容,並即時偵測模式 |
processLLMResponse | LLM 步驟完成並收集串流 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:
name?:
description?:
processorIndex?:
processDataParts?:
data-* chunk。預設為 false。onViolation?:
訊息引數訊息引數 的直接連結
大部分 processor 方法都會同時接收 messages 及 messageList。兩者指向相同的底層對話,但呈現方式不同。
messages 與 messageListmessages-vs-messagelist 的直接連結
messages: 以目前階段為範圍的普通MastraDBMessage物件陣列。processInput及processInputStep不包括系統訊息。processOutputResult及processOutputStep包括最新的 LLM 回應。 此陣列由messageList支援,因此就地編輯訊息的content.parts,下游 processor 及持久保存程序都會看見變更。messageList: 支援該次執行的即時MessageListinstance。它提供經篩選的檢視(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(或傳回同一個MessageListinstance)時,已記錄的修改會就地套用,因此儲存的對話會反映你的變更。 - 傳回
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。
方法方法 的直接連結
processInputprocessinput 的直接連結
在輸入訊息傳送至 LLM 前處理訊息。Agent 開始執行時只會運行一次。
processInput?(args: ProcessInputArgs): Promise<ProcessInputResult> | ProcessInputResult;
ProcessInputArgsprocessinputargs 的直接連結
messages:
systemMessages:
messageList:
abort:
retry: true 可要求 LLM 使用回饋重試該步驟。retryCount:
tracingContext?:
requestContext?:
ProcessInputResultprocessinputresult 的直接連結
此方法可傳回以下三種類型之一:
MastraDBMessage[]:
MessageList:
{ messages, systemMessages }:
processInputStepprocessinputstep 的直接連結
在 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(開始時執行一次)processInputStep(來自 inputProcessors) (每個步驟、LLM 呼叫之前)prepareStepcallback (作為 processInputStep pipeline 一部分執行,位於 inputProcessors 之後)processLLMRequest(來自 inputProcessors) (prompt 轉換後、Provider 呼叫前)- LLM 執行
processOutputStream(來自 outputProcessors) (每個串流 chunk)processLLMResponse(來自 inputProcessors) (串流完成後,與processLLMRequest配對)processOutputStep(來自 outputProcessors) (LLM 回應後、Tool 執行前)- Tool 執行(如需要)
- 如有呼叫 Tool,從步驟 2 重複
ProcessInputStepArgsprocessinputstepargs 的直接連結
messages:
messageList:
stepNumber:
steps:
systemMessages:
model:
toolChoice?:
activeTools?:
tools?:
providerOptions?:
modelSettings?:
structuredOutput?:
abort:
retry: true 可要求 LLM 使用回饋重試該步驟。retryCount:
ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。tracingContext?:
requestContext?:
ProcessInputStepResultprocessinputstepresult 的直接連結
processInputStep 可傳回多種結構:
ProcessInputStepResult物件: 覆寫此步驟下列屬性的任何組合(詳見下文)。MessageList: 傳回同一個messageListinstance,表示你已就地修改訊息。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 作出的修改只影響目前步驟,不影響後續步驟。
使用情境使用情境 的直接連結
- 按步驟編號或 context 動態切換 model
- 在指定步數後停用 Tool
- 按對話 context 動態新增或取代 Tool
- 在 Provider 之間轉換訊息 part 類型(例如為 Anthropic 將
reasoning轉為thinking) - 按步驟編號或累積的 context 修改訊息
- 新增步驟專用的系統指示
- 按步驟調整 Provider 選項(例如快取控制)
- 按步驟 context 修改結構化輸出 schema
processLLMRequestprocessllmrequest 的直接連結
Mastra 將 MessageList 轉換為 LanguageModelV2Prompt 後、呼叫 Provider 前,此方法會處理最終 LLM 請求。可用於只應影響目前送出請求、可識別 model 的暫時改寫。
傳回的 prompt 變更只會轉送至目前呼叫的 model,不會持久保存回 MessageList、memory、UI 記錄或後續 Provider 呼叫。
processLLMRequest?(
args: ProcessLLMRequestArgs,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult;
ProcessLLMRequestArgsprocessllmrequestargs 的直接連結
prompt:
model:
stepNumber:
steps:
state:
abort:
tripwire chunk。retryCount:
ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。requestContext?:
tracingContext?:
writer?:
writer.custom() 可發出 data-* chunk。abortSignal?:
傳回值傳回值 的直接連結
processLLMRequest 傳回 ProcessLLMRequestResult,即 { prompt?: LanguageModelV2Prompt } | undefined | void。
- 傳回
{ prompt }以取代目前 Provider 呼叫送出的 prompt。 - 傳回
undefined或void,以原樣轉送原始 prompt。
使用情境使用情境 的直接連結
- 在呼叫 model 前移除或重新建構 Provider 專用的 prompt part
- 標準化角色或內容,以符合 Provider 的輸入要求
- 在 loop 中途切換 Provider 時調整 Tool 結果格式
processLLMResponseprocessllmresponse 的直接連結
此方法會在步驟完成(或重播快取回應)且 output processor 收集回應 chunk 後處理 LLM 回應。此 hook 與 processLLMRequest 配對:在 Provider 呼叫前用 processLLMRequest 暫存 state(例如 cache key),再用 processLLMResponse 處理已完成的回應(例如寫入快取)。
state 物件與同一步驟傳入 processLLMRequest 的 instance 相同,因此 processor 可關聯呼叫前後的工作。
processLLMResponse?(
args: ProcessLLMResponseArgs,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult;
ProcessLLMResponseArgsprocessllmresponseargs 的直接連結
chunks:
{ type, payload })。model:
stepNumber:
steps:
state:
processLLMRequest 共用、每個 processor 獨立的 state。可用它在兩個 hook 之間傳遞資料(例如 cache key)。fromCache:
true 時,表示回應是透過 processLLMRequest 傳回 { response } 從快取重播。寫入快取的 processor 應在此值為 true 時略過寫入。warnings?:
request?:
rawResponse?:
abort:
retryCount:
0;用於限制 processor 觸發的重試。requestContext?:
tracingContext?:
writer?:
abortSignal?:
傳回值傳回值 的直接連結
processLLMResponse 傳回 ProcessLLMResponseResult,即 undefined | void。傳回值保留供日後擴充。
使用情境使用情境 的直接連結
- 即時呼叫後將 LLM 回應寫入快取(與
processLLMRequest中衍生 cache key 的操作配對) - 記錄完整回應以供分析
- 按已完成的回應觸發副作用
processAPIErrorprocessapierror 的直接連結
在 LLM API 拒絕錯誤成為最終錯誤前處理它。當 API 呼叫因不可重試的錯誤(例如 400 或 422 狀態碼)失敗時執行。processOutputStep 在成功回應後執行,而此方法會在 API 拒絕請求時執行。
請將實作 processAPIError 的 processor 加入 Agent 的 errorProcessors 陣列。
Processor 可檢查錯誤並修改請求,例如在 messageList 加入訊息。傳回 { retry: true } 可使用修改後的 state 重試。
processAPIError?(args: ProcessAPIErrorArgs): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void;
ProcessAPIErrorArgsprocessapierrorargs 的直接連結
error:
messages:
messageList:
stepNumber:
steps:
state:
retryCount:
abort:
writer?:
writer.custom() 可發出 data-* chunk。requestContext?:
abortSignal?:
ProcessAPIErrorResultprocessapierrorresult 的直接連結
retry:
使用情境使用情境 的直接連結
- 修改請求並重試,以處理 API 專用的拒絕
- 透過修改請求,將不可重試的錯誤轉為可重試
- 實作 model 專用的錯誤復原策略
範例: 自訂錯誤復原範例: 自訂錯誤復原 的直接連結
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'
export class ErrorRecoveryProcessor implements Processor {
id = 'error-recovery'
processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
// Only retry once
if (retryCount > 0) return
// Check for a specific API error
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
// Trim older messages to fit within context
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}
processOutputStreamprocessoutputstream 的直接連結
使用內置 state 管理來處理串流輸出 chunk,讓 processor 可累積 chunk,並按更完整的 context 作出決定。
processOutputStream?(args: ProcessOutputStreamArgs): Promise<ChunkType | null | undefined>;
ProcessOutputStreamArgsprocessoutputstreamargs 的直接連結
part:
streamParts:
state:
abort:
tripwire chunk。傳入 retry: true 可要求 LLM 再嘗試,而非結束。retryCount:
ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。messageList?:
tracingContext?:
requestContext?:
writer?:
傳回值傳回值 的直接連結
processOutputStream 傳回 Promise<ChunkType | null | undefined>。
- 傳回
ChunkType以發出 chunk。傳回原始part可原樣發出;傳回新的ChunkType則可發出修改後的 chunk。 - 傳回
null以捨棄 chunk。不會向下一個 processor 或 client 傳送任何內容。 - 傳回
undefined(包括return;陳述式或方法執行至結尾時隱含的undefined)以捨棄 chunk。null與undefined的行為相同。
捨棄 chunk 只影響該個 chunk。串流會繼續,下一個 chunk 仍會接受處理。要完全停止串流,請呼叫 abort()。
processOutputResultprocessoutputresult 的直接連結
在串流或生成完成後處理完整的輸出結果。
processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;
ProcessOutputResultArgsprocessoutputresultargs 的直接連結
messages:
messageList:
state:
result:
text(累積文字)、usage(包含 inputTokens、outputTokens、totalTokens 的 token 用量)、finishReason(生成結束原因)及 steps(所有 LLM 步驟結果,每項包含 toolCalls、toolResults、reasoning、sources、files 等)。abort:
tripwire chunk。retryCount:
ProcessorContext 的目前重試次數。初始值為 0;用於限制 processor 觸發的重試。tracingContext?:
requestContext?:
writer?:
processOutputStepprocessoutputstep 的直接連結
在 agentic loop 中每次 LLM 回應後、Tool 執行前處理輸出。processOutputResult 只在最後執行一次,此方法則會在每個步驟執行,最適合實作可觸發重試的 guardrail。
processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;
ProcessOutputStepArgsprocessoutputstepargs 的直接連結
messages:
messageList:
stepNumber:
finishReason?:
providerMetadata?:
steps 為空的 content-filter 封鎖情況。toolCalls?:
text?:
usage:
inputTokens、outputTokens、totalTokens)。systemMessages:
steps:
state:
abort:
retry: true 可要求 LLM 重試該步驟。retryCount:
tracingContext?:
requestContext?:
使用情境使用情境 的直接連結
- 實作可要求重試的質素 guardrail
- 在 Tool 執行前驗證 LLM 輸出
- 新增逐步記錄或指標
- 實作可重試的輸出內容審核
範例: 可重試的質素 guardrail範例: 可重試的質素 guardrail 的直接連結
import type { Processor } from '@mastra/core/processors'
export class QualityGuardrail implements Processor {
id = 'quality-guardrail'
async processOutputStep({ text, abort, retryCount }) {
const score = await evaluateResponseQuality(text)
if (score < 0.7) {
if (retryCount < 3) {
// Request retry with feedback for the LLM
abort('Response quality too low. Please provide more detail.', {
retry: true,
metadata: { qualityScore: score },
})
} else {
// Max retries reached, block the response
abort('Response quality too low after multiple attempts.')
}
}
return []
}
}
processToolResultprocesstoolresult 的直接連結
在 tool.execute() 傳回後、結果加入訊息清單或送入下一次 LLM 呼叫前處理 Tool 結果。此方法與 Tool 執行前觸發的 processOutputStep 對稱。可用此方法掃描 Tool 輸出中的 prompt injection、遮蓋敏感欄位,或用 abort('reason', { retry: true }) 中止執行。
要取代 Tool 結果,請透過 messageList.updateToolInvocation 就地修改 messageList。runtime 會從訊息清單重新讀取經 processor 處理的結果,並在下游 Tool 結果串流 chunk 排入佇列前覆寫它,讓串流 client 看見處理後的值。
tool.execute() 擲出錯誤時不會觸發此方法;只有 Tool 成功執行並有結果可用時才會呼叫。
processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;
ProcessToolResultArgsprocesstoolresultargs 的直接連結
messages:
messageList:
updateToolInvocation,以遮蓋或轉換後的值取代 Tool 結果。stepNumber:
toolName:
toolCallId:
args:
result:
tool.execute() 經過 ensureSerializable 後的輸出。對於 Provider 執行的 Tool(例如 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 的逐步 processorper-step-processor-with-processinputstep 的直接連結
import type {
Processor,
ProcessInputStepArgs,
ProcessInputStepResult,
} from '@mastra/core/processors'
export class DynamicModelProcessor implements Processor {
id = 'dynamic-model'
async processInputStep({
stepNumber,
steps,
toolChoice,
}: ProcessInputStepArgs): Promise<ProcessInputStepResult> {
// Use a fast model for initial response
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
// Switch to powerful model after tool calls
if (steps.length > 0 && steps[steps.length - 1].toolCalls?.length) {
return { model: 'openai/gpt-5.6-sol' }
}
// Disable tools after 5 steps to force completion
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
return {}
}
}
使用 processInputStep 的訊息轉換 processormessage-transformer-with-processinputstep 的直接連結
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class ReasoningTransformer implements Processor {
id = 'reasoning-transformer'
async processInputStep({ messages, messageList }) {
// Transform reasoning parts to thinking parts at each step
// This is useful when switching between model providers
for (const msg of messages) {
if (msg.role === 'assistant' && msg.content.parts) {
for (const part of msg.content.parts) {
if (part.type === 'reasoning') {
;(part as any).type = 'thinking'
}
}
}
}
return messageList
}
}
混合 processor(輸入及輸出)混合 processor(輸入及輸出) 的直接連結
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
import type { ChunkType } from '@mastra/core/stream'
export class ContentFilter implements Processor {
id = 'content-filter'
private blockedWords: string[]
constructor(blockedWords: string[]) {
this.blockedWords = blockedWords
}
async processInput({ messages, abort }): Promise<MastraDBMessage[]> {
for (const msg of messages) {
const text = msg.content.parts
?.filter(p => p.type === 'text')
.map(p => p.text)
.join(' ')
if (this.blockedWords.some(word => text?.includes(word))) {
abort('Blocked content detected in input')
}
}
return messages
}
async processOutputStream({ part, abort }): Promise<ChunkType | null> {
if (part.type === 'text-delta') {
if (this.blockedWords.some(word => part.payload.text.includes(word))) {
abort('Blocked content detected in output')
}
}
return part
}
}
使用 state 的串流累加器使用 state 的串流累加器 的直接連結
import type { Processor } from '@mastra/core/processors'
import type { ChunkType } from '@mastra/core/stream'
export class WordCounter implements Processor {
id = 'word-counter'
async processOutputStream({ part, state }): Promise<ChunkType> {
// Initialize state on first chunk
if (!state.wordCount) {
state.wordCount = 0
}
// Count words in text chunks
if (part.type === 'text-delta') {
const words = part.payload.text.split(/\s+/).filter(Boolean)
state.wordCount += words.length
}
// Log word count on finish
if (part.type === 'finish') {
console.log(`Total words: ${state.wordCount}`)
}
return part
}
}
State 生命週期State 生命週期 的直接連結
每個 processor 都會在 processLLMRequest、processLLMResponse、processOutputStream、processOutputStep、processOutputResult 及 processAPIError 中接收 state 物件。State 有三項重要特性:
- 每個 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: 附加至tripwirechunk、供下游使用者使用的選填結構化資料。
發出的 tripwire chunk 結構如下:
type TripwireChunk = {
type: 'tripwire'
runId: string
from: 'AGENT'
payload: {
reason: string
retry?: boolean
metadata?: unknown
processorId: string
}
}
在非串流呼叫(agent.generate())中,結果會透過 result.tripwire 及 result.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 後,從 processOutputStream 或 processOutputResult 發出的自訂 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 傳回 null 或 undefined 仍會捨棄 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() 接受 inputProcessors、outputProcessors、errorProcessors 及 maxProcessorRetries。呼叫時設定任何 processor 陣列,都會為該請求取代 Agent 上相應的陣列。Mastra 自動加入的 memory、Workspace、Skill、channel 及 browser processor 一律會保留,並在你的陣列前後執行。
await agent.stream('Summarize this', {
outputProcessors: [new StreamFilter()],
maxProcessorRetries: 5,
})
呼叫時傳入的 maxProcessorRetries 會覆寫 Agent 預設值。如兩者均未設定,processor 要求的重試會視為中止。
相關內容相關內容 的直接連結
- Processor 概覽: Processor 概念指南
- Guardrails: 安全及驗證 processor
- Memory Processor: Memory 專用 processor