Processor 接口
Processor 接口定义了 Mastra 中所有 Processor 都应遵循的约定。Processor 可以实现一个或多个方法,用于处理 Agent 执行管道的不同阶段。
Processor 方法的运行时机Processor 方法的运行时机的直接链接
Processor 方法会在 Agent 执行生命周期的不同节点运行:
┌────────────────────────────────────────────────────────────────────┐
│ Agent Execution Flow │
├────────────────────────────────────────────────────────────────────┤
│ │
│ User Input │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ processInput │ ← Runs ONCE at start │
│ └───────────┬────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────────────┐ │
│ │ Agentic Loop │ │
│ │ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processInputStep │ ← Runs at EACH step │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processLLMRequest │ ← Before provider call │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ LLM Execution ──── API Error? ───┐ │ │
│ │ │ │ │ │
│ │ │ ┌───────────┴──────────┐ │ │
│ │ │ │ processAPIError │ │ │
│ │ │ └──────────────────────┘ │ │
│ │ │ (retry loops back to LLM) │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processOutputStream │ ← Runs on EACH stream chunk │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processLLMResponse │ ← After stream completes │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processOutputStep │ ← Runs after EACH LLM step │ │
│ │ └───────────┬────────────┘ │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ Tool Execution (if needed) │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ ┌────────────────────────┐ │ │
│ │ │ processToolResult │ ← Runs per tool, after each │ │
│ │ └───────────┬────────────┘ tool.execute() returns │ │
│ │ │ │ │
│ │ └──────── Loop back if tools called ────────────│ │
│ │ │ │
│ └──────────────────────────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌────────────────────────┐ │
│ │ processOutputResult │ ← Runs ONCE after completion │
│ └────────────────────────┘ │
│ │ │
│ ▼ │
│ Final Response │
│ │
└────────────────────────────────────────────────────────────────────┘
| 方法 | 运行时机 | 用途 |
|---|---|---|
processInput | 开始时运行一次,在 agentic loop 之前 | 验证或转换初始用户输入、添加上下文 |
processInputStep | agentic loop 的每个步骤中,在每次 LLM 调用之前 | 在步骤之间转换消息、处理 Tool 结果 |
processLLMRequest | LLM 请求转换后、Provider 调用前 | 为当前调用重写传出的 LanguageModelV2Prompt,且不持久化更改 |
processAPIError | LLM API 调用失败时 | 检查 API 拒绝,可选择修改状态或消息,并请求重试 |
processOutputStream | LLM 响应期间的每个流式数据块 | 过滤或修改流式内容、实时检测模式 |
processLLMResponse | LLM 步骤完成且流式数据块收集完毕后 | 捕获或缓存完整响应,并运行与 processLLMRequest 配对的调用后副作用 |
processOutputStep | 每次 LLM 响应后、Tool 执行前 | 验证输出质量,实现支持重试的 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-* 数据块。默认为 false。onViolation?:
消息参数消息参数的直接链接
大多数 Processor 方法都会同时接收 messages 和 messageList。它们指向同一个底层对话,但提供不同的访问方式。
messages 与 messageListmessages-vs-messagelist的直接链接
messages:作用域限定为当前阶段的普通MastraDBMessage对象数组。对于processInput和processInputStep,其中不包含系统消息;对于processOutputResult和processOutputStep,其中包含最新的 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。
方法方法的直接链接
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(开始时运行一次)- inputProcessors 中的
processInputStep(每个步骤中,在 LLM 调用前) prepareStep回调(作为 processInputStep 管道的一部分,在 inputProcessors 之后运行)- inputProcessors 中的
processLLMRequest(prompt 转换后、Provider 调用前) - 执行 LLM
- outputProcessors 中的
processOutputStream(处理每个流式数据块) - inputProcessors 中的
processLLMResponse(流结束后运行,与processLLMRequest配对) - outputProcessors 中的
processOutputStep(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:返回同一个messageList实例,表示已就地修改消息。MastraDBMessage[]:返回转换后的消息数组,替换此步骤的消息。void或undefined:不返回任何内容,使此步骤保持不变。
对象形式可以返回以下属性的任意组合:
model?:
toolChoice?:
activeTools?:
tools?:
messages?:
messageList?:
systemMessages?:
providerOptions?:
modelSettings?:
structuredOutput?:
Processor 链式处理Processor 链式处理的直接链接
当多个 Processor 实现 processInputStep 时,它们会按顺序运行,并将更改依次传递:
Processor 1: receives { model: 'gpt-5.4' } → returns { model: 'gpt-5.4-mini' }
Processor 2: receives { model: 'gpt-5.4-mini' } → returns { toolChoice: 'none' }
Final: model = 'gpt-5.4-mini', toolChoice = 'none'
系统消息隔离系统消息隔离的直接链接
每个步骤开始时,系统消息都会重置为原始值。在 processInputStep 中所做的修改仅影响当前步骤,不影响后续步骤。
使用场景使用场景的直接链接
- 根据步骤编号或上下文动态切换模型
- 在达到一定步骤数后禁用 Tool
- 根据对话上下文动态添加或替换 Tool
- 在 Provider 之间转换消息部分类型(例如针对 Anthropic 将
reasoning转换为thinking) - 根据步骤编号或累积的上下文修改消息
- 添加步骤专用的系统指令
- 按步骤调整 Provider 选项(例如缓存控制)
- 根据步骤上下文修改结构化输出 schema
processLLMRequestprocessllmrequest的直接链接
在 Mastra 将 MessageList 转换为 LanguageModelV2Prompt 后、Provider 调用前处理最终 LLM 请求。此方法适用于只应影响当前传出请求、且能感知模型的临时重写。
返回的 prompt 更改只会转发给当前调用的模型,不会持久化回 MessageList、memory、UI 历史记录或后续 Provider 调用。
processLLMRequest?(
args: ProcessLLMRequestArgs,
): Promise<ProcessLLMRequestResult> | ProcessLLMRequestResult;
ProcessLLMRequestArgsprocessllmrequestargs的直接链接
prompt:
model:
stepNumber:
steps:
state:
abort:
tripwire 数据块。retryCount:
ProcessorContext 的当前重试次数。初始值为 0;可用它限制 Processor 触发的重试。requestContext?:
tracingContext?:
writer?:
writer.custom() 可发出 data-* 数据块。abortSignal?:
返回值返回值的直接链接
processLLMRequest 返回 ProcessLLMRequestResult,即 { prompt?: LanguageModelV2Prompt } | undefined | void。
- 返回
{ prompt },替换当前 Provider 调用传出的 prompt。 - 返回
undefined或void,不作更改地转发原始 prompt。
使用场景使用场景的直接链接
- 在模型调用前删除或重塑 Provider 专用的 prompt 部分
- 规范化角色或内容,以符合 Provider 的输入要求
- 在循环中途切换 Provider 时调整 Tool 结果格式
processLLMResponseprocessllmresponse的直接链接
在步骤完成(或重放缓存响应)且 output Processor 收集完响应数据块后处理 LLM 响应。此 hook 与 processLLMRequest 配对:在 Provider 调用前使用 processLLMRequest 暂存状态(例如缓存键),再使用 processLLMResponse 对完整响应执行操作(例如写入缓存)。
state 对象与同一步骤传给 processLLMRequest 的实例相同,因此 Processor 可以关联调用前后的工作。
processLLMResponse?(
args: ProcessLLMResponseArgs,
): Promise<ProcessLLMResponseResult> | ProcessLLMResponseResult;
ProcessLLMResponseArgsprocessllmresponseargs的直接链接
chunks:
{ type, payload })。model:
stepNumber:
steps:
state:
processLLMRequest 共享。可用它在两个 hook 之间传递数据(例如缓存键)。fromCache:
true 时,表示该响应是通过 processLLMRequest 返回 { response } 从缓存重放的。写入缓存的 Processor 应在此值为 true 时跳过写入。warnings?:
request?:
rawResponse?:
abort:
retryCount:
0;可用它限制 Processor 触发的重试。requestContext?:
tracingContext?:
writer?:
abortSignal?:
返回值返回值的直接链接
processLLMResponse 返回 ProcessLLMResponseResult,即 undefined | void。该返回值为未来扩展而保留。
使用场景使用场景的直接链接
- 在实时调用后将 LLM 响应写入缓存(与
processLLMRequest中的缓存键派生配对) - 为分析记录完整响应
- 根据完整响应触发副作用
processAPIErrorprocessapierror的直接链接
在 LLM API 拒绝错误成为最终错误前对其进行处理。当 API 调用因不可重试的错误(例如 400 或 422 状态码)失败时,此方法会运行。processOutputStep 在成功响应后运行,而此方法在 API 拒绝请求时运行。
将实现 processAPIError 的 Processor 添加到 Agent 的 errorProcessors 数组。
Processor 可以检查错误并修改请求,例如向 messageList 追加消息。返回 { retry: true } 可使用修改后的状态重试。
processAPIError?(args: ProcessAPIErrorArgs): Promise<ProcessAPIErrorResult | void> | ProcessAPIErrorResult | void;
ProcessAPIErrorArgsprocessapierrorargs的直接链接
error:
messages:
messageList:
stepNumber:
steps:
state:
retryCount:
abort:
writer?:
writer.custom() 可发出 data-* 数据块。requestContext?:
abortSignal?:
ProcessAPIErrorResultprocessapierrorresult的直接链接
retry:
使用场景使用场景的直接链接
- 通过修改请求并重试来处理 API 特有的拒绝
- 通过修改请求,将不可重试的错误转换为可重试错误
- 实现模型专用的错误恢复策略
示例:自定义错误恢复示例:自定义错误恢复的直接链接
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'
export class ErrorRecoveryProcessor implements Processor {
id = 'error-recovery'
processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
// Only retry once
if (retryCount > 0) return
// Check for a specific API error
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
// Trim older messages to fit within context
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}
processOutputStreamprocessoutputstream的直接链接
使用内置状态管理处理流式输出数据块。Processor 可以累积数据块,并根据更完整的上下文作出决策。
processOutputStream?(args: ProcessOutputStreamArgs): Promise<ChunkType | null | undefined>;
ProcessOutputStreamArgsprocessoutputstreamargs的直接链接
part:
streamParts:
state:
abort:
tripwire 数据块。传入 retry: true 可请求 LLM 再次尝试,而不是结束。retryCount:
ProcessorContext 的当前重试次数。初始值为 0;可用它限制 Processor 触发的重试。messageList?:
tracingContext?:
requestContext?:
writer?:
返回值返回值的直接链接
processOutputStream 返回 Promise<ChunkType | null | undefined>。
- 返回
ChunkType以发出数据块。返回原始part会原样发出;返回新的ChunkType会发出修改后的数据块。 - 返回
null以丢弃数据块。不会向下一个 Processor 或客户端发送任何内容。 - 返回
undefined(包括return;语句或方法运行到末尾时隐式返回的undefined)以丢弃数据块。null与undefined的行为相同。
丢弃数据块只影响该数据块。流会继续,下一个数据块仍会被处理。要完全停止流,请调用 abort()。
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 数据块。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 为空的内容筛选阻止情况。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。运行时会从消息列表中重新读取 Processor 处理后的结果,并在下游 Tool 结果流式数据块入队前将其覆盖,因此流式客户端看到的是处理后的值。
tool.execute() 抛出错误时不会触发此方法;只有 Tool 执行成功并且结果可用时才会调用。
processToolResult?(args: ProcessToolResultArgs): ProcessorMessageResult;
ProcessToolResultArgsprocesstoolresultargs的直接链接
messages:
messageList:
updateToolInvocation 可将 Tool 结果替换为经过遮盖或转换的值。stepNumber:
toolName:
toolCallId:
args:
result:
ensureSerializable 处理后的 tool.execute() 输出。对于 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()],
})
使用示例使用示例的直接链接
基础 input Processor基础 input 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
}
}
使用状态的流累加 Processor使用状态的流累加 Processor的直接链接
import type { Processor } from '@mastra/core/processors'
import type { ChunkType } from '@mastra/core/stream'
export class WordCounter implements Processor {
id = 'word-counter'
async processOutputStream({ part, state }): Promise<ChunkType> {
// Initialize state on first chunk
if (!state.wordCount) {
state.wordCount = 0
}
// Count words in text chunks
if (part.type === 'text-delta') {
const words = part.payload.text.split(/\s+/).filter(Boolean)
state.wordCount += words.length
}
// Log word count on finish
if (part.type === 'finish') {
console.log(`Total words: ${state.wordCount}`)
}
return part
}
}
状态生命周期状态生命周期的直接链接
每个 Processor 都会在 processLLMRequest、processLLMResponse、processOutputStream、processOutputStep、processOutputResult 和 processAPIError 中接收一个 state 对象。状态有三个重要特性:
- 每个 Processor 独立:每个 Processor 都有自己的
state对象,以该 Processor 的id为键。不同 id 的 Processor 无法读取或覆盖彼此的状态。 - 每个请求独立:每次调用
agent.generate()或agent.stream()时,都会在开始时创建新的状态对象。状态不会在请求之间或用户之间泄漏。 - 跨方法共享:在一次请求内,同一个
state对象会传给processLLMRequest(Provider 调用前)、processLLMResponse(步骤完成后)、processOutputStream(处理每个数据块)、processOutputStep(每个 LLM 步骤后)、processOutputResult(结束时一次)和processAPIError(LLM 调用失败时)。例如,processLLMRequest可以暂存缓存键,processLLMResponse随后可以读取该键并写入响应。
由于 state 初始为空对象,因此首次访问字段时请进行防御性初始化:
import type { Processor } from '@mastra/core/processors'
export class WordCounter implements Processor {
id = 'word-counter'
async processOutputStream({ part, state }) {
state.wordCount ??= 0
if (part.type === 'text-delta') {
state.wordCount += part.payload.text.split(/\s+/).filter(Boolean).length
}
return part
}
}
中止与 tripwire 数据块中止与 tripwire 数据块的直接链接
每个方法上的 abort 函数都会抛出 TripWire 错误,以停止处理并在输出流中发出 tripwire 数据块。客户端可以检测该数据块,从而区分被阻止的响应和正常结束。
abort('Blocked content detected', { retry: false, metadata: { category: 'pii' } })
reason:易于理解的说明,显示为tripwire.payload.reason。retry:为true时,Agent 会以reason作为反馈重试同一步骤。只有在 Agent 或调用中设置了maxProcessorRetries时才会执行重试,否则请求将中止。配置errorProcessors后,该调用的maxProcessorRetries默认为10。metadata:附加到tripwire数据块、供下游使用者处理的可选结构化数据。
发出的 tripwire 数据块结构如下:
type TripwireChunk = {
type: 'tripwire'
runId: string
from: 'AGENT'
payload: {
reason: string
retry?: boolean
metadata?: unknown
processorId: string
}
}
在非流式调用(agent.generate())中,结果通过 result.tripwire 和 result.finishReason === 'other' 公开相同的信息。
发出自定义数据块发出自定义数据块的直接链接
能够访问 writer 的 Processor 可以调用 writer.custom(chunk),将自定义 data-* 数据块流式传输给客户端。Tool 也可以通过自己的 writer 执行相同操作。这是 Processor 发出普通文本和 Tool 数据块以外内容的唯一方式。
await writer.custom({
type: 'data-moderation',
runId,
from: 'AGENT',
data: { level: 'warn', reason: 'Possibly unsafe' },
})
配置 memory 后,从 processOutputStream 或 processOutputResult 发出的自定义 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 返回 null 或 undefined 仍会丢弃数据块,因此 Processor 可以像筛选文本数据块一样检查、修改或筛选自定义数据。
在 Agent 上配置 Processor在 Agent 上配置 Processor的直接链接
通过三个数组将 Processor 附加到 Agent:
import { Agent } from '@mastra/core/agent'
import { PrefillErrorHandler } from '@mastra/core/processors'
const agent = new Agent({
id: 'support-agent',
name: 'support-agent',
model: 'openai/gpt-5',
instructions: '...',
inputProcessors: [new ContentFilter(['secret'])],
outputProcessors: [new WordCounter()],
errorProcessors: [new PrefillErrorHandler()],
maxProcessorRetries: 3,
})
inputProcessors:在 LLM 之前运行,接收输入消息。outputProcessors:在 LLM 响应期间或之后运行,接收输出数据块或消息。errorProcessors:在 LLM API 调用抛出错误时运行,接收原始错误。
每个数组也接受函数,以便针对每个请求从 RequestContext 构建 Processor:
new Agent({
id: 'processor-interface-agent',
inputProcessors: ({ requestContext }) => {
const blockedWords = requestContext.get('blockedWords') ?? []
return [new ContentFilter(blockedWords)]
},
})
单次调用覆盖项单次调用覆盖项的直接链接
agent.generate() 和 agent.stream() 接受 inputProcessors、outputProcessors、errorProcessors 和 maxProcessorRetries。在调用中设置任何 Processor 数组后,对于该请求,它会替换 Agent 上配置的对应数组。Mastra 自动添加的 memory、Workspace、Skill、channel 和 browser Processor 始终会保留,并在你的数组前后运行。
await agent.stream('Summarize this', {
outputProcessors: [new StreamFilter()],
maxProcessorRetries: 5,
})
调用中传入的 maxProcessorRetries 会覆盖 Agent 默认值。如果两处都未设置,Processor 请求的重试会被视为中止。
相关内容相关内容的直接链接
- Processor 概述:Processor 概念指南
- Guardrail:安全与验证 Processor
- Memory Processor:memory 专用 Processor