Processor
Processor 会在消息经过 Agent 时对其进行转换、验证或控制。它们在 Agent 执行管道的特定位置运行,让你可以在输入到达语言模型之前修改输入,或在输出返回用户之前修改输出。
Processor 的配置方式如下:
inputProcessors:在消息到达语言模型之前运行。outputProcessors:在语言模型生成响应之后、响应返回用户之前运行。
你可以使用单独的 Processor 对象,也可以使用 Mastra 的 Workflow 原语将它们组合成 Workflow。Workflow 让你能够更细致地控制 Processor 执行顺序、并行处理和条件逻辑。
某些 Processor 同时实现输入和输出逻辑,可以根据转换所需发生的位置用于任一数组。
某些内置 Processor 还会发送隐藏的系统提醒信号。这些信号会持久化到原始 Memory 历史中,并在下一次模型调用前转换为 <system-reminder>...</system-reminder> 上下文;但面向 UI 的标准消息转换和默认 Memory 召回会隐藏它们,除非你明确选择显示。要为当前调用传递信号而不保留它,请使用 transient: true 发送。
何时使用 Processor何时使用 Processor的直接链接
使用 Processor 可以:
- 规范化或验证用户输入
- 为 Agent 添加 Guardrail
- 检测并阻止提示词注入或越狱尝试
- 出于安全或合规目的审核内容
- 转换消息(例如翻译语言、过滤 Tool 调用)
- 限制 token 用量或消息历史长度
- 遮盖敏感信息(PII)
- 对消息应用自定义业务逻辑
Mastra 包含多个适用于常见用例的 Processor。你还可以根据应用的特定要求创建自定义 Processor。
快速开始快速开始的直接链接
导入并实例化 Processor,然后将其传入 Agent 的 inputProcessors 或 outputProcessors 数组:
import { Agent } from '@mastra/core/agent'
import { ModerationProcessor } from '@mastra/core/processors'
export const moderatedAgent = new Agent({
id: 'moderated-agent',
name: 'moderated-agent',
instructions: 'You are a helpful assistant',
model: 'openai/gpt-5-mini',
inputProcessors: [
new ModerationProcessor({
model: 'openai/gpt-5-mini',
categories: ['hate', 'harassment', 'violence'],
threshold: 0.7,
strategy: 'block',
}),
],
})
执行顺序执行顺序的直接链接
Processor 按照它们在数组中的顺序运行:
inputProcessors: [new UnicodeNormalizer(), new PromptInjectionDetector(), new ModerationProcessor()]
对于输出 Processor,此顺序决定对模型响应应用转换的先后次序。
启用 Memory 时启用 Memory 时的直接链接
在 Agent 上启用 Memory 后,Memory Processor 会自动添加到管道中:
输入 Processor:
[Memory Processors] → [Your inputProcessors]
Memory 先加载消息历史,然后运行你的 Processor。
输出 Processor:
[Your outputProcessors] → [Memory Processors]
你的 Processor 先运行,然后由 Memory 持久化消息。
按照此顺序,调用 abort() 的输出 Guardrail 会跳过 Memory Processor,并阻止保存消息。有关详细信息,请参阅 Memory Processor。
将 Processor 附加到 Agent将 Processor 附加到 Agent的直接链接
通过三个数组在 Agent 上配置 Processor:
import { Agent } from '@mastra/core/agent'
import { PrefillErrorHandler, TokenLimiter, ModerationProcessor } from '@mastra/core/processors'
const agent = new Agent({
id: 'support-agent',
name: 'support-agent',
model: 'openai/gpt-5',
instructions: '...',
inputProcessors: [
new TokenLimiter(4000),
new ModerationProcessor({ model: 'openai/gpt-5-nano' }),
],
outputProcessors: [new ModerationProcessor({ model: 'openai/gpt-5-nano' })],
errorProcessors: [new PrefillErrorHandler()],
})
inputProcessors在 LLM 之前运行。outputProcessors在 LLM 响应期间及之后运行。errorProcessors在 LLM API 调用抛出异常时运行,以便从 Provider 错误中恢复。
每个数组也接受返回数组的函数,因此可以根据 RequestContext 按请求构建 Processor:
new Agent({
id: 'processors-agent',
inputProcessors: ({ requestContext }) => {
const limit = requestContext.get('tokenLimit') ?? 4000
return [new TokenLimiter(limit)]
},
})
按调用覆盖 Processor按调用覆盖 Processor的直接链接
agent.generate() 和 agent.stream() 接受相同的三个数组。传入其中一个时,它仅针对该次调用替换 Agent 上匹配的数组。Memory、Workspace 和其他由框架管理的 Processor 仍会在你的数组前后运行。
await agent.stream('Summarize this', {
inputProcessors: [new TokenLimiter(2000)],
maxProcessorRetries: 5,
})
创建自定义 Processor创建自定义 Processor的直接链接
自定义 Processor 实现 Processor 接口。
Processor 方法接收两个用于访问对话的参数:
messages:当前阶段的MastraDBMessage对象快照数组。messageList:实时MessageList实例。使用它读取其他阶段,或就地添加、移除或替换消息。
文本位于 message.content.parts 中,而不是 message.content 本身。遍历 parts 并按 part.type === 'text' 过滤,以读取用户或助手文本。为兼容旧版,还提供扁平化的 message.content.content 字符串,可用作后备。有关完整详情,请参阅 Processor 参考中的消息参数。
转换输入消息转换输入消息的直接链接
import type { Processor, ProcessInputArgs } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class CustomInputProcessor implements Processor {
id = 'custom-input'
async processInput({ messages }: ProcessInputArgs): Promise<MastraDBMessage[]> {
// Transform messages before they reach the LLM.
// Text lives in content.parts — iterate parts and rewrite text parts only.
return messages.map(msg => ({
...msg,
content: {
...msg.content,
parts: msg.content.parts?.map(part =>
part.type === 'text' ? { ...part, text: part.text.toLowerCase() } : part,
),
},
}))
}
}
processInput() 方法接收 messages、systemMessages 和 abort() 函数。返回 MastraDBMessage[] 以替换消息,或返回 { messages, systemMessages } 以同时修改系统消息。
有关所有可用参数和返回类型,请参阅 Processor 参考。
控制每个步骤控制每个步骤的直接链接
processInput() 在 Agent 执行开始时运行一次,而 processInputStep() 会在 Agent 循环的每个步骤运行(包括 Tool 调用的后续步骤)。借助它,可以按步骤更改配置,例如在运行时切换模型或修改 Tool 选择。
import type {
Processor,
ProcessInputStepArgs,
ProcessInputStepResult,
} from '@mastra/core/processors'
export class DynamicModelProcessor implements Processor {
id = 'dynamic-model'
async processInputStep({
stepNumber,
model,
toolChoice,
messageList,
}: ProcessInputStepArgs): Promise<ProcessInputStepResult> {
// Use a fast model for initial response
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
// Disable tools after 5 steps to force completion
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
// No changes for other steps
return {}
}
}
该方法接收当前的 stepNumber、model、tools、toolChoice、messages 等。返回包含你希望在该步骤覆盖的任意属性的对象,例如 { model, toolChoice, tools, systemMessages }。
有关所有可用参数和返回类型,请参阅 Processor 参考。
在调用 Provider 前重写 LLM 请求在调用 Provider 前重写 LLM 请求的直接链接
需要重写 Mastra 发送给模型的最终提示词时,请使用 processLLMRequest()。此 hook 在 Mastra 将 MessageList 转换为面向 Provider 的提示词格式(LanguageModelV2Prompt)之后、调用 Provider 之前立即运行。
使用基于消息的 hook 更改对话:
processInput():在 Agent 循环开始前更改一次对话。processInputStep():在每次 LLM 调用前更改消息或步骤配置。processLLMRequest():仅更改当前 Provider 调用的出站提示词。
processLLMRequest() 返回的更改是临时的,不会持久化回 MessageList、Memory、UI 历史或未来的 Provider 调用。因此,此 hook 适合用于 Provider 兼容性重写、角色/内容规范化,或其他不应更改已存储对话历史的模型专用提示词更改。
该方法接收 prompt、model、stepNumber、steps、state 和共享的 Processor 上下文。从 processLLMRequest() 调用 abort() 会发出常规 tripwire 响应并停止调用。
有关所有可用参数和返回类型,请参阅 Processor 参考。
调用 Provider 后处理 LLM 响应调用 Provider 后处理 LLM 响应的直接链接
步骤完成且流块收集完毕后,使用 processLLMResponse() 处理完整的 LLM 响应。此 hook 与 processLLMRequest() 配对:在请求 hook 中保存状态(例如缓存键),然后在响应 hook 中读回状态,以执行写入缓存等副作用。
state 对象与同一步骤传递给 processLLMRequest() 的实例相同。当 fromCache 为 true 时,响应来自缓存重放,而不是实时模型调用;写入缓存的 Processor 此时应跳过写入。
该方法接收 chunks、model、stepNumber、steps、state、fromCache 和共享的 Processor 上下文。
有关所有可用参数和返回类型,请参阅 Processor 参考。
使用 prepareStep() 回调use-the-preparestep-callback的直接链接
generate() 或 stream() 上的 prepareStep() 回调是 processInputStep() 的简写。在内部,Mastra 将其封装到一个在每个步骤调用你的函数的 Processor 中。它接受与 processInputStep() 相同的参数和返回类型,但无需创建类:
await agent.generate('Complex task', {
prepareStep: async ({ stepNumber, model }) => {
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
},
})
转换输出消息转换输出消息的直接链接
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class CustomOutputProcessor implements Processor {
id = 'custom-output'
async processOutputResult({ messages }): Promise<MastraDBMessage[]> {
// Transform messages after the LLM generates them
return messages.filter(msg => msg.role !== 'system')
}
}
该方法还会接收包含完整生成数据的 result 对象、text、usage(token 计数)、finishReason 和 steps(每项都包含 toolCalls、toolResults 等)。使用这些数据跟踪用量或检查 Tool 调用:
import type { Processor } from '@mastra/core/processors'
export class UsageTracker implements Processor {
id = 'usage-tracker'
async processOutputResult({ messages, result }) {
console.log(`Tokens: ${result.usage.inputTokens} in, ${result.usage.outputTokens} out`)
console.log(`Finish reason: ${result.finishReason}`)
return messages
}
}
过滤流式输出过滤流式输出的直接链接
processOutputStream() 方法会在流块到达客户端之前转换或过滤它们:
import type { Processor } from '@mastra/core/processors'
import type { ChunkType } from '@mastra/core/stream'
export class StreamFilter implements Processor {
id = 'stream-filter'
async processOutputStream({ part }): Promise<ChunkType | null> {
// Drop text-delta chunks that contain the word "secret"
if (part.type === 'text-delta' && part.payload.text.includes('secret')) {
return null
}
// Return the (possibly modified) chunk to emit it
return part
}
}
返回值:
ChunkType会发出该块。返回原始part可原样传递。null或undefined会丢弃该块。两者行为相同,因此不返回任何值的方法也会丢弃该块。- 丢弃操作仅影响单个块。要完全停止流,请调用
abort()。
要同时接收 Tool 通过 writer.custom() 发出的自定义 data-* 块,请在 Processor 上设置 processDataParts = true。这样可以在 Tool 发出的数据块到达客户端之前检查、修改或阻止它们。
验证每个响应验证每个响应的直接链接
processOutputStep() 方法在每个 LLM 步骤后运行,让你可以验证响应,并选择请求重试:
import type { Processor } from '@mastra/core/processors'
export class ResponseValidator implements Processor {
id = 'response-validator'
async processOutputStep({ text, abort, retryCount }) {
const isValid = await validateResponse(text)
if (!isValid && retryCount < 3) {
abort('Response did not meet requirements. Try again.', { retry: true })
}
return []
}
}
有关重试行为的更多信息,请参阅高级模式中的重试机制。
跨块和步骤持久化数据跨块和步骤持久化数据的直接链接
输出方法接收一个在单次请求生命周期内持久存在的 state 对象。状态以 Processor 的 id 为键,因此每个 Processor 只能看到自己的数据,并在 processOutputStream、processOutputStep 和 processOutputResult 之间共享。每次新的 agent.generate() 或 agent.stream() 调用都会创建新的状态对象。
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
}
async processOutputResult({ messages, state }) {
console.log(`Total words: ${state.wordCount}`)
return messages
}
}
内置实用 Processor内置实用 Processor的直接链接
Mastra 为常见任务提供实用 Processor:
有关安全和验证 Processor,请参阅 Guardrail 页面,了解输入/输出 Guardrail 和内容审核 Processor。 有关 Memory 专用 Processor,请参阅 Memory Processor 页面,了解处理消息历史、语义召回和工作 Memory 的 Processor。
TokenLimitertokenlimiter的直接链接
当 token 总数超过指定限制时,通过移除较早的消息来防止上下文窗口溢出。优先保留近期消息和系统消息。
import { Agent } from '@mastra/core/agent'
import { TokenLimiter } from '@mastra/core/processors'
const agent = new Agent({
id: 'my-agent',
name: 'my-agent',
model: 'openai/gpt-5.6-sol',
inputProcessors: [new TokenLimiter(127000)],
})
有关自定义编码、策略和计数模式选项,请参阅 TokenLimiterProcessor 参考。
ToolCallFiltertoolcallfilter的直接链接
从发送到 LLM 的消息中移除 Tool 调用和结果,减少冗长 Tool 交互的 token 用量。也可以选择仅排除特定 Tool。此过滤器只影响 LLM 输入,过滤后的消息仍会保存到 Memory。
默认情况下,ToolCallFilter 会在 Agent 循环开始前过滤初始输入。使用 filterAfterToolSteps 还可在每个循环步骤中进行过滤,同时保留近期产生 Tool 的步骤。
new ToolCallFilter({
filterAfterToolSteps: 2,
})
设置 preserveModelOutput: true,为已过滤且完成的 Tool 结果保留精简的 toModelOutput 历史。过滤器仅保留面向模型的输出,并移除原始 Tool 参数和原始结果。
new ToolCallFilter({
preserveModelOutput: true,
})
有关配置选项,请参阅 ToolCallFilter 参考;有关 Memory 前过滤,请参阅 Memory Processor 页面。
ToolSearchProcessortoolsearchprocessor的直接链接
为拥有大型 Tool 库的 Agent 启用运行时 Tool 发现。Processor 不会预先提供所有 Tool,而是为 Agent 提供 search_tools 和 load_tool 元 Tool,按需根据关键字查找并加载 Tool,从而减少上下文 token 用量。
有关配置选项和用法示例,请参阅 ToolSearchProcessor 参考。
ProviderHistoryCompatproviderhistorycompat的直接链接
处理 Agent 在不同模型 Provider 之间复用消息时特定于 Provider 的历史不兼容问题。它可以在调用 Provider 前重写出站 LLM 请求,或从已知的 Provider API 错误中恢复并重试。
需要 Provider 历史兼容性规则、响应式 API 错误恢复、自定义兼容性规则或可预测的 Processor 顺序时,请显式添加 ProviderHistoryCompat。
有关设置、内置规则和自定义规则选项,请参阅 ProviderHistoryCompat 参考。
响应缓存响应缓存的直接链接
此功能处于 beta 阶段。在 API 稳定之前,可能会发生不伴随主版本号变更的破坏性更改。
当 Agent 收到相同请求时,响应缓存会跳过 LLM 调用并重放之前缓存的响应。使用它可降低延迟,并避免为重复调用付费。
缓存通过 ResponseCache 输入 Processor 实现。Mastra 不提供 Agent 级选项。要启用缓存,请显式注册该 Processor。这样能在 Mastra 收集反馈期间保持较小的 API 表面积。每次调用的覆盖项通过 RequestContext 传递。
何时使用响应缓存何时使用响应缓存的直接链接
当相同请求形状在多个用户或会话中重复出现时,请使用响应缓存,例如提示词模板、建议提示词按钮、Agent 搜索中的重复提问,或反复对相同输入进行分类的 Guardrail LLM。当调用通过 Tool 触发外部副作用时,请勿使用,因为缓存命中会重放 Tool 调用而不重新执行。
快速开始快速开始的直接链接
向 Agent 的 inputProcessors 添加 ResponseCache,并传入任意 MastraServerCache 作为后端。对于开发环境,InMemoryServerCache 开箱即用:
import { Agent } from '@mastra/core/agent'
import { InMemoryServerCache } from '@mastra/core/cache'
import { ResponseCache } from '@mastra/core/processors'
const cache = new InMemoryServerCache()
export const searchAgent = new Agent({
id: 'search-agent',
name: 'Search Agent',
instructions: 'You answer questions concisely.',
model: 'openai/gpt-5',
inputProcessors: [new ResponseCache({ cache, ttl: 600 })], // 10 minutes
})
首次调用会正常运行 LLM 并将响应写入缓存。随后使用相同已解析提示词的调用会返回缓存响应,而不会调用 LLM。
通过 RequestContext 按调用覆盖通过 RequestContext 按调用覆盖的直接链接
每次调用的配置通过 RequestContext 传递。使用 ResponseCache.context() 构建新上下文,或使用 ResponseCache.applyContext() 合并到现有上下文:
import { ResponseCache } from '@mastra/core/processors'
import { RequestContext } from '@mastra/core/request-context'
// Fresh context with the override
await agent.stream('hello', {
requestContext: ResponseCache.context({ key: 'custom-key', bust: true }),
})
// Or merge into an existing context
const ctx = new RequestContext()
ctx.set('caller-meta', { userId: 'u-123' })
ResponseCache.applyContext(ctx, { bust: true })
await agent.stream('hello', { requestContext: ctx })
以下字段可按调用覆盖:
key:字符串或函数。仅覆盖此请求自动派生的缓存键。scope:字符串或null。仅覆盖此请求的租户/用户作用域。null表示不使用作用域。bust:布尔值。跳过缓存读取,但完成后仍会写入(适用于“强制刷新”按钮)。
cache、ttl 和 agentId 保留在构造函数上。它们属于实例级配置,不宜按调用变化。
租户作用域租户作用域的直接链接
默认情况下,ResponseCache 会在请求上下文中查找 MASTRA_RESOURCE_ID_KEY,并将其用作缓存作用域。这意味着已经填充资源 ID(例如通过 Memory)的 Agent 会自动获得按用户隔离。用户永远不会看到彼此的缓存响应。
需要不同作用域时,请显式覆盖:
new Agent({
id: 'processors-agent',
inputProcessors: [
new ResponseCache({
cache,
scope: 'org-123', // explicit tenant scope
}),
],
})
传入 scope: null 可有意在所有调用者之间共享条目。仅对已知公开且非个性化的内容使用此设置。
自定义缓存后端自定义缓存后端的直接链接
ResponseCache 接受任意 MastraServerCache。在生产环境中,请使用 @mastra/redis 中的 RedisCache:
import { Agent } from '@mastra/core/agent'
import { ResponseCache } from '@mastra/core/processors'
import { RedisCache } from '@mastra/redis'
const cache = new RedisCache({ url: process.env.REDIS_URL })
export const agent = new Agent({
id: 'cached-agent',
name: 'Cached Agent',
instructions: '...',
model: 'openai/gpt-5',
inputProcessors: [new ResponseCache({ cache })],
})
对于自定义后端,请扩展 MastraServerCache 并实现其抽象方法(该 Processor 只调用 get 和 set)。
缓存的实现方式缓存的实现方式的直接链接
ResponseCache 挂接到 processLLMRequest(查找缓存,命中时短路)和 processLLMResponse(完成时写入缓存)。二者都在 Agent 循环中运行,位置是在 Memory 加载完毕且前面的输入 Processor 转换提示词_之后_。
这意味着缓存键派生自 Mastra 即将发送给模型的已解析 LanguageModelV2Prompt。该键在 Memory 加载完毕且前面的输入 Processor 运行_之后_创建,Agent Tool 循环中的每个步骤会单独缓存。
缓存键包含的内容缓存键包含的内容的直接链接
如果未提供 key,Processor 会根据会改变此步骤 LLM 响应的输入,确定性地派生缓存键:agentId、stepNumber(因此 Tool 循环中的每个步骤都有自己的缓存条目)、scope、模型标识(provider、modelId、规范版本),以及已解析的 prompt(Memory 后 + Processor 后)。这些输入的任何变化都会自动使缓存失效。
多模态提示词也包括在内。图像和文件部分按值进入缓存键:URL 会贡献完整的 href,内联二进制数据(Uint8Array、ArrayBuffer)会贡献其字节的摘要。因此,仅引用图像不同的两个请求会得到不同的缓存条目。
自定义缓存键自定义缓存键的直接链接
在构造函数上或按调用将 key 作为函数传入,以根据这些输入的任意子集派生自定义缓存键。该函数接收确定性哈希原本会使用的相同输入,并返回字符串(或 Promise<string>):
import { ResponseCache, buildResponseCacheKey } from '@mastra/core/processors'
await agent.stream(input, {
requestContext: ResponseCache.context({
// Cache only on the model id and the resolved prompt tail — ignore
// step number, scope, etc.
key: ({ model, prompt }) => `qa:${model.modelId}:${JSON.stringify(prompt).slice(-200)}`,
}),
})
// Or reuse the deterministic helper while overriding individual fields:
await agent.stream(input, {
requestContext: ResponseCache.context({
key: inputs => buildResponseCacheKey({ ...inputs, scope: 'global' }),
}),
})
如果函数抛出异常,Processor 会回退到默认缓存键派生方式,使调用仍能受益于缓存。
缓存命中的工作方式缓存命中的工作方式的直接链接
Processor 发现缓存命中时,会从 processLLMRequest 返回缓存块,使 LLM 调用短路。Agent 循环会根据这些块合成流,而不是调用模型。agent.generate() 将它们收集到 FullOutput 中;agent.stream() 返回块来自缓存缓冲区的 MastraModelOutput,因此遍历 fullStream 或等待 text、usage 和 finishReason 的使用者会看到缓存值。
响应完成后才会写入缓存。失败的运行(错误、tripwire 激活)不会缓存,因此下次调用可以正常重试。
高级模式高级模式的直接链接
使用 maxSteps 确保最终响应ensure-a-final-response-with-maxsteps的直接链接
使用 maxSteps 限制 Agent 执行时,如果 Agent 尝试在最后一步调用 Tool,可能会返回空响应。请结合 sendSignal 使用 processInputStep(),在最后一步注入响应式提醒。此方法通过附加信号而不是修改系统消息来保留提示词缓存。
import type { Processor, ProcessInputStepArgs } from '@mastra/core/processors'
export class EnsureFinalResponseProcessor implements Processor {
readonly id = 'ensure-final-response'
private maxSteps: number
constructor(maxSteps: number) {
this.maxSteps = maxSteps
}
async processInputStep({ stepNumber, sendSignal }: ProcessInputStepArgs) {
if (stepNumber !== this.maxSteps - 1) {
return
}
await sendSignal?.({
type: 'reactive',
contents:
`This is your final step (step ${stepNumber + 1} of ${this.maxSteps}). ` +
`Do not call any more tools. Summarize what you have found and give the user a complete final answer now.`,
attributes: { reason: 'max-steps-reached', step: stepNumber + 1 },
})
}
}
信号以 <system-reminder> 用户消息的形式传递,模型会在上下文中看到它:
<system-reminder reason="max-steps-reached" step="5">This is your final step (step 5 of 5). Do not call any more tools. Summarize what you have found and give the user a complete final answer now.</system-reminder>
将 Processor 添加到 inputProcessors,加入解释信号标签的系统提示词,并向 generate() 或 stream() 传入相同的 maxSteps 值:
import { Agent } from '@mastra/core/agent'
import { EnsureFinalResponseProcessor } from '../processors/ensure-final-response'
const MAX_STEPS = 5
const agent = new Agent({
id: 'agent',
instructions: `You are a helpful assistant.
Some messages you receive may contain <system-reminder>...</system-reminder> tags.
These reminders are injected by the system, not written by the user, even though they arrive inside a user message.
Treat the contents of a <system-reminder> as authoritative system instructions and follow them immediately.
Do not mention the reminder to the user or quote the tags back to them.`,
inputProcessors: [new EnsureFinalResponseProcessor(MAX_STEPS)],
// ...
})
await agent.generate('Your prompt', { maxSteps: MAX_STEPS })
响应式信号默认使用 tagName: 'system-reminder'。有关 Processor 发出的信号的更多信息,请参阅 Signal。
传递提醒但不保留传递提醒但不保留的直接链接
默认情况下,从 Processor 发送的信号会成为对话的一部分:它会写入存储,并在后续轮次重新进入提示词。对于每轮都重新注入的指令,这并不理想,因为副本会不断累积,模型也会开始把自己过去的提醒视为要模仿的先前上下文。设置 transient: true,仅在当前调用中将信号传递给模型,而不保留它。
何时使用: 随着对话增长,你希望在模型的近期窗口中保留一条简短的引导指令,例如“专注于当前任务”“回答不超过三句话”,或依赖实时应用状态的每轮约束。每轮重新注入,使其保持在最新消息附近。
import type { Processor, ProcessInputStepArgs } from '@mastra/core/processors'
export class SteeringReminderProcessor implements Processor {
readonly id = 'steering-reminder'
async processInputStep({ sendSignal }: ProcessInputStepArgs) {
await sendSignal?.({
type: 'reactive',
contents: 'Stay on the current task and keep answers under three sentences.',
transient: true,
})
}
}
临时信号仍会出现在当前调用的提示词中,因此模型会在最新轮次附近看到它。由于它不会保留,每轮重新发送只会在上下文中保留一个最新副本,而不会累积历史;它也不会出现在已存储的线程历史中。由于没有写入任何内容,还能在各轮之间保持稳定的提示词缓存前缀。
发出自定义流事件发出自定义流事件的直接链接
输出 Processor 接收一个 writer 对象,让你可以在流式传输期间向客户端发回自定义数据块。这适合流式传输审核结果,或在不阻塞原始流的情况下发送 UI 更新信号等用例。
import type { Processor } from '@mastra/core/processors'
export class ModerationProcessor implements Processor {
id = 'moderation'
async processOutputResult({ messages, writer }) {
// Run moderation on the final output
const text = messages
.filter(m => m.role === 'assistant')
.flatMap(m => m.content.parts?.filter(p => p.type === 'text'))
.map(p => p.text)
.join(' ')
const result = await runModeration(text)
if (result.requiresChange) {
// Emit a custom event to the client with the moderated text
await writer?.custom({
type: 'data-moderation-update',
data: {
originalText: text,
moderatedText: result.moderatedText,
reason: result.reason,
},
})
}
return messages
}
}
在客户端监听流中的自定义块类型:
const stream = await agent.stream('Hello')
for await (const chunk of stream.fullStream) {
if (chunk.type === 'data-moderation-update') {
// Update the UI with moderated text
updateDisplayedMessage(chunk.data.moderatedText)
}
}
自定义块类型必须使用 data- 前缀(例如 data-moderation-update、data-status)。
默认情况下,processOutputStream() 会跳过 data-* 块,以免意外处理 Tool 遥测数据或其他 Processor 的输出。要在 Processor 中检查、修改或阻止这些块,请在该 Processor 上设置 processDataParts = true:
class ModerationCollector implements Processor {
id = 'moderation-collector'
processDataParts = true
async processOutputStream({ part, state }) {
if (part.type === 'data-moderation-update') {
state.warnings ??= []
state.warnings.push(part.data)
}
return part
}
}
向消息添加元数据向消息添加元数据的直接链接
可以在 processOutputResult 中向消息添加自定义元数据。可通过响应对象访问这些元数据:
import type { Processor } from '@mastra/core/processors'
import type { MastraDBMessage } from '@mastra/core/memory'
export class MetadataProcessor implements Processor {
id = 'metadata-processor'
async processOutputResult({
messages,
}: {
messages: MastraDBMessage[]
}): Promise<MastraDBMessage[]> {
return messages.map(msg => {
if (msg.role === 'assistant') {
return {
...msg,
content: {
...msg.content,
metadata: {
...msg.content.metadata,
processedAt: new Date().toISOString(),
customData: 'your data here',
},
},
}
}
return msg
})
}
}
使用 generate() 访问元数据:
const result = await agent.generate('Hello')
// The response includes uiMessages with processor-added metadata
const assistantMessage = result.response?.uiMessages?.find(m => m.role === 'assistant')
console.log(assistantMessage?.metadata?.customData)
对于流式传输,请从 finish 块 payload 或 stream.response Promise 访问元数据。
将 Workflow 用作 Processor将 Workflow 用作 Processor的直接链接
可以将 Mastra Workflow 用作 Processor,创建具有并行执行、条件分支和错误处理能力的复杂处理管道:
import { createWorkflow, createStep } from '@mastra/core/workflows'
import {
ProcessorStepSchema,
PromptInjectionDetector,
PIIDetector,
ModerationProcessor,
} from '@mastra/core/processors'
import { Agent } from '@mastra/core/agent'
// Create a workflow that runs multiple checks in parallel
const moderationWorkflow = createWorkflow({
id: 'moderation-pipeline',
inputSchema: ProcessorStepSchema,
outputSchema: ProcessorStepSchema,
})
.parallel([
createStep(
new PIIDetector({
strategy: 'redact',
}),
),
createStep(
new PromptInjectionDetector({
strategy: 'block',
}),
),
createStep(
new ModerationProcessor({
strategy: 'block',
}),
),
])
.map(async ({ inputData }) => {
return inputData['processor:pii-detector']
})
.commit()
// Use the workflow as an input processor
const agent = new Agent({
id: 'moderated-agent',
name: 'Moderated Agent',
model: 'openai/gpt-5.6-sol',
inputProcessors: [moderationWorkflow],
})
在 .parallel() 步骤后,每个分支结果都以其 Processor ID 为键(例如 processor:pii-detector)。使用 .map() 选择下一步骤应接收哪个分支的输出。
如果某个分支使用 redact 等变更策略,请映射到该分支,使其转换后的消息继续传递。如果所有分支都只使用 block,则任意分支均可。因为它们都不修改消息,可以任选其一。
在 Mastra 中注册 Agent 后,Processor Workflow 会自动注册为 Workflow,供你在 Studio 中查看和调试。
重试机制重试机制的直接链接
Processor 可以请求 LLM 根据反馈重试响应。这适合实现质量检查、输出验证或迭代优化:
import type { Processor } from '@mastra/core/processors'
export class QualityChecker implements Processor {
id = 'quality-checker'
async processOutputStep({ text, abort, retryCount }) {
const qualityScore = await evaluateQuality(text)
if (qualityScore < 0.7 && retryCount < 3) {
// Request a retry with feedback for the LLM
abort('Response quality score too low. Please provide a more detailed answer.', {
retry: true,
metadata: { score: qualityScore },
})
}
return []
}
}
const agent = new Agent({
id: 'quality-agent',
name: 'Quality Agent',
model: 'openai/gpt-5.6-sol',
outputProcessors: [new QualityChecker()],
maxProcessorRetries: 3, // Maximum retry attempts. If unset, retries are disabled (unless errorProcessors are configured, in which case it defaults to 10).
})
重试机制:
- 可在
processOutputStep()和processInputStep()方法中使用 - 重新执行步骤,并将中止原因作为 LLM 的上下文
- 通过
retryCount参数跟踪重试次数 - 需要在 Agent 或调用上显式设置
maxProcessorRetries限制
错误 Processor 重试限制错误 Processor 重试限制的直接链接
processAPIError() 有单独的默认值:配置 errorProcessors 且省略 maxProcessorRetries 时,运行时最多允许重试 10 次。需要限定重试预算时,请显式设置限制。
对于 StreamErrorRetryProcessor,还应将其 maxRetries 设置为相同值。它自身默认为 1,否则可能低于 Agent 上限。当 Processor 是请求的唯一重试机制时,请将模型重试次数保持为 0。
违规回调违规回调的直接链接
所有 Processor 都公开 onViolation 属性,每当检测到策略违规时触发,包括调用 abort()(阻止策略)和 Processor 发出警告(警告策略)的情况。使用它执行警报、日志记录或副作用,而不影响 Processor 的主要逻辑:
import { ModerationProcessor, CostGuardProcessor } from '@mastra/core/processors'
const moderation = new ModerationProcessor({
model: 'openai/gpt-5-nano',
strategy: 'block',
})
moderation.onViolation = ({ processorId, message, detail }) => {
// Log to external monitoring, send alerts, update dashboards
monitor.track('processor_violation', { processorId, message, detail })
}
const costGuard = new CostGuardProcessor({
maxCost: 10.0,
scope: 'resource',
window: '30d',
})
costGuard.onViolation = ({ processorId, message, detail }) => {
alertSystem.notify(`[${processorId}] ${message}`)
}
回调接收包含以下内容的 ProcessorViolation 对象:
processorId:检测到违规的 Processor IDmessage:对违规内容的可读描述detail:Processor 专用元数据(例如成本用量、检测到的 PII 类型、审核类别)
onViolation 是基础 Processor 接口的一部分,因此任何自定义 Processor 也都可以使用。任何 Processor 调用 abort() 时,运行程序会自动调用它。回调中抛出的错误会被静默捕获,以免干扰 Processor 管道。
中止和 tripwire 块中止和 tripwire 块的直接链接
调用 abort(reason, options) 会抛出结束处理的 TripWire 错误。在流中,Mastra 会发出客户端可检测的 tripwire 块:
for await (const chunk of stream.fullStream) {
if (chunk.type === 'tripwire') {
console.log('Blocked by', chunk.payload.processorId, '-', chunk.payload.reason)
break
}
}
对于 agent.generate(),结果通过 result.tripwire 公开相同信息,同时 result.finishReason === 'other'。
abort 接受第二个 options 参数:
retry: true请求 Agent 重试,而不是结束。输入和输出 Processor 重试要求在 Agent 或调用上设置maxProcessorRetries。metadata将结构化数据附加到tripwire块,使下游使用者可以根据pii、quality或moderation等类别进行分支。
API 错误处理API 错误处理的直接链接
processAPIError 方法处理 LLM API 拒绝,即 API 拒绝请求(例如状态码 400 或 422)的错误,而不是网络或服务器故障。这样可以在 API 拒绝消息格式时修改请求并重试。
import { APICallError } from '@ai-sdk/provider'
import type { Processor, ProcessAPIErrorArgs, ProcessAPIErrorResult } from '@mastra/core/processors'
export class ContextLengthHandler implements Processor {
id = 'context-length-handler'
processAPIError({
error,
messageList,
retryCount,
}: ProcessAPIErrorArgs): ProcessAPIErrorResult | void {
if (retryCount > 0) return
if (APICallError.isInstance(error) && error.message.includes('context length exceeded')) {
const messages = messageList.get.all.db()
if (messages.length > 4) {
messageList.removeByIds([messages[1]!.id, messages[2]!.id])
return { retry: true }
}
}
}
}
Mastra 包含内置的 PrefillErrorHandler,可自动处理 Anthropic 的“assistant message prefill”错误。此 Processor 会自动注入,无需配置。
相关文档相关文档的直接链接
- Guardrail:安全和验证 Processor
- Memory Processor:Memory 专用 Processor 和自动集成
- Processor 接口:Processor 的完整 API 参考
- ToolSearchProcessor 参考:运行时 Tool 搜索的 API 参考