跳到主要内容

Signal Provider

加入版本: @mastra/core@1.39.0

beta

此功能目前处于 beta 阶段。在 API 稳定之前,可能会在不提升 major 版本的情况下发生破坏性变更。

Signal Provider 会监控 GitHub、Slack、持续集成(CI)或你自己的 API 等外部来源,并向已订阅的 Agent thread 推送通知 signal

何时使用 Signal Provider
何时使用 Signal Provider的直接链接

当外部系统产生 Agent 应响应的事件,并且希望由 Mastra 管理订阅记录时,请使用 Signal Provider。

  • 来源会发出与 thread 所关注资源相关的事件,例如 pull request、channel 或构建。
  • 需要集中跟踪哪些 thread 正在监控哪些外部资源。
  • 希望通过轮询、webhook 或二者接收事件。

如果只需要向 thread 推送一次性事件,请改为直接调用 agent.sendNotificationSignal()

Signal Provider 的工作原理
Signal Provider 的工作原理的直接链接

Signal Provider 是 signal 系统的生产端。它将外部事件传入 thread,signal API 则控制 thread 如何消费事件。

Signal Provider 组合了三种能力:

  • 订阅跟踪: SignalProvider 基类维护内存 Registry,将每个 Agent thread 映射到它所监控的外部资源。
  • 摄取: 对于拉取型来源覆盖 poll(),对于推送型来源覆盖 handleWebhook()
  • 发送: 当事件与订阅匹配时,调用 protected notify() helper,将通知 signal 转发到已连接 Agent 的 thread。

将 Provider 传给 Agent 即可注册。

设置 pollInterval 后,Agent 会连接 Provider 并开始轮询。它还会合并 Provider 暴露的所有 Processor 或 Tool。

src/mastra/agents/support-agent.ts
import { Agent } from '@mastra/core/agent'
import { CiSignals } from '../signals/ci-signals'

export const supportAgent = new Agent({
id: 'support-agent',
name: 'Support Agent',
instructions: 'Help the user triage updates.',
model: 'openai/gpt-5.6-sol',
signals: [new CiSignals()],
})
备注

发送通知需要支持通知的 Storage adapter,例如 libSQLPostgreSQLMongoDB。请在 Mastra 实例上配置 Storage,以便 notify() 存储通知记录。

快速入门
快速入门的直接链接

以下示例演示一个轮询 Provider:它监控 CI pipeline,并在已订阅的 pipeline 失败时发出通知。

src/mastra/signals/ci-signals.ts
import { SignalProvider } from '@mastra/core/signals'
import type { SignalProviderTarget, SignalSubscription } from '@mastra/core/signals'

type BuildStatus = {
id: string
status: 'passed' | 'failed'
}

const builds = new Map<string, BuildStatus>([
['acme-app-main', { id: 'build_123', status: 'failed' }],
])

async function fetchBuildStatus(pipeline: string): Promise<BuildStatus> {
return builds.get(pipeline) ?? { id: 'build_unknown', status: 'passed' }
}

export class CiSignals extends SignalProvider<'ci-signals'> {
readonly id = 'ci-signals' as const
readonly pollInterval = 30_000

watch(target: SignalProviderTarget, pipeline: string) {
return this.subscribe(target, pipeline)
}

unwatch(target: SignalProviderTarget, pipeline: string) {
return this.unsubscribe(target, pipeline)
}

async poll(subscriptions: SignalSubscription[]) {
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 },
)
}
}
}

向 Agent 注册 Provider,并让 thread 订阅要监控的 pipeline。

src/mastra/agents/support-agent.ts
import { Agent } from '@mastra/core/agent'
import { CiSignals } from '../signals/ci-signals'

export const ciSignals = new CiSignals()

export const supportAgent = new Agent({
id: 'support-agent',
name: 'Support Agent',
instructions: 'Help the user triage CI updates.',
model: 'openai/gpt-5.6-sol',
signals: [ciSignals],
})

