构建 Signal Provider
在本指南中,你将构建一个 Signal Provider,它会按固定时间间隔轮询外部服务,并在受监视的资源发生变化时向 Agent 线程推送通知。你将学习如何扩展 SignalProvider 基类、跟踪订阅、从轮询循环发出通知,以及在 Agent 上注册 Provider。
本示例会监视一个模拟 CI 服务中的构建流水线,但此模式适用于任何基于拉取的数据源:问题跟踪器、状态 API、队列或你自己的后端。
Signal Provider 目前处于 beta 阶段。在 API 稳定之前,可能会出现不伴随主版本升级的破坏性变更。
前提条件前提条件的直接链接
- 已安装 Node.js
v22.13.0或更高版本 - 受支持的 Model Provider 提供的 API 密钥
- 现有的 Mastra 项目。如有需要,请按照安装指南操作。
本指南还假设你对 Signal 有基本了解。有关完整 API,请参阅 SignalProvider 参考文档。
添加通知存储添加通知存储的直接链接
Signal Provider 会将通知 Signal推送到线程,而通知需要支持通知域的存储 adapter。请在 Mastra 实例上配置存储。
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创建外部服务 client的直接链接
真实 Provider 会调用外部 API。为了让本指南可独立运行,请创建一个小型模拟 CI client,用于返回流水线的构建状态。之后再将它替换成真实的 API client。
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<BuildStatus> {
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构建 Signal Provider的直接链接
扩展 SignalProvider,实现抽象的 id 字段、设置 pollInterval,并覆盖 poll()。基类会按设定间隔调用 poll(),并传入所有活动订阅。只对你关注的构建发出通知。
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<void> {
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在 Agent 上注册 Provider的直接链接
通过 signals 将 Provider 传给 Agent。Agent 会连接它并自动启动轮询循环。
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,并添加第一步中的存储。
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() 有资源可供检查。
import { ciSignals } from './agents/dev-agent'
ciSignals.watch({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme/app:main')
在 Agent 注册后运行一次这段代码,例如通过设置脚本或 API 路由。订阅位于 Provider 的内存注册表中,因此重启后需要重新订阅。
测试 Signal Provider测试 Signal Provider的直接链接
启动开发服务器:
- npm
- pnpm
- Yarn
- Bun
npm run dev
pnpm run dev
yarn dev
bun run dev
确保某个线程已订阅,然后查看日志。模拟 client 每次轮询都会返回随机状态,因此在几个周期内,你就会看到构建失败为已订阅线程触发通知。
要查看 Agent 对通知的反应,请订阅该线程并进行流式传输:
const subscription = await devAgent.subscribeToThread({
resourceId: 'user_123',
threadId: 'thread_456',
})
for await (const chunk of subscription.stream) {
console.log(chunk)
}
构建失败时,模型会把通知作为上下文接收:
<notification source="ci-signals" type="ci-status" priority="high" status="delivered">Build failed for acme/app:main</notification>
由于模拟 client 会随机生成状态,且模型可以自由组织回复,因此输出并不确定,确切措辞也会有所不同。
后续步骤后续步骤的直接链接
你可以扩展此 Signal Provider:
- 将
fetchBuildStatus()替换为真实 API client。 - 持久化订阅,使其在重启后仍然存在,然后在
start()中恢复它们。 - 为基于推送的数据源添加 webhook 入口点,并使用
handleWebhook()。 - 通过
getTools()提供subscribe和unsubscribeTool,让 Agent 可以管理自己的订阅。 - 当通知需要去重或分批时,使用
dedupeKey和coalesceKey。
了解更多: