> Discover all available pages from the documentation index: https://mastra.zisheng.pro/llms.txt
# Processor
Processor 会在消息经过 Agent 时对其进行转换、验证或控制。它们在 Agent 执行管道的特定位置运行,让你可以在输入到达语言模型之前修改输入,或在输出返回用户之前修改输出。
Processor 的配置方式如下:
- **`inputProcessors`**:在消息到达语言模型之前运行。
- **`outputProcessors`**:在语言模型生成响应之后、响应返回用户之前运行。
你可以使用单独的 [`Processor`](https://mastra.zisheng.pro/reference/processors/processor-interface) 对象,也可以使用 Mastra 的 Workflow 原语将它们组合成 Workflow。Workflow 让你能够更细致地控制 Processor 执行顺序、并行处理和条件逻辑。
某些 Processor 同时实现输入和输出逻辑,可以根据转换所需发生的位置用于任一数组。
某些内置 Processor 还会发送隐藏的系统提醒信号。这些信号会持久化到原始 Memory 历史中,并在下一次模型调用前转换为 `...` 上下文;但面向 UI 的标准消息转换和默认 Memory 召回会隐藏它们,除非你明确选择显示。要为当前调用传递信号而不保留它,请使用 `transient: true` 发送。
## 何时使用 Processor
使用 Processor 可以:
- 规范化或验证用户输入
- 为 Agent 添加 Guardrail
- 检测并阻止提示词注入或越狱尝试
- 出于安全或合规目的审核内容
- 转换消息(例如翻译语言、过滤 Tool 调用)
- 限制 token 用量或消息历史长度
- 遮盖敏感信息(PII)
- 对消息应用自定义业务逻辑
Mastra 包含多个适用于常见用例的 Processor。你还可以根据应用的特定要求创建自定义 Processor。
## 快速开始
导入并实例化 Processor,然后将其传入 Agent 的 `inputProcessors` 或 `outputProcessors` 数组:
```typescript
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 按照它们在数组中的顺序运行:
```typescript
inputProcessors: [new UnicodeNormalizer(), new PromptInjectionDetector(), new ModerationProcessor()]
```
对于输出 Processor,此顺序决定对模型响应应用转换的先后次序。
### 启用 Memory 时
在 Agent 上启用 Memory 后,Memory Processor 会自动添加到管道中:
**输入 Processor:**
```text
[Memory Processors] → [Your inputProcessors]
```
Memory 先加载消息历史,然后运行你的 Processor。
**输出 Processor:**
```text
[Your outputProcessors] → [Memory Processors]
```
你的 Processor 先运行,然后由 Memory 持久化消息。
按照此顺序,调用 `abort()` 的输出 Guardrail 会跳过 Memory Processor,并阻止保存消息。有关详细信息,请参阅 [Memory Processor](https://mastra.zisheng.pro/docs/memory/memory-processors)。
## 将 Processor 附加到 Agent
通过三个数组在 Agent 上配置 Processor:
```typescript
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:
```typescript
new Agent({
id: 'processors-agent',
inputProcessors: ({ requestContext }) => {
const limit = requestContext.get('tokenLimit') ?? 4000
return [new TokenLimiter(limit)]
},
})
```
### 按调用覆盖 Processor
`agent.generate()` 和 `agent.stream()` 接受相同的三个数组。传入其中一个时,它仅针对该次调用**替换** Agent 上匹配的数组。Memory、Workspace 和其他由框架管理的 Processor 仍会在你的数组前后运行。
```typescript
await agent.stream('Summarize this', {
inputProcessors: [new TokenLimiter(2000)],
maxProcessorRetries: 5,
})
```
## 创建自定义 Processor
自定义 Processor 实现 `Processor` 接口。
Processor 方法接收两个用于访问对话的参数:
- `messages`:当前阶段的 `MastraDBMessage` 对象快照数组。
- `messageList`:实时 `MessageList` 实例。使用它读取其他阶段,或就地添加、移除或替换消息。
文本位于 `message.content.parts` 中,而不是 `message.content` 本身。遍历 `parts` 并按 `part.type === 'text'` 过滤,以读取用户或助手文本。为兼容旧版,还提供扁平化的 `message.content.content` 字符串,可用作后备。有关完整详情,请参阅 `Processor` 参考中的[消息参数](https://mastra.zisheng.pro/reference/processors/processor-interface)。
### 转换输入消息
```typescript
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 {
// 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` 参考](https://mastra.zisheng.pro/reference/processors/processor-interface)。
### 控制每个步骤
`processInput()` 在 Agent 执行开始时运行一次,而 `processInputStep()` 会在 Agent 循环的**每个步骤**运行(包括 Tool 调用的后续步骤)。借助它,可以按步骤更改配置,例如在运行时切换模型或修改 Tool 选择。
```typescript
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 {
// 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` 参考](https://mastra.zisheng.pro/reference/processors/processor-interface)。
### 在调用 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` 参考](https://mastra.zisheng.pro/reference/processors/processor-interface)。
### 调用 Provider 后处理 LLM 响应
步骤完成且流块收集完毕后,使用 `processLLMResponse()` 处理完整的 LLM 响应。此 hook 与 `processLLMRequest()` 配对:在请求 hook 中保存状态(例如缓存键),然后在响应 hook 中读回状态,以执行写入缓存等副作用。
`state` 对象与同一步骤传递给 `processLLMRequest()` 的实例相同。当 `fromCache` 为 `true` 时,响应来自缓存重放,而不是实时模型调用;写入缓存的 Processor 此时应跳过写入。
该方法接收 `chunks`、`model`、`stepNumber`、`steps`、`state`、`fromCache` 和共享的 Processor 上下文。
有关所有可用参数和返回类型,请参阅 [`Processor` 参考](https://mastra.zisheng.pro/reference/processors/processor-interface)。
### 使用 `prepareStep()` 回调
`generate()` 或 `stream()` 上的 `prepareStep()` 回调是 `processInputStep()` 的简写。在内部,Mastra 将其封装到一个在每个步骤调用你的函数的 Processor 中。它接受与 `processInputStep()` 相同的参数和返回类型,但无需创建类:
```typescript
await agent.generate('Complex task', {
prepareStep: async ({ stepNumber, model }) => {
if (stepNumber === 0) {
return { model: 'openai/gpt-5-mini' }
}
if (stepNumber > 5) {
return { toolChoice: 'none' }
}
},
})
```
### 转换输出消息
```typescript
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 {
// Transform messages after the LLM generates them
return messages.filter(msg => msg.role !== 'system')
}
}
```
该方法还会接收包含完整生成数据的 `result` 对象、`text`、`usage`(token 计数)、`finishReason` 和 `steps`(每项都包含 `toolCalls`、`toolResults` 等)。使用这些数据跟踪用量或检查 Tool 调用:
```typescript
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()` 方法会在流块到达客户端之前转换或过滤它们:
```typescript
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 {
// 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 步骤后运行,让你可以验证响应,并选择请求重试:
```typescript
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 []
}
}
```
有关重试行为的更多信息,请参阅高级模式中的[重试机制](#retry-mechanism)。
### 跨块和步骤持久化数据
输出方法接收一个在单次请求生命周期内持久存在的 `state` 对象。状态以 Processor 的 `id` 为键,因此每个 Processor 只能看到自己的数据,并在 `processOutputStream`、`processOutputStep` 和 `processOutputResult` 之间共享。每次新的 `agent.generate()` 或 `agent.stream()` 调用都会创建新的状态对象。
```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
}
async processOutputResult({ messages, state }) {
console.log(`Total words: ${state.wordCount}`)
return messages
}
}
```
## 内置实用 Processor
Mastra 为常见任务提供实用 Processor:
**有关安全和验证 Processor**,请参阅 [Guardrail](https://mastra.zisheng.pro/docs/agents/guardrails) 页面,了解输入/输出 Guardrail 和内容审核 Processor。 **有关 Memory 专用 Processor**,请参阅 [Memory Processor](https://mastra.zisheng.pro/docs/memory/memory-processors) 页面,了解处理消息历史、语义召回和工作 Memory 的 Processor。
### `TokenLimiter`
当 token 总数超过指定限制时,通过移除较早的消息来防止上下文窗口溢出。优先保留近期消息和系统消息。
```typescript
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` 参考](https://mastra.zisheng.pro/reference/processors/token-limiter-processor)。
### `ToolCallFilter`
从发送到 LLM 的消息中移除 Tool 调用和结果,减少冗长 Tool 交互的 token 用量。也可以选择仅排除特定 Tool。此过滤器只影响 LLM 输入,过滤后的消息仍会保存到 Memory。
默认情况下,`ToolCallFilter` 会在 Agent 循环开始前过滤初始输入。使用 `filterAfterToolSteps` 还可在每个循环步骤中进行过滤,同时保留近期产生 Tool 的步骤。
```typescript
new ToolCallFilter({
filterAfterToolSteps: 2,
})
```
设置 `preserveModelOutput: true`,为已过滤且完成的 Tool 结果保留精简的 `toModelOutput` 历史。过滤器仅保留面向模型的输出,并移除原始 Tool 参数和原始结果。
```typescript
new ToolCallFilter({
preserveModelOutput: true,
})
```
有关配置选项,请参阅 [`ToolCallFilter` 参考](https://mastra.zisheng.pro/reference/processors/tool-call-filter);有关 Memory 前过滤,请参阅 [Memory Processor](https://mastra.zisheng.pro/docs/memory/memory-processors) 页面。
### `ToolSearchProcessor`
为拥有大型 Tool 库的 Agent 启用运行时 Tool 发现。Processor 不会预先提供所有 Tool,而是为 Agent 提供 `search_tools` 和 `load_tool` 元 Tool,按需根据关键字查找并加载 Tool,从而减少上下文 token 用量。
有关配置选项和用法示例,请参阅 [`ToolSearchProcessor` 参考](https://mastra.zisheng.pro/reference/processors/tool-search-processor)。
### `ProviderHistoryCompat`
处理 Agent 在不同模型 Provider 之间复用消息时特定于 Provider 的历史不兼容问题。它可以在调用 Provider 前重写出站 LLM 请求,或从已知的 Provider API 错误中恢复并重试。
需要 Provider 历史兼容性规则、响应式 API 错误恢复、自定义兼容性规则或可预测的 Processor 顺序时,请显式添加 `ProviderHistoryCompat`。
有关设置、内置规则和自定义规则选项,请参阅 [`ProviderHistoryCompat` 参考](https://mastra.zisheng.pro/reference/processors/provider-history-compat)。
## 响应缓存
> **Beta:** 此功能处于 beta 阶段。在 API 稳定之前,可能会发生不伴随主版本号变更的破坏性更改。
当 Agent 收到相同请求时,响应缓存会跳过 LLM 调用并重放之前缓存的响应。使用它可降低延迟,并避免为重复调用付费。
缓存通过 [`ResponseCache`](https://mastra.zisheng.pro/reference/processors/response-cache) 输入 Processor 实现。Mastra 不提供 Agent 级选项。要启用缓存,请显式注册该 Processor。这样能在 Mastra 收集反馈期间保持较小的 API 表面积。每次调用的覆盖项通过 `RequestContext` 传递。
### 何时使用响应缓存
当相同请求形状在多个用户或会话中重复出现时,请使用响应缓存,例如提示词模板、建议提示词按钮、Agent 搜索中的重复提问,或反复对相同输入进行分类的 Guardrail LLM。当调用通过 Tool 触发外部副作用时,请勿使用,因为缓存命中会重放 Tool 调用而不重新执行。
### 快速开始
向 Agent 的 `inputProcessors` 添加 `ResponseCache`,并传入任意 `MastraServerCache` 作为后端。对于开发环境,`InMemoryServerCache` 开箱即用:
```typescript
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` 传递。使用 `ResponseCache.context()` 构建新上下文,或使用 `ResponseCache.applyContext()` 合并到现有上下文:
```typescript
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 会自动获得按用户隔离。用户永远不会看到彼此的缓存响应。
需要不同作用域时,请显式覆盖:
```typescript
new Agent({
id: 'processors-agent',
inputProcessors: [
new ResponseCache({
cache,
scope: 'org-123', // explicit tenant scope
}),
],
})
```
传入 `scope: null` 可有意在所有调用者之间共享条目。仅对已知公开且非个性化的内容使用此设置。
### 自定义缓存后端
`ResponseCache` 接受任意 `MastraServerCache`。在生产环境中,请使用 `@mastra/redis` 中的 `RedisCache`:
```typescript
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`):
```typescript
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` 确保最终响应
使用 `maxSteps` 限制 Agent 执行时,如果 Agent 尝试在最后一步调用 Tool,可能会返回空响应。请结合 `sendSignal` 使用 `processInputStep()`,在最后一步注入响应式提醒。此方法通过附加信号而不是修改系统消息来保留提示词缓存。
```typescript
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 },
})
}
}
```
信号以 `` 用户消息的形式传递,模型会在上下文中看到它:
```xml
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.
```
将 Processor 添加到 `inputProcessors`,加入解释信号标签的系统提示词,并向 `generate()` 或 `stream()` 传入相同的 `maxSteps` 值:
```typescript
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 ... 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 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](https://mastra.zisheng.pro/docs/long-running-agents/signals)。
### 传递提醒但不保留
默认情况下,从 Processor 发送的信号会成为对话的一部分:它会写入存储,并在后续轮次重新进入提示词。对于每轮都重新注入的指令,这并不理想,因为副本会不断累积,模型也会开始把自己过去的提醒视为要模仿的先前上下文。设置 `transient: true`,仅在当前调用中将信号传递给模型,而不保留它。
**何时使用:** 随着对话增长,你希望在模型的近期窗口中保留一条简短的引导指令,例如“专注于当前任务”“回答不超过三句话”,或依赖实时应用状态的每轮约束。每轮重新注入,使其保持在最新消息附近。
```typescript
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 更新信号等用例。
```typescript
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
}
}
```
在客户端监听流中的自定义块类型:
```typescript
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`:
```typescript
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` 中向消息添加自定义元数据。可通过响应对象访问这些元数据:
```typescript
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 {
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()` 访问元数据:
```typescript
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
可以将 Mastra Workflow 用作 Processor,创建具有并行执行、条件分支和错误处理能力的复杂处理管道:
```typescript
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](https://mastra.zisheng.pro/docs/studio/overview) 中查看和调试。
### 重试机制
Processor 可以请求 LLM 根据反馈重试响应。这适合实现质量检查、输出验证或迭代优化:
```typescript
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 重试限制
`processAPIError()` 有单独的默认值:配置 `errorProcessors` 且省略 `maxProcessorRetries` 时,运行时最多允许重试 `10` 次。需要限定重试预算时,请显式设置限制。
对于 `StreamErrorRetryProcessor`,还应将其 `maxRetries` 设置为相同值。它自身默认为 `1`,否则可能低于 Agent 上限。当 Processor 是请求的唯一重试机制时,请将模型重试次数保持为 `0`。
### 违规回调
所有 Processor 都公开 `onViolation` 属性,每当检测到策略违规时触发,包括调用 `abort()`(阻止策略)和 Processor 发出警告(警告策略)的情况。使用它执行警报、日志记录或副作用,而不影响 Processor 的主要逻辑:
```typescript
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 ID
- `message`:对违规内容的可读描述
- `detail`:Processor 专用元数据(例如成本用量、检测到的 PII 类型、审核类别)
`onViolation` 是基础 [`Processor` 接口](https://mastra.zisheng.pro/reference/processors/processor-interface)的一部分,因此任何自定义 Processor 也都可以使用。任何 Processor 调用 `abort()` 时,运行程序会自动调用它。回调中抛出的错误会被静默捕获,以免干扰 Processor 管道。
### 中止和 tripwire 块
调用 `abort(reason, options)` 会抛出结束处理的 `TripWire` 错误。在流中,Mastra 会发出客户端可检测的 `tripwire` 块:
```typescript
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 错误处理
`processAPIError` 方法处理 LLM API 拒绝,即 API 拒绝请求(例如状态码 400 或 422)的错误,而不是网络或服务器故障。这样可以在 API 拒绝消息格式时修改请求并重试。
```typescript
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`](https://mastra.zisheng.pro/reference/processors/prefill-error-handler),可自动处理 Anthropic 的“assistant message prefill”错误。此 Processor 会自动注入,无需配置。
## 相关文档
- [Guardrail](https://mastra.zisheng.pro/docs/agents/guardrails):安全和验证 Processor
- [Memory Processor](https://mastra.zisheng.pro/docs/memory/memory-processors):Memory 专用 Processor 和自动集成
- [Processor 接口](https://mastra.zisheng.pro/reference/processors/processor-interface):Processor 的完整 API 参考
- [ToolSearchProcessor 参考](https://mastra.zisheng.pro/reference/processors/tool-search-processor):运行时 Tool 搜索的 API 参考