跳到主要内容

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 之前验证或转换初始用户输入、添加上下文
processInputStepagentic loop 的每个步骤中,在每次 LLM 调用之前在步骤之间转换消息、处理 Tool 结果
processLLMRequestLLM 请求转换后、Provider 调用前为当前调用重写传出的 LanguageModelV2Prompt,且不持久化更改
processAPIErrorLLM API 调用失败时检查 API 拒绝,可选择修改状态或消息,并请求重试
processOutputStreamLLM 响应期间的每个流式数据块过滤或修改流式内容、实时检测模式
processLLMResponseLLM 步骤完成且流式数据块收集完毕后捕获或缓存完整响应,并运行与 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 列表中的位置。当 Processor 与 memory、Workspace 和单次调用覆盖项合并时,由 Mastra 在运行时设置。无需自行设置。

processDataParts?:

boolean
为 true 时,processOutputStream 方法也会接收 Tool 通过 writer.custom() 发出的 data-* 数据块。默认为 false。

onViolation?:

(violation: ProcessorViolation) => void | Promise<void>
Processor 检测到策略违规时调用的可选回调,无论采用何种策略(阻止或警告)都会调用。可用于告警、记录到外部系统或向用户发送电子邮件等副作用。此回调抛出的错误会被静默捕获,以免干扰 Processor 逻辑。violation 对象包含 processorId、message 和 Processor 专用的 detail 字段。

消息参数
消息参数的直接链接

大多数 Processor 方法都会同时接收 messagesmessageList。它们指向同一个底层对话,但提供不同的访问方式。

messagesmessageList
messages-vs-messagelist的直接链接

  • messages:作用域限定为当前阶段的普通 MastraDBMessage 对象数组。对于 processInputprocessInputStep,其中不包含系统消息;对于 processOutputResultprocessOutputStep,其中包含最新的 LLM 响应。该数组由 messageList 支持,因此就地编辑消息的 content.parts 后,下游 Processor 和持久化流程都能看到更改。
  • messageList:支撑本次运行的实时 MessageList 实例。它提供经过筛选的视图(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 实例)时,记录的修改会就地应用,因此保存的对话会反映这些更改。
  • 返回 MastraDBMessage[]{ messages, systemMessages } 时,Mastra 会针对当前阶段将返回的数组与 messageList 协调一致,删除缺失的消息并替换系统消息。

返回另一个 MessageList 实例会导致错误。请始终修改传给 Processor 的实例。

从消息中读取文本
从消息中读取文本的直接链接

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 是主要来源。一条消息可以包含多个部分,其中包括 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:

MastraDBMessage[]
要处理的用户和 assistant 消息(不包含系统消息)。

systemMessages:

CoreMessage[]
所有系统消息(Agent 指令、memory 上下文和用户提供的消息)。可以修改并返回。

messageList:

MessageList
用于高级消息管理的完整 MessageList 实例。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
用于中止处理的函数。它会抛出 TripWire 错误并停止执行。传入 retry: true 可请求 LLM 根据反馈重试此步骤。

retryCount:

number
本次生成中 Processor 触发重试的次数。可用它限制重试次数。始终由 Mastra 传入,初始值为 0。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

requestContext?:

RequestContext
请求作用域的上下文,包含 threadId 和 resourceId 等执行元数据。

ProcessInputResult
processinputresult的直接链接

该方法可以返回以下三种类型之一:

MastraDBMessage[]:

array
转换后的消息数组。系统消息保持不变。

MessageList:

MessageList
传入的同一个 messageList 实例,表示已直接修改该实例。