ciSignals.watch({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme-app-main')

Mastra 会按 pollInterval 调用 poll(),并传入所有活跃订阅。没有订阅时会跳过一轮,并且不会让轮次重叠,因此运行较慢的 poll() 不会与自身并发执行。

备注

如需完整构建轮询 Provider,包括通知 Storage、Agent 注册、thread 订阅和测试,请参阅构建 Signal Provider

轮询和 Webhook Provider
轮询和 Webhook Provider的直接链接

当外部来源不会向应用推送事件时,请使用轮询。设置 pollInterval 并覆盖 poll(subscriptions)。每项订阅都包含 thread 目标和要检查的外部资源 ID。

当外部来源可以调用应用时,请使用 webhook。覆盖 handleWebhook(request),解析 payload,查找匹配的订阅,并针对每个匹配项调用 notify()

src/mastra/signals/ci-signals.ts
import { SignalProvider } from '@mastra/core/signals'
import type { SignalProviderWebhookRequest } from '@mastra/core/signals'

export class CiSignals extends SignalProvider<'ci-signals'> {
readonly id = 'ci-signals' as const

async handleWebhook(request: SignalProviderWebhookRequest) {
const payload = request.body as { pipeline: string; status: string }
const subscriptions = this.getSubscriptionsForResource(payload.pipeline)

for (const sub of subscriptions) {
await this.notify(
{
source: this.id,
kind: 'ci-status',
priority: 'high',
summary: `Build ${payload.status} for ${payload.pipeline}`,
payload,
},
{ resourceId: sub.resourceId, threadId: sub.threadId },
)
}

return { status: 200, body: { matched: subscriptions.length } }
}
}

handleWebhook() 是 Provider 方法,不是自动挂载的 HTTP 路由。请从自己的 endpoint 调用它,并传入请求 body、header 和任何路由参数。有关订阅、轮询、生命周期和 notify() 的详情,请参阅 SignalProvider Reference。有关完整通知 payload 结构(包括去重和合并字段),请参阅 Agent.sendNotificationSignal() Reference

内置 Webhook Provider
内置 Webhook Provider的直接链接

对于通用 webhook 来源,请使用 WebhookSignalProvider,而不必编写子类。使用一个从 payload 中提取资源 ID 的函数进行配置,并可选择提供构建通知的函数。

src/mastra/index.ts
import { Agent } from '@mastra/core/agent'
import { WebhookSignalProvider } from '@mastra/core/signals'

const webhooks = new WebhookSignalProvider({
extractResourceId: payload => (payload as { repository: string }).repository,
buildNotification: (payload, sub) => ({
source: 'ci',
kind: 'build-status',
priority: 'medium',
summary: `Build ${(payload as { status: string }).status} for ${sub.externalResourceId}`,
}),
})

export const supportAgent = new Agent({
id: 'support-agent',
name: 'Support Agent',
instructions: 'Help the user triage updates.',
model: 'openai/gpt-5.6-sol',
signals: [webhooks],
})

webhooks.subscribeThread({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme/app')

Webhook 到达时,请从路由调用 webhooks.handleWebhook({ body, headers })。Provider 会将提取的资源 ID 与其订阅进行匹配,并通知每个匹配的 thread。

高级 Provider 能力
高级 Provider 能力的直接链接

Provider 可以支持事件摄取以外的能力。请只添加来源需要的能力。

  • 持久订阅: 基础 Registry 位于内存中且限定于单个进程。如果订阅必须在重启后继续存在,请自行持久化,然后在 start() 中恢复。
  • 生命周期 hook: 覆盖 start() 进行异步设置,并覆盖 stop() 进行清理。覆盖 stop() 时请调用 super.stop(),以便基础 Provider 停止轮询并清空 Registry。
  • Processor 和 Tool:getInputProcessors()getOutputProcessors() 返回 Processor,并从 getTools() 返回 Agent 可调用的 Tool。

@mastra/github-signals 包是生产级 Signal Provider,用于监控 GitHub pull request,并将评论、审查状态、持续集成状态和合并通知到 thread。可以将它作为轮询、持久订阅、Tool、Processor 和生命周期 hook 的参考。

src/mastra/agents/dev-agent.ts
import { Agent } from '@mastra/core/agent'
import { GithubSignals } from '@mastra/github-signals'

export const devAgent = new Agent({
id: 'dev-agent',
name: 'Dev Agent',
instructions: 'Help triage pull request activity.',
model: 'openai/gpt-5.6-sol',
signals: [new GithubSignals()],
})