> 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 参考