訊號 Provider
新增於: @mastra/core@1.39.0
此功能目前為 beta。在 API 穩定前,即使未提高主要版本,也可能發生破壞性變更。
訊號 Provider 會監控 GitHub、Slack、持續整合(CI)或你自己的 API 等外部來源,並將通知訊號推送至已訂閱的 Agent 對話串。
何時使用訊號 Provider「何時使用訊號 Provider」的直接連結
當外部系統產生 Agent 應回應的事件,而且你希望由 Mastra 管理訂閱記錄時,請使用訊號 Provider。
- 來源發出的事件與對話串關注的資源相關,例如 Pull Request、頻道或建置。
- 你希望集中追蹤哪些對話串正在監看哪些外部資源。
- 你希望透過輪詢、Webhook 或兩者接收事件。
若只需將一次性事件推送至對話串,請改為直接呼叫 agent.sendNotificationSignal()。
訊號 Provider 的運作方式「訊號 Provider 的運作方式」的直接連結
訊號 Provider 是訊號系統的產生端。它將外部事件帶入對話串,而訊號 API 則控制對話串如何取用這些事件。
訊號 Provider 結合三種能力:
- 訂閱追蹤:
SignalProvider基底類別維護記憶體內登錄,將每個 Agent 對話串對應至其監看的外部資源。 - 擷取: 針對提取式來源覆寫
poll(),或針對推送式來源覆寫handleWebhook()。 - 傳送: 事件符合訂閱時,呼叫受保護的
notify()Helper,將通知訊號轉送至已連線 Agent 的對話串。
將 Provider 傳入 Agent 即可註冊。
若設有 pollInterval,Agent 會連接 Provider 並開始輪詢,也會合併 Provider 公開的任何 Processor 或 Tool。
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()],
})
通知傳送需要支援通知的儲存空間配接器,例如 libSQL、PostgreSQL 或 MongoDB。請在 Mastra 執行個體上設定儲存空間,讓 notify() 能儲存通知記錄。
快速開始「快速開始」的直接連結
下列範例示範監看 CI Pipeline,並在已訂閱的 Pipeline 失敗時發出通知的輪詢 Provider。
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,並讓對話串訂閱你要監看的 Pipeline。
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() 不會與自身並行執行。
如需包含通知儲存、Agent 註冊、對話串訂閱及測試的完整輪詢 Provider 建構方式,請參閱建構訊號 Provider。
輪詢與 Webhook Provider「輪詢與 Webhook Provider」的直接連結
外部來源不會將事件推送至應用程式時,請使用輪詢。設定 pollInterval 並覆寫 poll(subscriptions)。每筆訂閱都包含要檢查的對話串目標及外部資源 ID。
外部來源能呼叫應用程式時,請使用 Webhook。覆寫 handleWebhook(request)、剖析 payload、尋找相符訂閱,再針對每筆相符項目呼叫 notify()。
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 路由。請從你自己的端點叫用它,並傳入請求內文、標頭及任何路由參數。訂閱、輪詢、生命週期及 notify() 的詳細資訊請參閱 SignalProvider 參考。完整通知 payload 結構(包括去重複及合併欄位)請參閱 Agent.sendNotificationSignal() 參考。
內建 Webhook Provider「內建 Webhook Provider」的直接連結
針對一般 Webhook 來源,請使用 WebhookSignalProvider,不必自行撰寫子類別。請使用從 payload 擷取資源 ID 的函式進行設定,並可選擇提供建立通知的函式。
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 與訂閱比對,並通知每個相符的對話串。
進階 Provider 能力「進階 Provider 能力」的直接連結
Provider 除了事件擷取外,也能支援更多功能。請只加入來源所需的能力。
- 持久型訂閱: 基底登錄位於記憶體內,且各處理程序彼此獨立。訂閱必須承受重新啟動時,請自行持久化,然後在
start()中重新載入。 - 生命週期掛鉤: 覆寫
start()進行非同步設定,並覆寫stop()進行清理。覆寫stop()時請呼叫super.stop(),讓基底 Provider 能停止輪詢並清除其登錄。 - Processor 與 Tool: 從
getInputProcessors()或getOutputProcessors()傳回 Processor,並從getTools()傳回可由 Agent 呼叫的 Tool。
@mastra/github-signals 套件是監看 GitHub Pull Request,並將留言、審查狀態、持續整合狀態與合併通知對話串的正式環境訊號 Provider。你可以將它當作輪詢、持久型訂閱、Tool、Processor 及生命週期掛鉤的參考。
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()],
})