建立訊號 Provider
在本指南中,你會建立一個定期輪詢外部服務的訊號 Provider,並在受監察資源變更時將通知推送至 Agent 線程。你將學習如何擴充 SignalProvider 基礎類別、追蹤訂閱、從輪詢迴圈發出通知,以及在 Agent 上註冊 Provider。
此範例會監察模擬 CI 服務中的建置管線,但此模式適用於任何拉取型來源:問題追蹤器、狀態 API、佇列或你自己的後端。
訊號 Provider 目前處於 beta 階段。在 API 穩定前,即使主要版本沒有提升,也可能出現破壞性變更。
先決條件先決條件 的直接連結
- 已安裝 Node.js
v22.13.0或更新版本 - 受支援 Model Provider 的 API 金鑰
- 現有的 Mastra 項目。如有需要,請按照安裝指南操作。
本指南亦假設你對訊號有概括了解。如要查看完整 API 介面,請參閱 SignalProvider 參考。
加入通知儲存加入通知儲存 的直接連結
訊號 Provider 會將通知訊號推送至線程,而通知需要支援通知領域的儲存配接器。請在 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() 呼叫會在執行階段擲出錯誤。
建立外部服務用戶端建立外部服務用戶端 的直接連結
實際的 Provider 會呼叫外部 API。為使本指南自成一體,請建立一個小型模擬 CI 用戶端,以傳回管線的建置狀態。你可稍後將其換成實際 API 用戶端。
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 }
}
此用戶端每次呼叫都會傳回隨機狀態。接駁實際 API 時,只需變更此檔案。
建立訊號 Provider建立訊號 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 方法。請以你自己的公開方法(watch/unwatch)包裝它們,讓呼叫者可以管理訂閱。externalResourceId可以是任何 Provider 專用字串。此處是管線名稱;GitHub Provider 可能會使用"github:owner/repo#123"。notify()會將通知訊號轉送至已連接 Agent 的線程。如果 Provider 從未在 Agent 上註冊,便會擲出錯誤。dedupeKey可避免重複儲存相同失敗事件。
在 Agent 上註冊 Provider在 Agent 上註冊 Provider 的直接連結
透過 signals 將 Provider 傳給 Agent。Agent 會連接 Provider 並自動啟動輪詢迴圈。
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 的記憶體登錄中,因此重新啟動後需要再次訂閱。
測試訊號 Provider測試訊號 Provider 的直接連結
啟動開發伺服器:
- npm
- pnpm
- Yarn
- Bun
npm run dev
pnpm run dev
yarn dev
bun run dev
確保線程已訂閱,然後觀察記錄。模擬用戶端每次輪詢都會傳回隨機狀態,因此在幾個週期內,你會看到建置失敗為已訂閱線程觸發通知。
如要查看 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>
由於模擬用戶端會隨機產生狀態,而模型亦可自由措辭,因此輸出並非確定,確切字句會有所不同。
後續步驟後續步驟 的直接連結
你可以擴充此訊號 Provider,以:
- 以實際 API 用戶端取代
fetchBuildStatus()。 - 持久保存訂閱,使其在重新啟動後仍然存在,然後在
start()中重新載入。 - 為推送型來源加入 webhook 進入點,並使用
handleWebhook()處理請求。 - 使用
getTools()公開subscribe及unsubscribeTool,讓 Agent 管理自己的訂閱。 - 通知需要去重或批次處理時,使用
dedupeKey及coalesceKey。
了解更多: