Signal provider
新增於: @mastra/core@1.39.0
此功能目前為 beta 版本。在 API 穩定之前,即使主版本號沒有提升,也可能會有破壞性變更。
Signal provider 會監察 GitHub、Slack、持續整合(CI)或你自己的 API 等外部來源,並將通知訊號推送至已訂閱的 Agent 對話串。
何時使用 Signal provider何時使用 Signal provider 的直接連結
當外部系統產生 Agent 應該回應的事件,而你希望由 Mastra 管理訂閱記錄時,便可使用 Signal provider。
- 來源會發出與對話串關注資源相關的事件,例如 pull request、頻道或 build。
- 你希望集中追蹤哪些對話串正在監察哪些外部資源。
- 你希望透過輪詢或 webhook 接收事件,亦可同時使用兩者。
如果只需向對話串推送一次性事件,請直接呼叫 agent.sendNotificationSignal()。
Signal provider 的運作方式Signal provider 的運作方式 的直接連結
Signal provider 是訊號系統的產生端。它把外部事件帶入對話串,而訊號 API 則控制對話串如何取用這些事件。
Signal provider 結合三項能力:
- 訂閱追蹤:
SignalProviderbase class 會保存記憶體內 registry,將每個 Agent 對話串對應至其監察的外部資源。 - 擷取: pull-based 來源可覆寫
poll(),push-based 來源則可覆寫handleWebhook()。 - 傳送: 當事件符合訂閱時,呼叫 protected
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()],
})
通知傳送功能需要支援通知的儲存空間 adapter,例如 libSQL、PostgreSQL 或 MongoDB。請在 Mastra instance 上設定儲存空間,讓 notify() 可以儲存通知記錄。
快速開始快速開始 的直接連結
以下例子示範一個輪詢 provider,它會監察 CI pipeline,並在已訂閱的 pipeline 失敗時發出通知。
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 建立教學,請參閱建立 Signal provider。
輪詢及 webhook provider輪詢及 webhook provider 的直接連結
如果外部來源不會向你的應用程式推送事件,請使用輪詢。設定 pollInterval 並覆寫 poll(subscriptions);每項訂閱都包含要檢查的對話串 target 及外部資源 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 路由。請從你自己的 endpoint 呼叫它,並傳入 request body、header 及所有路由參數。訂閱、輪詢、生命週期及 notify() 的詳情請參閱 SignalProvider 參考資料。完整通知 payload 格式(包括去重及合併欄位)請參閱 Agent.sendNotificationSignal() 參考資料。
內置 webhook provider內置 webhook provider 的直接連結
如要處理通用 webhook 來源,請使用 WebhookSignalProvider,毋須自行編寫 subclass。設定時請提供一個從 payload 擷取資源 ID 的 function,亦可選擇提供建立通知的 function。
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 除了擷取事件,亦可支援其他功能。只需加入來源所需的功能。
- 持久訂閱: base registry 位於記憶體內,並按 process 分開。需要讓訂閱在重新啟動後仍然保留時,請自行持久儲存訂閱,然後在
start()中重新載入。 - 生命週期 hook: 覆寫
start()以進行 async 設定,並覆寫stop()以執行清理。覆寫stop()時請呼叫super.stop(),讓 base provider 可以停止輪詢並清除 registry。 - Processor 及 Tool: 從
getInputProcessors()或getOutputProcessors()傳回 processor,並從getTools()傳回 Agent 可呼叫的 Tool。
@mastra/github-signals 套件是可供生產環境使用的 Signal provider,會監察 GitHub pull request,並將留言、審查狀態、持續整合狀態及合併情況通知對話串。你可以參考它如何處理輪詢、持久訂閱、Tool、processor 及生命週期 hook。
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()],
})