{ 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. inputProcessors 中的 processInputStep(每个步骤中,在 LLM 调用前)
  3. prepareStep 回调(作为 processInputStep 管道的一部分,在 inputProcessors 之后运行)
  4. inputProcessors 中的 processLLMRequest(prompt 转换后、Provider 调用前)
  5. 执行 LLM
  6. outputProcessors 中的 processOutputStream(处理每个流式数据块)
  7. inputProcessors 中的 processLLMResponse(流结束后运行,与 processLLMRequest 配对)
  8. outputProcessors 中的 processOutputStep(LLM 响应后、Tool 执行前)
  9. 执行 Tool(如有需要)
  10. 如果调用了 Tool,则从第 2 步重复

ProcessInputStepArgs
processinputstepargs的直接链接

messages:

MastraDBMessage[]
所有消息,包括前面步骤中的 Tool 调用和结果(只读快照)。

messageList:

MessageList
用于管理消息的 MessageList 实例。可以直接修改,也可以在结果中返回。

stepNumber:

number
当前步骤编号(从 0 开始)。步骤 0 是初始 LLM 调用。

steps:

StepResult[]
前面步骤的结果,包括 text、toolCalls 和 toolResults。

systemMessages:

CoreMessage[]
所有系统消息(只读快照)。在结果中返回可将其替换。

model:

MastraLanguageModelV2
当前使用的模型。在结果中返回其他模型可进行切换。

toolChoice?:

ToolChoice
当前 Tool 选择设置('auto'、'none'、'required' 或特定 Tool)。

activeTools?:

string[]
当前处于活动状态的 Tool 名称。返回筛选后的数组可限制 Tool。

tools?:

ToolSet
当前步骤可用的 Tool。在结果中返回可添加或替换 Tool。

providerOptions?:

SharedV2ProviderOptions
Provider 专用选项(例如 Anthropic cacheControl、OpenAI reasoningEffort)。

modelSettings?:

CallSettings
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 上下文。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

ProcessInputStepResult
processinputstepresult的直接链接

processInputStep 可以返回多种结构:

  • ProcessInputStepResult 对象:为此步骤覆盖下列属性的任意组合(详见下文)。
  • MessageList:返回同一个 messageList 实例,表示已就地修改消息。
  • MastraDBMessage[]:返回转换后的消息数组,替换此步骤的消息。
  • voidundefined:不返回任何内容,使此步骤保持不变。

对象形式可以返回以下属性的任意组合:

model?:

LanguageModelV2 | string
更改此步骤的模型。可以是模型实例,也可以是类似 'openai/gpt-5.5' 的路由 ID。

toolChoice?:

ToolChoice
更改此步骤的 Tool 选择行为。

activeTools?:

string[]
筛选此步骤可用的 Tool。

tools?:

ToolSet
替换或修改此步骤的 Tool。使用展开语法合并:{ tools: { ...tools, newTool } }。

messages?:

MastraDBMessage[]
替换所有消息。不能与 messageList 同时使用。

messageList?:

MessageList
返回同一个 messageList 实例(表示已修改该实例)。不能与 messages 同时使用。

systemMessages?:

CoreMessage[]
仅替换此步骤的所有系统消息。

providerOptions?:

SharedV2ProviderOptions
更改此步骤的 Provider 专用选项。

modelSettings?:

CallSettings
更改此步骤的模型设置。

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 中所做的修改仅影响当前步骤,不影响后续步骤。

使用场景
使用场景的直接链接

  • 根据步骤编号或上下文动态切换模型
  • 在达到一定步骤数后禁用 Tool
  • 根据对话上下文动态添加或替换 Tool
  • 在 Provider 之间转换消息部分类型(例如针对 Anthropic 将 reasoning 转换为 thinking
  • 根据步骤编号或累积的上下文修改消息
  • 添加步骤专用的系统指令
  • 按步骤调整 Provider 选项(例如缓存控制)
  • 根据步骤上下文修改结构化输出 schema

processLLMRequest
processllmrequest的直接链接

在 Mastra 将 MessageList 转换为 LanguageModelV2Prompt 后、Provider 调用前处理最终 LLM 请求。此方法适用于只应影响当前传出请求、且能感知模型的临时重写。

返回的 prompt 更改只会转发给当前调用的模型,不会持久化回 MessageList、memory、UI 历史记录或后续 Provider 调用。

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

ProcessLLMRequestArgs
processllmrequestargs的直接链接

prompt:

LanguageModelV2Prompt
本次调用将发送给 Provider 的 LLM 请求 prompt。

model:

MastraLanguageModel
将接收 prompt 的已解析模型。可用它限定 Provider 专用重写的作用域。

stepNumber:

number
当前步骤编号(从 0 开始)。步骤 0 是首次 LLM 调用。

steps:

StepResult[]
之前步骤的结果,包括 text、toolCalls 和 toolResults。

state:

Record<string, unknown>
每个 Processor 独立的状态,在此请求内的所有方法调用之间持续存在。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
用于中止处理的函数。它会抛出 TripWire 错误、停止执行并发出 tripwire 数据块。

retryCount:

number
来自 ProcessorContext 的当前重试次数。初始值为 0;可用它限制 Processor 触发的重试。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

writer?:

ProcessorStreamWriter
用于在流式传输期间发出自定义数据块的 stream writer。调用 writer.custom() 可发出 data-* 数据块。

abortSignal?:

AbortSignal
用于取消操作的信号。

返回值
返回值的直接链接

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

  • 返回 { prompt },替换当前 Provider 调用传出的 prompt。
  • 返回 undefinedvoid,不作更改地转发原始 prompt。

使用场景
使用场景的直接链接

  • 在模型调用前删除或重塑 Provider 专用的 prompt 部分
  • 规范化角色或内容,以符合 Provider 的输入要求
  • 在循环中途切换 Provider 时调整 Tool 结果格式

processLLMResponse
processllmresponse的直接链接

在步骤完成(或重放缓存响应)且 output Processor 收集完响应数据块后处理 LLM 响应。此 hook 与 processLLMRequest 配对:在 Provider 调用前使用 processLLMRequest 暂存状态(例如缓存键),再使用 processLLMResponse 对完整响应执行操作(例如写入缓存)。

state 对象与同一步骤传给 processLLMRequest 的实例相同,因此 Processor 可以关联调用前后的工作。

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

ProcessLLMResponseArgs
processllmresponseargs的直接链接

chunks:

CachedLLMStepChunk[]
此步骤中由 LLM 调用生成(或从缓存重放)的数据块,采用精简格式({ type, payload })。

model:

MastraLanguageModel
生成(或原本会生成)该响应的模型。

stepNumber:

number
当前步骤编号(从 0 开始)。

steps:

StepResult[]
截至目前已完成的所有步骤,包括当前步骤。

state:

Record<string, unknown>
每个 Processor 独立的状态,与同一步骤的 processLLMRequest 共享。可用它在两个 hook 之间传递数据(例如缓存键)。

fromCache:

boolean
true 时,表示该响应是通过 processLLMRequest 返回 { response } 从缓存重放的。写入缓存的 Processor 应在此值为 true 时跳过写入。

warnings?:

LanguageModelV2CallWarning[]
语言模型调用报告的警告(例如不支持的设置)。

request?:

unknown
Provider 请求体(如果可用),可用于 tracing。

rawResponse?:

unknown
原始 Provider 响应(如果可用),可用于 tracing。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
用于中止处理的函数。它会抛出 TripWire 错误并停止执行。

retryCount:

number
当前重试次数。初始值为 0;可用它限制 Processor 触发的重试。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

writer?:

ProcessorStreamWriter
用于发出自定义数据块的 stream writer。

abortSignal?:

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:

unknown
LLM API 调用期间发生的错误。

messages:

MastraDBMessage[]
发生错误时的所有消息。

messageList:

MessageList
用于管理消息的 MessageList 实例。修改该实例可在重试前更改请求。

stepNumber:

number
当前步骤编号(从 0 开始)。

steps:

StepResult[]
截至目前已完成的所有步骤。

state:

Record<string, unknown>
每个 Processor 独立的状态,在此请求内的所有方法调用之间持续存在。

retryCount:

number
错误处理程序的当前重试次数。可用它限制重试次数。

abort:

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

writer?:

ProcessorStreamWriter
用于在流式传输期间发出自定义数据块的 stream writer。调用 writer.custom() 可发出 data-* 数据块。

requestContext?:

RequestContext
从 Agent 调用传入的请求上下文。

abortSignal?:

AbortSignal
用于取消操作的信号。

ProcessAPIErrorResult
processapierrorresult的直接链接

retry:

boolean
应用修改后是否重试 LLM 调用。

使用场景
使用场景的直接链接

  • 通过修改请求并重试来处理 API 特有的拒绝
  • 通过修改请求,将不可重试的错误转换为可重试错误
  • 实现模型专用的错误恢复策略

示例:自定义错误恢复
示例:自定义错误恢复的直接链接

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的直接链接

使用内置状态管理处理流式输出数据块。Processor 可以累积数据块,并根据更完整的上下文作出决策。

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

ProcessOutputStreamArgs
processoutputstreamargs的直接链接

part:

ChunkType
当前正在处理的流式数据块。

streamParts:

ChunkType[]
截至目前在流中收到的所有数据块。

state:

Record<string, unknown>
每个 Processor 独立的可变状态,在单次请求的所有数据块和所有方法调用之间持续存在。每次新的 generate 或 stream 调用都会创建新的状态对象。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
用于中止流的函数。它会抛出 TripWire 错误、结束流并发出 tripwire 数据块。传入 retry: true 可请求 LLM 再次尝试,而不是结束。

retryCount:

number
来自 ProcessorContext 的当前重试次数。初始值为 0;可用它限制 Processor 触发的重试。

messageList?:

MessageList
用于访问对话历史记录的 MessageList 实例。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

writer?:

ProcessorStreamWriter
用于向客户端发回自定义数据块的 stream writer。调用 writer.custom() 可发出 data-* 类型的数据块。在流式传输期间可用。

返回值
返回值的直接链接

processOutputStream 返回 Promise<ChunkType | null | undefined>

  • 返回 ChunkType 以发出数据块。返回原始 part 会原样发出;返回新的 ChunkType 会发出修改后的数据块。
  • 返回 null 以丢弃数据块。不会向下一个 Processor 或客户端发送任何内容。
  • 返回 undefined(包括 return; 语句或方法运行到末尾时隐式返回的 undefined)以丢弃数据块。nullundefined 的行为相同。

丢弃数据块只影响该数据块。流会继续,下一个数据块仍会被处理。要完全停止流,请调用 abort()


processOutputResult
processoutputresult的直接链接

在流式传输或生成完成后处理完整的输出结果。

processOutputResult?(args: ProcessOutputResultArgs): ProcessorMessageResult;

ProcessOutputResultArgs
processoutputresultargs的直接链接

messages:

MastraDBMessage[]
生成的响应消息。

messageList:

MessageList
用于管理消息的 MessageList 实例。

state:

Record<string, unknown>
每个 Processor 独立的状态,在此请求内的所有方法调用之间持续存在,并与 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 数据块。

retryCount:

number
来自 ProcessorContext 的当前重试次数。初始值为 0;可用它限制 Processor 触发的重试。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

writer?:

ProcessorStreamWriter
用于向客户端发回自定义数据块的 stream writer。调用 writer.custom() 可发出 data-* 类型的数据块。在流式传输期间可用。

processOutputStep
processoutputstep的直接链接

在 agentic loop 中每次 LLM 响应后、Tool 执行前处理输出。processOutputResult 仅在结束时运行一次,而此方法会在每个步骤运行。它非常适合实现能够触发重试的 guardrail。

processOutputStep?(args: ProcessOutputStepArgs): ProcessorMessageResult;

ProcessOutputStepArgs
processoutputstepargs的直接链接

messages:

MastraDBMessage[]
包括最新 LLM 响应在内的所有消息。

messageList:

MessageList
用于管理消息的 MessageList 实例。

stepNumber:

number
当前步骤编号(从 0 开始)。

finishReason?:

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

providerMetadata?:

ProviderMetadata
结束步骤的 Provider 专用元数据(例如 AWS Bedrock guardrail trace)。模型步骤生成 Provider 元数据时会提供此值,包括 steps 为空的内容筛选阻止情况。

toolCalls?:

ToolCallInfo[]
此步骤中发起的 Tool 调用(如果有)。

text?:

string
此步骤生成的文本。

usage:

LanguageModelUsage
当前步骤的 token 使用量(inputTokensoutputTokenstotalTokens)。

systemMessages:

CoreMessage[]
所有系统消息,可供读取和修改。

steps:

StepResult[]
截至目前已完成的所有步骤,包括当前步骤。

state:

Record<string, unknown>
每个 Processor 独立的状态,在此请求内的所有方法调用之间持续存在,并与 processOutputStream 和 processOutputResult 共享。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
用于中止处理的函数。传入 retry: true 可请求 LLM 重试此步骤。

retryCount:

number
Processor 触发重试的次数。可用它限制重试次数。始终由 Mastra 传入,初始值为 0。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

使用场景
使用场景的直接链接

  • 实现能够请求重试的质量 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。运行时会从消息列表中重新读取 Processor 处理后的结果,并在下游 Tool 结果流式数据块入队前将其覆盖,因此流式客户端看到的是处理后的值。

tool.execute() 抛出错误时不会触发此方法;只有 Tool 执行成功并且结果可用时才会调用。

processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;

ProcessToolResultArgs
processtoolresultargs的直接链接

messages:

MastraDBMessage[]
包括当前含 Tool 调用的 assistant 消息在内的所有消息。

messageList:

MessageList
用于管理消息的 MessageList 实例。调用 updateToolInvocation 可将 Tool 结果替换为经过遮盖或转换的值。

stepNumber:

number
当前步骤编号(从 0 开始)。

toolName:

string
已执行 Tool 的名称。

toolCallId:

string
此次特定 Tool 调用的唯一标识符。

args:

unknown
LLM 传给 Tool 的参数。

result:

unknown
Tool 返回的值。对于客户端执行的 Tool,这是经过 ensureSerializable 处理后的 tool.execute() 输出。对于 Provider 执行的 Tool(例如 Anthropic web_search),这是来自 Provider 流的原始结果,不会经过 ensureSerializable 处理。

providerExecuted?:

boolean
此结果是否来自 Anthropic web_search 等由 Provider 执行的 Tool。对于客户端执行的 Tool,默认为 undefined。

systemMessages:

CoreMessage[]
所有系统消息,可供读取。

steps:

StepResult[]
截至目前已完成的所有步骤。

state:

Record<string, unknown>
每个 Processor 独立的状态,在此请求内的所有方法调用之间持续存在,并与同一 Processor 上的其他方法共享。

abort:

(reason?: string, options?: { retry?: boolean; metadata?: unknown }) => never
用于中止运行的函数。传入 retry: true 可请求 LLM 以中止原因为反馈重试此步骤。

retryCount:

number
Processor 触发重试的次数,初始值为 0。

tracingContext?:

TracingContext
用于可观测性的 tracing 上下文。

requestContext?:

RequestContext
包含执行元数据的请求作用域上下文。

使用场景
使用场景的直接链接

  • 在 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()],
})

使用示例
使用示例的直接链接

基础 input Processor
基础 input 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
}
}

使用状态的流累加 Processor
使用状态的流累加 Processor的直接链接

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
}
}

状态生命周期
状态生命周期的直接链接

每个 Processor 都会在 processLLMRequestprocessLLMResponseprocessOutputStreamprocessOutputStepprocessOutputResultprocessAPIError 中接收一个 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.tripwireresult.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' },
})

配置 memory 后,从 processOutputStreamprocessOutputResult 发出的自定义 data-* 数据块会保存为 assistant 消息的一部分。在数据块对象上设置 transient: true,可流式传输数据块而不将其保存到 memory:

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 返回 nullundefined 仍会丢弃数据块,因此 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() 接受 inputProcessorsoutputProcessorserrorProcessorsmaxProcessorRetries。在调用中设置任何 Processor 数组后,对于该请求,它会替换 Agent 上配置的对应数组。Mastra 自动添加的 memory、Workspace、Skill、channel 和 browser Processor 始终会保留,并在你的数组前后运行。

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

调用中传入的 maxProcessorRetries 会覆盖 Agent 默认值。如果两处都未设置,Processor 请求的重试会被视为中止。