跳到主要内容

构建 Signal Provider

在本指南中,你将构建一个 Signal Provider,它会按固定时间间隔轮询外部服务,并在受监视的资源发生变化时向 Agent 线程推送通知。你将学习如何扩展 SignalProvider 基类、跟踪订阅、从轮询循环发出通知,以及在 Agent 上注册 Provider。

本示例会监视一个模拟 CI 服务中的构建流水线,但此模式适用于任何基于拉取的数据源:问题跟踪器、状态 API、队列或你自己的后端。

beta

Signal Provider 目前处于 beta 阶段。在 API 稳定之前,可能会出现不伴随主版本升级的破坏性变更。

前提条件
前提条件的直接链接

  • 已安装 Node.js v22.13.0 或更高版本
  • 受支持的 Model Provider 提供的 API 密钥
  • 现有的 Mastra 项目。如有需要,请按照安装指南操作。

本指南还假设你对 Signal 有基本了解。有关完整 API,请参阅 SignalProvider 参考文档

添加通知存储
添加通知存储的直接链接

Signal Provider 会将通知 Signal推送到线程,而通知需要支持通知域的存储 adapter。请在 Mastra 实例上配置存储。

src/mastra/index.ts
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。

src/mastra/signals/ci-client.ts
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(),并传入所有活动订阅。只对你关注的构建发出通知。

src/mastra/signals/ci-signals.ts
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 会连接它并自动启动轮询循环。

src/mastra/agents/dev-agent.ts
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,并添加第一步中的存储。

src/mastra/index.ts
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() 有资源可供检查。

src/mastra/subscribe.ts
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 run dev

确保某个线程已订阅,然后查看日志。模拟 client 每次轮询都会返回随机状态,因此在几个周期内,你就会看到构建失败为已订阅线程触发通知。

要查看 Agent 对通知的反应,请订阅该线程并进行流式传输:

src/mastra/watch-thread.ts
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() 提供 subscribeunsubscribe Tool,让 Agent 可以管理自己的订阅。
  • 当通知需要去重或分批时,使用 dedupeKeycoalesceKey

了解更多: