跳至主要內容

建構 signal provider

在本指南中,你將建構一個 signal provider,定期輪詢外部服務,並在受監看的資源發生變更時,將通知推送至 Agent thread。你將學會如何擴充 SignalProvider 基底類別、追蹤訂閱、從輪詢迴圈發出通知,以及在 Agent 上註冊 provider。

範例會監看模擬 CI 服務中的建置 pipeline,但這個模式適用於任何拉取型來源,例如 issue tracker、狀態 API、queue 或你自己的後端。

beta

Signal provider 目前為 beta 版。在 API 穩定之前,即使未提高 major version,也可能會有 breaking changes。

先決條件
「先決條件」的直接連結

  • 已安裝 Node.js v22.13.0 或更新版本
  • 具備受支援 Model Provider 的 API key
  • 已有 Mastra 專案。如有需要,請依照安裝指南操作。

本指南也假設你對 signal 有概略了解。如需完整 API 功能,請參閱 SignalProvider 參考文件

新增通知儲存空間
「新增通知儲存空間」的直接連結

Signal provider 會將通知 signal 推送至 thread,而通知需要支援 notifications domain 的 storage adapter。請在 Mastra instance 上設定 storage。

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 都支援通知紀錄。如果沒有 notification storage,provider 的 notify() 呼叫會在 runtime 拋出錯誤。

建立外部服務 client
「建立外部服務 client」的直接連結

實際的 provider 會呼叫外部 API。為了讓本指南可以獨立完成,請建立一個小型的模擬 CI client,用來傳回 pipeline 的建置狀態。之後再將它替換成實際使用的 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 method(watch / unwatch)包裝它們,讓呼叫端可以管理訂閱。
  • externalResourceId 可以是任何 provider 專用字串。此處使用 pipeline 名稱;GitHub provider 則可能使用 "github:owner/repo#123"
  • notify() 會將通知 signal 轉送至已連線 Agent 的 thread。如果 provider 從未在 Agent 上註冊,就會拋出錯誤。
  • dedupeKey 可避免重複儲存同一筆失敗紀錄。

在 Agent 上註冊 provider
「在 Agent 上註冊 provider」的直接連結

透過 signals 將 provider 傳給 Agent。Agent 會自動連接 provider 並啟動輪詢迴圈。

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,並加入第一個步驟中的 storage。

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',
}),
})

訂閱 thread
「訂閱 thread」的直接連結

Provider 只會輪詢 thread 正在監看的資源。讓 thread 訂閱 pipeline,poll() 才有內容可檢查。

src/mastra/subscribe.ts
import { ciSignals } from './agents/dev-agent'

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

請在 Agent 註冊後執行一次,例如透過設定 script 或 API route 執行。訂閱會存放在 provider 的記憶體內 registry,因此重新啟動後需要再次訂閱。

測試 signal provider
「測試 signal provider」的直接連結

啟動開發伺服器:

npm run dev

請確認已有 thread 訂閱,然後查看 log。模擬 client 每次輪詢都會傳回隨機狀態,因此在幾輪之內,你就會看到建置失敗觸發一則傳給已訂閱 thread 的通知。

若要查看 Agent 如何回應通知,請訂閱 thread 並以串流方式接收內容:

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)
}

建置失敗時,model 會收到作為 context 的通知:

<notification source="ci-signals" type="ci-status" priority="high" status="delivered">Build failed for acme/app:main</notification>

由於模擬 client 會隨機產生狀態,而 model 也會自由組織回覆內容,因此輸出並非 deterministic,實際措辭會有所不同。

後續步驟
「後續步驟」的直接連結

你可以進一步擴充這個 signal provider:

  • fetchBuildStatus() 替換成實際的 API client。
  • 持久化訂閱,讓訂閱能在重新啟動後保留,接著在 start() 中重新載入。
  • 為推送型來源新增 webhook 進入點,並使用 handleWebhook()
  • 公開 subscribeunsubscribe tools,並使用 getTools(),讓 Agent 能管理自己的訂閱。
  • 當通知需要去重或批次處理時,使用 dedupeKeycoalesceKey

深入了解: