> Discover all available pages from the documentation index: https://mastra.zisheng.pro/zh-HK/llms.txt # Processor 介面 `Processor` 介面定義 Mastra 所有 processor 必須遵守的合約。Processor 可實作一個或多個方法,以處理 Agent 執行 pipeline 的不同階段。 ## Processor 方法何時執行 Processor 方法會在 Agent 執行生命週期的不同時間點執行: ```text ┌────────────────────────────────────────────────────────────────────┐ │ Agent Execution Flow │ ├────────────────────────────────────────────────────────────────────┤ │ │ │ User Input │ │ │ │ │ ▼ │ │ ┌────────────────────────┐ │ │ │ processInput │ ← Runs ONCE at start │ │ └───────────┬────────────┘ │ │ │ │ │ ▼ │ │ ┌──────────────────────────────────────────────────────────────┐ │ │ │ Agentic Loop │ │ │ │ │ │ │ │ ┌────────────────────────┐ │ │ │ │ │ processInputStep │ ← Runs at EACH step │ │ │ │ └───────────┬────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌────────────────────────┐ │ │ │ │ │ processLLMRequest │ ← Before provider call │ │ │ │ └───────────┬────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ LLM Execution ──── API Error? ───┐ │ │ │ │ │ │ │ │ │ │ │ ┌───────────┴──────────┐ │ │ │ │ │ │ processAPIError │ │ │ │ │ │ └──────────────────────┘ │ │ │ │ │ (retry loops back to LLM) │ │ │ │ ▼ │ │ │ │ ┌────────────────────────┐ │ │ │ │ │ processOutputStream │ ← Runs on EACH stream chunk │ │ │ │ └───────────┬────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌────────────────────────┐ │ │ │ │ │ processLLMResponse │ ← After stream completes │ │ │ │ └───────────┬────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌────────────────────────┐ │ │ │ │ │ processOutputStep │ ← Runs after EACH LLM step │ │ │ │ └───────────┬────────────┘ │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ Tool Execution (if needed) │ │ │ │ │ │ │ │ │ ▼ │ │ │ │ ┌────────────────────────┐ │ │ │ │ │ processToolResult │ ← Runs per tool, after each │ │ │ │ └───────────┬────────────┘ tool.execute() returns │ │ │ │ │ │ │ │ │ └──────── Loop back if tools called ────────────│ │ │ │ │ │ │ └──────────────────────────────────────────────────────────────┘ │ │ │ │ │ ▼ │ │ ┌────────────────────────┐ │ │ │ processOutputResult │ ← Runs ONCE after completion │ │ └────────────────────────┘ │ │ │ │ │ ▼ │ │ Final Response │ │ │ └────────────────────────────────────────────────────────────────────┘ ``` | 方法 | 執行時間 | 使用情境 | | --------------------- | --------------------------------------------- | ---------------------------------------------- | | `processInput` | 在開始時執行一次,位於 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` | 生成完成後執行一次 | 後處理最終回應並記錄結果 | ## 介面定義 ```typescript interface Processor { readonly id: TId readonly name?: string readonly description?: string /** Index of this processor in the workflow (set at runtime when combining processors). */ processorIndex?: number /** When true, processOutputStream also receives `data-*` chunks. Default: false. */ processDataParts?: boolean /** Callback invoked when this processor detects a violation, regardless of strategy. */ onViolation?: (violation: ProcessorViolation) => void | Promise processInput?( args: ProcessInputArgs, ): Promise | ProcessInputResult processInputStep?( args: ProcessInputStepArgs, ): | Promise | ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined processLLMRequest?( args: ProcessLLMRequestArgs, ): Promise | ProcessLLMRequestResult processLLMResponse?( args: ProcessLLMResponseArgs, ): Promise | ProcessLLMResponseResult processAPIError?( args: ProcessAPIErrorArgs, ): Promise | ProcessAPIErrorResult | void processOutputStream?( args: ProcessOutputStreamArgs, ): Promise processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult } ``` ## 屬性 **id** (`string`): Processor 的唯一識別碼,用於 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`): 當 processor 偵測到違反政策時呼叫的選填 callback,無論策略為封鎖還是警告均會執行。可用於發出警示、記錄至外部系統或向使用者傳送電郵等副作用。此 callback 擲出的錯誤會被靜默擷取,以免干擾 processor 邏輯。violation 物件包含 processorId、message,以及 processor 專用的 detail 欄位。 ## 訊息引數 大部分 processor 方法都會同時接收 `messages` 及 `messageList`。兩者指向相同的底層對話,但呈現方式不同。 ### `messages` 與 `messageList` - `messages`: 以目前階段為範圍的普通 `MastraDBMessage` 物件陣列。 `processInput` 及 `processInputStep` 不包括系統訊息。 `processOutputResult` 及 `processOutputStep` 包括最新的 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`: ```typescript import type { MastraDBMessage } from '@mastra/core/memory' function getText(message: MastraDBMessage): string { let text = '' if (message.content.parts) { for (const part of message.content.parts) { if (part.type === 'text' && typeof part.text === 'string') { text += part.text } } } // Fallback for legacy messages that only have the flattened `content` string if (!text && typeof message.content.content === 'string') { text = message.content.content } return text } ``` 重點: - `message.content.parts` 是主要來源。單一訊息可包含多個 part,包括 Tool 呼叫、Tool 結果及檔案 part 等非文字 part。 讀取 `part.text` 前,請先按 `part.type === 'text'` 篩選。 - `message.content.content` 是為向後兼容而保留的扁平字串。只有在 `parts` 為空或缺少時才用作後備。 - `message.content` 本身在 `MastraDBMessage` 上絕不會是普通字串。舊有 `CoreMessage` 結構可能是字串,但 processor 一律接收 `MastraDBMessage`。 ## 方法 ### `processInput` 在輸入訊息傳送至 LLM 前處理訊息。Agent 開始執行時只會運行一次。 ```typescript processInput?(args: ProcessInputArgs): Promise | ProcessInputResult; ``` #### `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` 此方法可傳回以下三種類型之一: **MastraDBMessage\[]** (`array`): 轉換後的訊息陣列。系統訊息維持不變。 **MessageList** (`MessageList`): 傳入的同一個 messageList instance,表示你已直接修改它。 **{ messages, systemMessages }** (`object`): 同時包含轉換後訊息及修改後系統訊息的物件。 *** ### `processInputStep` 在 agentic loop 的每個步驟、輸入訊息傳送至 LLM 前處理訊息。`processInput` 只在開始時執行一次,此方法則會在每個步驟執行,包括延續 Tool 呼叫的步驟。 ```typescript processInputStep?( args: ProcessInputStepArgs, ): | Promise | ProcessInputStepResult | MessageList | MastraDBMessage[] | void | undefined; ``` #### 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` **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` `processInputStep` 可傳回多種結構: - **`ProcessInputStepResult` 物件**: 覆寫此步驟下列屬性的任何組合(詳見下文)。 - **`MessageList`**: 傳回同一個 `messageList` instance,表示你已就地修改訊息。 - **`MastraDBMessage[]`**: 傳回轉換後的訊息陣列,取代該步驟的訊息。 - **`void` 或 `undefined`**: 不傳回任何內容,讓該步驟維持不變。 物件形式可傳回以下屬性的任何組合: **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 實作 `processInputStep` 時,會依次執行並串接變更: ```text Processor 1: receives { model: 'gpt-5.4' } → returns { model: 'gpt-5.4-mini' } Processor 2: receives { model: 'gpt-5.4-mini' } → returns { toolChoice: 'none' } Final: model = 'gpt-5.4-mini', toolChoice = 'none' ``` #### 系統訊息隔離 每個步驟開始時,系統訊息都會**重設為原始值**。在 `processInputStep` 作出的修改只影響目前步驟,不影響後續步驟。 #### 使用情境 - 按步驟編號或 context 動態切換 model - 在指定步數後停用 Tool - 按對話 context 動態新增或取代 Tool - 在 Provider 之間轉換訊息 part 類型(例如為 Anthropic 將 `reasoning` 轉為 `thinking`) - 按步驟編號或累積的 context 修改訊息 - 新增步驟專用的系統指示 - 按步驟調整 Provider 選項(例如快取控制) - 按步驟 context 修改結構化輸出 schema *** ### `processLLMRequest` Mastra 將 `MessageList` 轉換為 `LanguageModelV2Prompt` 後、呼叫 Provider 前,此方法會處理最終 LLM 請求。可用於只應影響目前送出請求、可識別 model 的暫時改寫。 傳回的 prompt 變更只會轉送至目前呼叫的 model,不會持久保存回 `MessageList`、memory、UI 記錄或後續 Provider 呼叫。 ```typescript processLLMRequest?( args: ProcessLLMRequestArgs, ): Promise | ProcessLLMRequestResult; ``` #### `ProcessLLMRequestArgs` **prompt** (`LanguageModelV2Prompt`): 這次呼叫將傳送至 Provider 的 LLM 請求 prompt。 **model** (`MastraLanguageModel`): 將接收 prompt 的已解析 model。用此值限定 Provider 專用改寫的範圍。 **stepNumber** (`number`): 目前步驟編號(從 0 開始)。步驟 0 是最初的 LLM 呼叫。 **steps** (`StepResult[]`): 先前步驟的結果,包括 text、toolCalls 及 toolResults。 **state** (`Record`): 每個 processor 獨立的 state,會在此請求內所有方法呼叫之間持續保留。 **abort** (`(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never`): 中止處理的函數。它會擲出 TripWire 錯誤以停止執行,並發出 tripwire 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。 - 傳回 `undefined` 或 `void`,以原樣轉送原始 prompt。 #### 使用情境 - 在呼叫 model 前移除或重新建構 Provider 專用的 prompt part - 標準化角色或內容,以符合 Provider 的輸入要求 - 在 loop 中途切換 Provider 時調整 Tool 結果格式 *** ### `processLLMResponse` 此方法會在步驟完成(或重播快取回應)且 output processor 收集回應 chunk 後處理 LLM 回應。此 hook 與 `processLLMRequest` 配對:在 Provider 呼叫前用 `processLLMRequest` 暫存 state(例如 cache key),再用 `processLLMResponse` 處理已完成的回應(例如寫入快取)。 `state` 物件與同一步驟傳入 `processLLMRequest` 的 instance 相同,因此 processor 可關聯呼叫前後的工作。 ```typescript processLLMResponse?( args: ProcessLLMResponseArgs, ): Promise | ProcessLLMResponseResult; ``` #### `ProcessLLMResponseArgs` **chunks** (`CachedLLMStepChunk[]`): 此步驟由 LLM 呼叫產生(或從快取重播)的 chunk,採精簡格式({ type, payload })。 **model** (`MastraLanguageModel`): 產生(或原應產生)回應的 model。 **stepNumber** (`number`): 目前步驟編號(從 0 開始)。 **steps** (`StepResult[]`): 目前為止所有已完成的步驟,包括此步驟。 **state** (`Record`): 與同一步驟的 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` 在 LLM API 拒絕錯誤成為最終錯誤前處理它。當 API 呼叫因不可重試的錯誤(例如 400 或 422 狀態碼)失敗時執行。`processOutputStep` 在成功回應後執行,而此方法會在 API 拒絕請求時執行。 請將實作 `processAPIError` 的 processor 加入 Agent 的 `errorProcessors` 陣列。 Processor 可檢查錯誤並修改請求,例如在 `messageList` 加入訊息。傳回 `{ retry: true }` 可使用修改後的 state 重試。 ```typescript processAPIError?(args: ProcessAPIErrorArgs): Promise | ProcessAPIErrorResult | void; ``` #### `ProcessAPIErrorArgs` **error** (`unknown`): LLM API 呼叫期間發生的錯誤。 **messages** (`MastraDBMessage[]`): 發生錯誤時的所有訊息。 **messageList** (`MessageList`): 用於管理訊息的 MessageList instance。修改它可在重試前變更請求。 **stepNumber** (`number`): 目前步驟編號(從 0 開始)。 **steps** (`StepResult[]`): 目前為止所有已完成的步驟。 **state** (`Record`): 每個 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` **retry** (`boolean`): 套用修改後是否重試 LLM 呼叫。 #### 使用情境 - 修改請求並重試,以處理 API 專用的拒絕 - 透過修改請求,將不可重試的錯誤轉為可重試 - 實作 model 專用的錯誤復原策略 #### 範例: 自訂錯誤復原 ```typescript import { APICallError } from '@ai-sdk/provider' import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors' export class ErrorRecoveryProcessor implements Processor { id = 'error-recovery' processAPIError({ error, messageList, retryCount, }: ProcessAPIErrorArgs): ProcessAPIErrorResult | void { // Only retry once if (retryCount > 0) return // Check for a specific API error if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) { // Trim older messages to fit within context const messages = messageList.get.all.db() if (messages.length > 4) { messageList.removeByIds([messages[1]!.id, messages[2]!.id]) return { retry: true } } } } } ``` *** ### `processOutputStream` 使用內置 state 管理來處理串流輸出 chunk,讓 processor 可累積 chunk,並按更完整的 context 作出決定。 ```typescript processOutputStream?(args: ProcessOutputStreamArgs): Promise; ``` #### `ProcessOutputStreamArgs` **part** (`ChunkType`): 目前正在處理的串流 chunk。 **streamParts** (`ChunkType[]`): 目前為止在串流中看見的所有 chunk。 **state** (`Record`): 可修改且每個 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` 以發出 chunk。傳回原始 `part` 可原樣發出;傳回新的 `ChunkType` 則可發出修改後的 chunk。 - 傳回 `null` 以捨棄 chunk。不會向下一個 processor 或 client 傳送任何內容。 - 傳回 `undefined`(包括 `return;` 陳述式或方法執行至結尾時隱含的 `undefined`)以捨棄 chunk。`null` 與 `undefined` 的行為相同。 捨棄 chunk 只影響該個 chunk。串流會繼續,下一個 chunk 仍會接受處理。要完全停止串流,請呼叫 `abort()`。 *** ### `processOutputResult` 在串流或生成完成後處理完整的輸出結果。 ```typescript processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult; ``` #### `ProcessOutputResultArgs` **messages** (`MastraDBMessage[]`): 已生成的回應訊息。 **messageList** (`MessageList`): 用於管理訊息的 MessageList instance。 **state** (`Record`): 每個 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` 在 agentic loop 中每次 LLM 回應後、Tool 執行前處理輸出。`processOutputResult` 只在最後執行一次,此方法則會在每個步驟執行,最適合實作可觸發重試的 guardrail。 ```typescript processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult; ``` #### `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 用量(inputTokens、outputTokens、totalTokens)。 **systemMessages** (`CoreMessage[]`): 所有系統訊息,可供讀取及修改。 **steps** (`StepResult[]`): 目前為止所有已完成的步驟,包括目前步驟。 **state** (`Record`): 每個 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 ```typescript import type { Processor } from '@mastra/core/processors' export class QualityGuardrail implements Processor { id = 'quality-guardrail' async processOutputStep({ text, abort, retryCount }) { const score = await evaluateResponseQuality(text) if (score < 0.7) { if (retryCount < 3) { // Request retry with feedback for the LLM abort('Response quality too low. Please provide more detail.', { retry: true, metadata: { qualityScore: score }, }) } else { // Max retries reached, block the response abort('Response quality too low after multiple attempts.') } } return [] } } ``` ### `processToolResult` 在 `tool.execute()` 傳回後、結果加入訊息清單或送入下一次 LLM 呼叫前處理 Tool 結果。此方法與 Tool 執行前觸發的 `processOutputStep` 對稱。可用此方法掃描 Tool 輸出中的 prompt injection、遮蓋敏感欄位,或用 `abort('reason', { retry: true })` 中止執行。 要取代 Tool 結果,請透過 `messageList.updateToolInvocation` 就地修改 `messageList`。runtime 會從訊息清單重新讀取經 processor 處理的結果,並在下游 Tool 結果串流 chunk 排入佇列前覆寫它,讓串流 client 看見處理後的值。 `tool.execute()` 擲出錯誤時不會觸發此方法;只有 Tool 成功執行並有結果可用時才會呼叫。 ```typescript processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult; ``` #### `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`): 每個 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 傳回值,以供合規或審計。 #### 範例: 遮蓋敏感欄位 ```typescript import type { Processor } from '@mastra/core/processors' export class RedactToolResult implements Processor { id = 'redact-tool-result' async processToolResult({ toolName, toolCallId, args, result, messageList }) { if (toolName !== 'lookup-customer') return const redacted = { ...(result as Record), ssn: '[REDACTED]', email: '[REDACTED]', } messageList.updateToolInvocation({ type: 'tool-invocation', toolInvocation: { state: 'result', toolCallId, toolName, args, result: redacted, }, }) } } ``` #### 範例: 封鎖 Tool 輸出中的 prompt injection ```typescript import type { Processor } from '@mastra/core/processors' export class ScanToolResult implements Processor { id = 'scan-tool-result' async processToolResult({ result, abort }) { const text = typeof result === 'string' ? result : JSON.stringify(result) if (containsPromptInjection(text)) { abort('blocked by scan-tool-result: suspected prompt injection') } } } function containsPromptInjection(text: string): boolean { return /ignore (all )?(previous|prior) instructions/i.test(text) } ``` ## Processor 類型 Mastra 提供類型別名,確保 processor 實作所需方法: ```typescript // Must implement processInput, processInputStep, processLLMRequest, or processLLMResponse (or any combination) type InputProcessor = Processor & ( | { processInput: required } | { processInputStep: required } | { processLLMRequest: required } | { processLLMResponse: required } ) // Must implement processOutputStream, processOutputStep, OR processOutputResult (or any combination) type OutputProcessor = Processor & ( | { processOutputStream: required } | { processOutputStep: required } | { processOutputResult: required } ) // Must implement processAPIError type ErrorProcessor = Processor & { processAPIError: required } ``` 在 `errorProcessors` 設定實作 `processAPIError` 的 processor: ```typescript const agent = new Agent({ id: 'agent', errorProcessors: [new PrefillErrorHandler()], }) ``` ## 使用範例 ### 基本輸入 processor ```typescript import type { Processor } from '@mastra/core/processors' import type { MastraDBMessage } from '@mastra/core/memory' export class LowercaseProcessor implements Processor { id = 'lowercase' async processInput({ messages }): Promise { return messages.map(msg => ({ ...msg, content: { ...msg.content, parts: msg.content.parts?.map(part => part.type === 'text' ? { ...part, text: part.text.toLowerCase() } : part, ), }, })) } } ``` ### 使用 `processInputStep` 的逐步 processor ```typescript import type { Processor, ProcessInputStepArgs, ProcessInputStepResult, } from '@mastra/core/processors' export class DynamicModelProcessor implements Processor { id = 'dynamic-model' async processInputStep({ stepNumber, steps, toolChoice, }: ProcessInputStepArgs): Promise { // Use a fast model for initial response if (stepNumber === 0) { return { model: 'openai/gpt-5-mini' } } // Switch to powerful model after tool calls if (steps.length > 0 && steps[steps.length - 1].toolCalls?.length) { return { model: 'openai/gpt-5.6-sol' } } // Disable tools after 5 steps to force completion if (stepNumber > 5) { return { toolChoice: 'none' } } return {} } } ``` ### 使用 `processInputStep` 的訊息轉換 processor ```typescript import type { Processor } from '@mastra/core/processors' import type { MastraDBMessage } from '@mastra/core/memory' export class ReasoningTransformer implements Processor { id = 'reasoning-transformer' async processInputStep({ messages, messageList }) { // Transform reasoning parts to thinking parts at each step // This is useful when switching between model providers for (const msg of messages) { if (msg.role === 'assistant' && msg.content.parts) { for (const part of msg.content.parts) { if (part.type === 'reasoning') { ;(part as any).type = 'thinking' } } } } return messageList } } ``` ### 混合 processor(輸入及輸出) ```typescript import type { Processor } from '@mastra/core/processors' import type { MastraDBMessage } from '@mastra/core/memory' import type { ChunkType } from '@mastra/core/stream' export class ContentFilter implements Processor { id = 'content-filter' private blockedWords: string[] constructor(blockedWords: string[]) { this.blockedWords = blockedWords } async processInput({ messages, abort }): Promise { for (const msg of messages) { const text = msg.content.parts ?.filter(p => p.type === 'text') .map(p => p.text) .join(' ') if (this.blockedWords.some(word => text?.includes(word))) { abort('Blocked content detected in input') } } return messages } async processOutputStream({ part, abort }): Promise { if (part.type === 'text-delta') { if (this.blockedWords.some(word => part.payload.text.includes(word))) { abort('Blocked content detected in output') } } return part } } ``` ### 使用 state 的串流累加器 ```typescript import type { Processor } from '@mastra/core/processors' import type { ChunkType } from '@mastra/core/stream' export class WordCounter implements Processor { id = 'word-counter' async processOutputStream({ part, state }): Promise { // Initialize state on first chunk if (!state.wordCount) { state.wordCount = 0 } // Count words in text chunks if (part.type === 'text-delta') { const words = part.payload.text.split(/\s+/).filter(Boolean) state.wordCount += words.length } // Log word count on finish if (part.type === 'finish') { console.log(`Total words: ${state.wordCount}`) } return part } } ``` ## State 生命週期 每個 processor 都會在 `processLLMRequest`、`processLLMResponse`、`processOutputStream`、`processOutputStep`、`processOutputResult` 及 `processAPIError` 中接收 `state` 物件。State 有三項重要特性: - **每個 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` 初始為空物件,首次存取時應以防禦方式初始化欄位: ```typescript import type { Processor } from '@mastra/core/processors' export class WordCounter implements Processor { id = 'word-counter' async processOutputStream({ part, state }) { state.wordCount ??= 0 if (part.type === 'text-delta') { state.wordCount += part.payload.text.split(/\s+/).filter(Boolean).length } return part } } ``` ## 中止及 tripwire chunk 每個方法的 `abort` 函數都會擲出 `TripWire` 錯誤,以停止處理並在輸出串流發出 `tripwire` chunk。Client 可偵測該 chunk,以區分被封鎖的回應與正常完成。 ```typescript abort('Blocked content detected', { retry: false, metadata: { category: 'pii' } }) ``` - `reason`: 易於理解的說明,會顯示為 `tripwire.payload.reason`。 - `retry`: 設為 `true` 時,Agent 會重試相同步驟,並將 `reason` 作為回饋。只有在 Agent 或呼叫上設定 `maxProcessorRetries` 時才會重試,否則請求會中止。設定 `errorProcessors` 後,該次呼叫的 `maxProcessorRetries` 預設為 `10`。 - `metadata`: 附加至 `tripwire` chunk、供下游使用者使用的選填結構化資料。 發出的 `tripwire` chunk 結構如下: ```typescript type TripwireChunk = { type: 'tripwire' runId: string from: 'AGENT' payload: { reason: string retry?: boolean metadata?: unknown processorId: string } } ``` 在非串流呼叫(`agent.generate()`)中,結果會透過 `result.tripwire` 及 `result.finishReason === 'other'` 提供相同資料。 ## 發出自訂資料 chunk 可存取 `writer` 的 processor 能呼叫 `writer.custom(chunk)`,向 client 串流傳送自訂 `data-*` chunk。Tool 也可透過自己的 writer 執行相同操作。這是 processor 在一般文字及 Tool chunk 以外發出內容的唯一方式。 ```typescript 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: ```typescript 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` 即可啟用: ```typescript class ModerationCollector implements Processor { id = 'moderation-collector' processDataParts = true async processOutputStream({ part, state }) { if (part.type === 'data-moderation') { state.warnings ??= [] state.warnings.push(part.data) } return part } } ``` Chunk 的 `type` 必須以 `data-` 開頭,才會視為自訂資料 chunk。從 `processOutputStream` 傳回 `null` 或 `undefined` 仍會捨棄 chunk,因此 processor 可像篩選文字 chunk 一樣檢查、修改或篩選自訂資料。 ## 在 Agent 設定 processor Processor 透過三個陣列附加至 Agent: ```typescript import { Agent } from '@mastra/core/agent' import { PrefillErrorHandler } from '@mastra/core/processors' const agent = new Agent({ id: 'support-agent', name: 'support-agent', model: 'openai/gpt-5', instructions: '...', inputProcessors: [new ContentFilter(['secret'])], outputProcessors: [new WordCounter()], errorProcessors: [new PrefillErrorHandler()], maxProcessorRetries: 3, }) ``` - `inputProcessors`: 在 LLM 前執行,接收輸入訊息。 - `outputProcessors`: 在 LLM 回應期間或之後執行,接收輸出 chunk 或訊息。 - `errorProcessors`: 在 LLM API 呼叫擲出錯誤時執行,接收原始錯誤。 每個陣列亦接受函數,讓 processor 可按請求從 `RequestContext` 建立: ```typescript new Agent({ id: 'processor-interface-agent', inputProcessors: ({ requestContext }) => { const blockedWords = requestContext.get('blockedWords') ?? [] return [new ContentFilter(blockedWords)] }, }) ``` ### 每次呼叫的覆寫設定 `agent.generate()` 及 `agent.stream()` 接受 `inputProcessors`、`outputProcessors`、`errorProcessors` 及 `maxProcessorRetries`。呼叫時設定任何 processor 陣列,都會為該請求**取代** Agent 上相應的陣列。Mastra 自動加入的 memory、Workspace、Skill、channel 及 browser processor 一律會保留,並在你的陣列前後執行。 ```typescript await agent.stream('Summarize this', { outputProcessors: [new StreamFilter()], maxProcessorRetries: 5, }) ``` 呼叫時傳入的 `maxProcessorRetries` 會覆寫 Agent 預設值。如兩者均未設定,processor 要求的重試會視為中止。 ## 相關內容 - [Processor 概覽](https://mastra.zisheng.pro/zh-HK/docs/agents/processors): Processor 概念指南 - [Guardrails](https://mastra.zisheng.pro/zh-HK/docs/agents/guardrails): 安全及驗證 processor - [Memory Processor](https://mastra.zisheng.pro/zh-HK/docs/memory/memory-processors): Memory 專用 processor