> Discover all available pages from the documentation index: https://mastra.zisheng.pro/llms.txt # 构建 Signal Provider 在本指南中,你将构建一个 Signal Provider,它会按固定时间间隔轮询外部服务,并在受监视的资源发生变化时向 Agent 线程推送通知。你将学习如何扩展 `SignalProvider` 基类、跟踪订阅、从轮询循环发出通知,以及在 Agent 上注册 Provider。 本示例会监视一个模拟 CI 服务中的构建流水线,但此模式适用于任何基于拉取的数据源:问题跟踪器、状态 API、队列或你自己的后端。 > **Beta:** Signal Provider 目前处于 beta 阶段。在 API 稳定之前,可能会出现不伴随主版本升级的破坏性变更。 ## 前提条件 - 已安装 Node.js `v22.13.0` 或更高版本 - 受支持的 [Model Provider](https://mastra.zisheng.pro/models) 提供的 API 密钥 - 现有的 Mastra 项目。如有需要,请按照[安装指南](https://mastra.zisheng.pro/guides/getting-started/quickstart)操作。 本指南还假设你对 [Signal](https://mastra.zisheng.pro/docs/long-running-agents/signals) 有基本了解。有关完整 API,请参阅 [`SignalProvider` 参考文档](https://mastra.zisheng.pro/reference/signals/signal-provider)。 ## 添加通知存储 Signal Provider 会将[通知 Signal](https://mastra.zisheng.pro/docs/long-running-agents/signals)推送到线程,而通知需要支持通知域的存储 adapter。请在 Mastra 实例上配置存储。 ```typescript import { Mastra } from '@mastra/core' import { LibSQLStore } from '@mastra/libsql' export const mastra = new Mastra({ storage: new LibSQLStore({ id: 'mastra-storage', url: 'file:./mastra.db', }), }) ``` LibSQL、PostgreSQL 和 MongoDB 都支持通知记录。如果没有通知存储,Provider 的 `notify()` 调用会在运行时抛出错误。 ## 创建外部服务 client 真实 Provider 会调用外部 API。为了让本指南可独立运行,请创建一个小型模拟 CI client,用于返回流水线的构建状态。之后再将它替换成真实的 API client。 ```typescript export type BuildStatus = { id: string pipeline: string status: 'passed' | 'failed' | 'running' } // Returns a random status so you can see notifications fire while testing. export async function fetchBuildStatus(pipeline: string): Promise { const states: BuildStatus['status'][] = ['passed', 'failed', 'running'] const status = states[Math.floor(Math.random() * states.length)]! return { id: `build_${Date.now()}`, pipeline, status } } ``` 该 client 每次调用时都会返回随机状态。接入真实 API 后,只需更改此文件。 ## 构建 Signal Provider 扩展 `SignalProvider`,实现抽象的 `id` 字段、设置 `pollInterval`,并覆盖 `poll()`。基类会按设定间隔调用 `poll()`,并传入所有活动订阅。只对你关注的构建发出通知。 ```typescript import { SignalProvider } from '@mastra/core/signals' import type { SignalProviderTarget, SignalSubscription } from '@mastra/core/signals' import { fetchBuildStatus } from './ci-client' export class CiSignals extends SignalProvider<'ci-signals'> { readonly id = 'ci-signals' as const readonly pollInterval = 10_000 // poll every 10 seconds // Public API so callers can subscribe a thread to a pipeline. watch(target: SignalProviderTarget, pipeline: string): SignalSubscription { return this.subscribe(target, pipeline) } unwatch(target: SignalProviderTarget, pipeline: string): boolean { return this.unsubscribe(target, pipeline) } async poll(subscriptions: SignalSubscription[]): Promise { for (const sub of subscriptions) { const build = await fetchBuildStatus(sub.externalResourceId) if (build.status !== 'failed') continue await this.notify( { source: this.id, kind: 'ci-status', priority: 'high', summary: `Build failed for ${sub.externalResourceId}`, payload: build, dedupeKey: `${this.id}:${sub.externalResourceId}:${build.id}`, }, { resourceId: sub.resourceId, threadId: sub.threadId }, ) } } } ``` 需要注意以下几点: - `subscribe()` 和 `unsubscribe()` 在基类中是 protected 方法。请用你自己的 public 方法(`watch` / `unwatch`)封装它们,以便调用方管理订阅。 - `externalResourceId` 可以是 Provider 专用的任意字符串。这里使用流水线名称;GitHub Provider 可能会使用 `"github:owner/repo#123"`。 - `notify()` 会将通知 Signal 转发到已连接 Agent 的线程。如果 Provider 从未注册到 Agent,它会抛出错误。 - `dedupeKey` 可防止同一次失败被存储两次。 ## 在 Agent 上注册 Provider 通过 `signals` 将 Provider 传给 Agent。Agent 会连接它并自动启动轮询循环。 ```typescript import { Agent } from '@mastra/core/agent' import { CiSignals } from '../signals/ci-signals' export const ciSignals = new CiSignals() export const devAgent = new Agent({ id: 'dev-agent', name: 'Dev Agent', instructions: 'Help the user triage CI build failures.', model: 'openai/gpt-5.6-sol', signals: [ciSignals], }) ``` 向 Mastra 注册 Agent,并添加第一步中的存储。 ```typescript import { Mastra } from '@mastra/core' import { LibSQLStore } from '@mastra/libsql' import { devAgent } from './agents/dev-agent' export const mastra = new Mastra({ agents: { devAgent }, storage: new LibSQLStore({ id: 'mastra-storage', url: 'file:./mastra.db', }), }) ``` ## 订阅线程 Provider 只会轮询线程正在监视的资源。订阅某个流水线,让 `poll()` 有资源可供检查。 ```typescript import { ciSignals } from './agents/dev-agent' ciSignals.watch({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme/app:main') ``` 在 Agent 注册后运行一次这段代码,例如通过设置脚本或 API 路由。订阅位于 Provider 的内存注册表中,因此重启后需要重新订阅。 ## 测试 Signal Provider 启动开发服务器: **npm**: ```bash npm run dev ``` **pnpm**: ```bash pnpm run dev ``` **Yarn**: ```bash yarn dev ``` **Bun**: ```bash bun run dev ``` 确保某个线程已订阅,然后查看日志。模拟 client 每次轮询都会返回随机状态,因此在几个周期内,你就会看到构建失败为已订阅线程触发通知。 要查看 Agent 对通知的反应,请订阅该线程并进行流式传输: ```typescript const subscription = await devAgent.subscribeToThread({ resourceId: 'user_123', threadId: 'thread_456', }) for await (const chunk of subscription.stream) { console.log(chunk) } ``` 构建失败时,模型会把通知作为上下文接收: ```xml Build failed for acme/app:main ``` 由于模拟 client 会随机生成状态,且模型可以自由组织回复,因此输出并不确定,确切措辞也会有所不同。 ## 后续步骤 你可以扩展此 Signal Provider: - 将 `fetchBuildStatus()` 替换为真实 API client。 - 持久化订阅,使其在重启后仍然存在,然后在 [`start()`](https://mastra.zisheng.pro/reference/signals/signal-provider) 中恢复它们。 - 为基于推送的数据源添加 [webhook](https://mastra.zisheng.pro/docs/long-running-agents/signal-providers) 入口点,并使用 [`handleWebhook()`](https://mastra.zisheng.pro/reference/signals/signal-provider)。 - 通过 [`getTools()`](https://mastra.zisheng.pro/reference/signals/signal-provider) 提供 `subscribe` 和 `unsubscribe` Tool,让 Agent 可以管理自己的订阅。 - 当通知需要去重或分批时,使用 [`dedupeKey` 和 `coalesceKey`](https://mastra.zisheng.pro/reference/agents/agent)。 了解更多: - [构建 Signal Provider](https://mastra.zisheng.pro/docs/long-running-agents/signal-providers) - [Signal](https://mastra.zisheng.pro/docs/long-running-agents/signals) - [`SignalProvider` 参考文档](https://mastra.zisheng.pro/reference/signals/signal-provider) - [`WebhookSignalProvider` 参考文档](https://mastra.zisheng.pro/reference/signals/webhook-signal-provider)