跳至主要內容

SignalProvider

新增於: @mastra/core@1.39.0

用於建立 signal Provider 的抽象基礎類別。Signal Provider 會監察外部來源(API、webhook、事件串流),並透過內置訂閱 registry 將通知 signal 推送至 Agent thread。

Signal Provider 預設並非 processor。需要攔截 Agent 執行的 Provider 會從 getInputProcessors()getOutputProcessors() 傳回 processor。提供可由 Agent 呼叫之 Tool 的 Provider,則會從 getTools() 傳回這些 Tool。

如需即用的 webhook 型 Provider,請參閱 WebhookSignalProvider

使用範例
使用範例 的直接連結

每 30 秒檢查一次 API 的輪詢 Provider:

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

class SlackSignals extends SignalProvider<'slack-signals'> {
readonly id = 'slack-signals'
readonly pollInterval = 30_000

async poll(subscriptions: SignalSubscription[]) {
for (const sub of subscriptions) {
const messages = await fetchSlackMessages(sub.externalResourceId)
if (messages.length > 0) {
await this.notify(
{
source: 'slack',
kind: 'new-messages',
summary: `${messages.length} new messages in ${sub.externalResourceId}`,
},
{ threadId: sub.threadId, resourceId: sub.resourceId },
)
}
}
}
}

向 Agent 註冊:

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

const agent = new Agent({
id: 'agent',
signals: [new SlackSignals()],
})

Agent 會呼叫 connect(this),並註冊 Provider 傳回的任何 processor 或 Tool,然後開始輪詢。

建構函數參數
建構函數參數 的直接連結

SignalProvider 是抽象類別。子類別呼叫 super() 時無需傳入參數。

屬性
屬性 的直接連結

id:

TId extends string
此 Provider 的唯一識別碼。子類別必須將其實作為 readonly 屬性。

name?:

string
Provider 便於閱讀的顯示名稱。

pollInterval?:

number
輪詢間隔(毫秒)。設定後,框架會按此間隔呼叫 poll()。僅使用 webhook 的 Provider 請保留為 undefined0

isConnected:

boolean
此 Provider 是否已連接至 Agent。呼叫 connect() 後傳回 true。內部用於在 Agent.__fork() 期間略過重新連接。

方法
方法 的直接連結

連接
連接 的直接連結

connect(agent)
connectagent 的直接連結

由 Agent 建構函數呼叫。建立雙向連結,讓 Provider 可將 signal 傳回 Agent。如要在建立連結後執行額外設定,請覆寫此方法。必須一律呼叫 super.connect(agent)

class MySignals extends SignalProvider<'my-signals'> {
readonly id = 'my-signals'

override connect(agent) {
super.connect(agent)
// additional setup after agent link is established
}
}

__registerMastra(mastra)
__registermastramastra 的直接連結

當 Provider 的 Agent 向 Mastra 實例註冊時呼叫。如要存取儲存空間或其他 Mastra 服務,請覆寫此方法。必須一律呼叫 super.__registerMastra(mastra)

override __registerMastra(mastra) {
super.__registerMastra(mastra)
// this.mastra is now available
}

Processor 和 Tool 整合
Processor 和 Tool 整合 的直接連結

getInputProcessors()
getinputprocessors 的直接連結

傳回此 Provider 需要向 Agent 註冊的 input processor。當 Provider 會攔截 Agent 輸入步驟時覆寫此方法(例如注入情境提示或偵測 Tool 呼叫)。

getInputProcessors() {
return [this]
}

傳回:InputProcessorOrWorkflow[]

getOutputProcessors()
getoutputprocessors 的直接連結

傳回此 Provider 需要向 Agent 註冊的 output processor。當 Provider 會攔截 Agent 輸出步驟時覆寫此方法。

getOutputProcessors() {
return [this]
}

傳回:OutputProcessorOrWorkflow[]

getTools()
gettools 的直接連結

傳回此 Provider 向 Agent 提供的 Tool。當 Provider 加入可由 Agent 呼叫的 Tool(例如訂閱或取消訂閱命令)時覆寫此方法。

getTools() {
return {
subscribe_pr: createTool({ /* ... */ }),
unsubscribe_pr: createTool({ /* ... */ }),
}
}

傳回:Record<string, unknown>

訂閱追蹤
訂閱追蹤 的直接連結

subscribe(target, externalResourceId, metadata?)
subscribetarget-externalresourceid-metadata 的直接連結

讓 thread 訂閱外部資源。這是 protected 方法:請從 Provider 實作內部呼叫。

const sub = this.subscribe(
{ threadId: 'thread-1', resourceId: 'user-1' },
'github:mastra-ai/mastra#123',
{ pr: 123 },
)

傳回:SignalSubscription:新建立的訂閱;如訂閱已存在,則傳回已合併 metadata 的現有訂閱。

target:

SignalProviderTarget
要訂閱的 thread。必須包含 threadIdresourceId

externalResourceId:

string
Provider 專用的外部資源識別碼(例如 "github:owner/repo#123")。

metadata?:

Record<string, unknown>
與訂閱一併儲存的額外資料。重複訂閱時,會合併至現有 metadata。

unsubscribe(target, externalResourceId)
unsubscribetarget-externalresourceid 的直接連結

移除訂閱。

const removed = this.unsubscribe(
{ threadId: 'thread-1', resourceId: 'user-1' },
'github:mastra-ai/mastra#123',
)

傳回:boolean:如已移除則為 true;如沒有相符的訂閱則為 false

getSubscriptions()
getsubscriptions 的直接連結

傳回此 Provider 的所有有效訂閱。

const allSubs = this.getSubscriptions()

傳回:SignalSubscription[]

getSubscriptionsForResource(externalResourceId)
getsubscriptionsforresourceexternalresourceid 的直接連結

傳回特定外部資源的所有訂閱。

const subs = this.getSubscriptionsForResource('github:mastra-ai/mastra#123')
for (const sub of subs) {
await this.notify(
{ source: 'my-provider', kind: 'update', summary: 'Resource updated' },
{ threadId: sub.threadId, resourceId: sub.resourceId },
)
}

傳回:SignalSubscription[]

getSubscriptionsForThread(target)
getsubscriptionsforthreadtarget 的直接連結

傳回特定 thread 的所有訂閱。

const subs = this.getSubscriptionsForThread({
threadId: 'thread-1',
resourceId: 'user-1',
})

傳回:SignalSubscription[]

hasSubscription(target, externalResourceId)
hassubscriptiontarget-externalresourceid 的直接連結

檢查訂閱是否存在。

if (this.hasSubscription(target, 'github:mastra-ai/mastra#123')) {
// already subscribed
}

傳回:boolean

unsubscribeAll(target)
unsubscribealltarget 的直接連結

移除 thread 的所有訂閱。

const removed = this.unsubscribeAll({
threadId: 'thread-1',
resourceId: 'user-1',
})

傳回:number:已移除的訂閱數目。

subscriptionCount
subscriptioncount 的直接連結

此 Provider 的有效訂閱總數。

if (this.subscriptionCount === 0) {
// nothing to poll
}

傳回:number

輪詢
輪詢 的直接連結

poll(subscriptions)
pollsubscriptions 的直接連結

每個輪詢週期均會呼叫此方法,並傳入所有有效訂閱。覆寫此方法以檢查外部來源並發出通知。框架會防止輪詢週期重疊:如 poll() 呼叫所需時間超過 pollInterval,便會略過下一個週期。

async poll(subscriptions: SignalSubscription[]) {
for (const sub of subscriptions) {
const events = await checkExternalSource(sub.externalResourceId)
for (const event of events) {
await this.notify(
{ source: 'my-provider', kind: event.type, summary: event.message },
{ threadId: sub.threadId, resourceId: sub.resourceId },
)
}
}
}

startPolling()
startpolling 的直接連結

啟動輪詢 timer。Agent 會在 connect() 後呼叫此方法。此方法具冪等性:呼叫多次不會產生額外效果。

provider.startPolling()

stopPolling()
stoppolling 的直接連結

停止輪詢 timer。

provider.stopPolling()

Webhook
Webhook 的直接連結

handleWebhook(request)
handlewebhookrequest 的直接連結

處理傳入的 webhook request。覆寫此方法以解析 payload 並配對訂閱,然後發出通知 signal。如需即用的實作,請參閱 WebhookSignalProvider

驗證 webhook request 後,請從應用程式定義的 HTTP endpoint 呼叫此方法。

async handleWebhook(request) {
const payload = request.body as { repo: string, event: string }
const subs = this.getSubscriptionsForResource(payload.repo)

for (const sub of subs) {
await this.notify(
{ source: 'github', kind: payload.event, summary: `Event on ${payload.repo}` },
{ threadId: sub.threadId, resourceId: sub.resourceId },
)
}

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

傳回:Promise<{ status?: number; body?: unknown }>

生命週期
生命週期 的直接連結

start()
start 的直接連結

connect() 後呼叫以執行非同步初始化。如設定需要 Agent 或 Mastra 實例可供使用,請覆寫此方法。

async start() {
await this.loadInitialState()
}

stop()
stop 的直接連結

關閉時呼叫。預設實作會停止輪詢並清除所有訂閱。

provider.stop()

通知
通知 的直接連結

notify(notification, target)
notifynotification-target 的直接連結

向已連接的 Agent 傳送通知 signal。這是包裝 agent.sendNotificationSignal() 的 protected 便捷方法。

await this.notify(
{
source: 'my-provider',
kind: 'pr-updated',
summary: 'PR #123 was updated',
priority: 'high',
payload: { prNumber: 123 },
},
{ threadId: 'thread-1', resourceId: 'user-1' },
)

notification:

object
通知 payload。
object

source:

string
通知來源的識別碼。

kind:

string
事件類型(例如 "pr-updated""new-message")。

summary:

string
便於閱讀的事件摘要。

priority?:

"high" | "medium" | "low"
通知優先級。

payload?:

unknown
附加至通知的任意資料。

target:

SignalProviderTarget
要通知的 thread。必須包含 threadIdresourceId

類型
類型 的直接連結

SignalSubscription
signalsubscription 的直接連結

subscribe() 傳回的訂閱物件。

id:

string
訂閱的唯一識別碼。

providerId:

string
擁有此訂閱的 Provider。

threadId:

string
接收 signal 的 thread。

resourceId:

string
擁有該 thread 的資源。

externalResourceId:

string
Provider 專用的外部資源識別碼(例如 "github:owner/repo#123")。

subscribedAt:

Date
建立訂閱的時間。

metadata:

Record<string, unknown>
與訂閱一併儲存的 Provider 專用 metadata。

SignalProviderTarget
signalprovidertarget 的直接連結

識別特定 Agent thread。

threadId:

string
目標 thread。

resourceId:

string
擁有該 thread 的資源。

agentId?:

string
Agent 識別碼。

Type guard
Type guard 的直接連結

isSignalProvider(obj)
issignalproviderobj 的直接連結

SignalProvider 實例進行 runtime 檢查。

import { isSignalProvider } from '@mastra/core/signals'

if (isSignalProvider(obj)) {
obj.connect(agent)
}

傳回:boolean