SignalProvider
追加バージョン: @mastra/core@1.39.0
Signal Provider を構築するための抽象基底クラスです。Signal Provider は外部 Source(API、Webhook、Event Stream)を監視し、組み込みの Subscription Registry を通じて Agent Thread に Notification Signal を送出します。
Signal Provider はデフォルトでは Processor ではありません。Agent の実行をインターセプトする必要がある Provider は、getInputProcessors() または getOutputProcessors() から Processor を返します。Agent から呼び出せる Tool を公開する Provider は、getTools() から Tool を返します。
すぐに使用できる Webhook ベースの Provider については、WebhookSignalProvider を参照してください。
使用例使用例への直接リンク
30 秒ごとに API を確認する Polling Provider の例です。
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 に登録します。
import { Agent } from '@mastra/core/agent'
const agent = new Agent({
id: 'agent',
signals: [new SlackSignals()],
})
Agent は connect(this) を呼び出し、Provider が返すすべての Processor または Tool を登録します。その後、Polling を開始します。
コンストラクターのパラメーターコンストラクターのパラメーターへの直接リンク
SignalProvider は抽象クラスです。サブクラスは引数なしで super() を呼び出します。
プロパティプロパティへの直接リンク
id:
name?:
pollInterval?:
poll() を呼び出します。Webhook 専用 Provider では undefined または 0 のままにします。isConnected:
connect() の呼び出し後は true を返します。Agent.__fork() 中の再接続をスキップするために内部で使用されます。メソッドメソッドへの直接リンク
接続接続への直接リンク
connect(agent)connectagentへの直接リンク
Agent のコンストラクターから呼び出されます。Provider が Agent に Signal を送り返せるよう、双方向リンクを確立します。リンクの確立後に追加のセットアップを実行するにはオーバーライドします。必ず 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 Service にアクセスするにはオーバーライドします。必ず 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 の入力 Step をインターセプトする場合(たとえば、Context Hint の注入や Tool 呼び出しの検出)にオーバーライドします。
getInputProcessors() {
return [this]
}
戻り値: InputProcessorOrWorkflow[]
getOutputProcessors()getoutputprocessorsへの直接リンク
この Provider が Agent に登録する必要のある Output Processor を返します。Provider が Agent の出力 Step をインターセプトする場合にオーバーライドします。
getOutputProcessors() {
return [this]
}
戻り値: OutputProcessorOrWorkflow[]
getTools()gettoolsへの直接リンク
この Provider が Agent に公開する Tool を返します。購読や購読解除コマンドなど、Agent から呼び出せる Tool を Provider が追加する場合にオーバーライドします。
getTools() {
return {
subscribe_pr: createTool({ /* ... */ }),
unsubscribe_pr: createTool({ /* ... */ }),
}
}
戻り値: Record<string, unknown>
Subscription の追跡Subscription の追跡への直接リンク
subscribe(target, externalResourceId, metadata?)subscribetarget-externalresourceid-metadataへの直接リンク
Thread で外部 Resource を購読します。これは protected メソッドです。Provider の実装内から呼び出してください。
const sub = this.subscribe(
{ threadId: 'thread-1', resourceId: 'user-1' },
'github:mastra-ai/mastra#123',
{ pr: 123 },
)
戻り値: SignalSubscription。作成された Subscription、または Metadata が統合された既存の Subscription です。
target:
threadId と resourceId を含める必要があります。externalResourceId:
"github:owner/repo#123")。metadata?:
unsubscribe(target, externalResourceId)unsubscribetarget-externalresourceidへの直接リンク
Subscription を削除します。
const removed = this.unsubscribe(
{ threadId: 'thread-1', resourceId: 'user-1' },
'github:mastra-ai/mastra#123',
)
戻り値: boolean。削除した場合は true、一致する Subscription が存在しなかった場合は false です。
getSubscriptions()getsubscriptionsへの直接リンク
この Provider のアクティブな Subscription をすべて返します。
const allSubs = this.getSubscriptions()
戻り値: SignalSubscription[]
getSubscriptionsForResource(externalResourceId)getsubscriptionsforresourceexternalresourceidへの直接リンク
特定の外部 Resource に対するすべての Subscription を返します。
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 に対するすべての Subscription を返します。
const subs = this.getSubscriptionsForThread({
threadId: 'thread-1',
resourceId: 'user-1',
})
戻り値: SignalSubscription[]
hasSubscription(target, externalResourceId)hassubscriptiontarget-externalresourceidへの直接リンク
Subscription が存在するか確認します。
if (this.hasSubscription(target, 'github:mastra-ai/mastra#123')) {
// already subscribed
}
戻り値: boolean
unsubscribeAll(target)unsubscribealltargetへの直接リンク
Thread のすべての Subscription を削除します。
const removed = this.unsubscribeAll({
threadId: 'thread-1',
resourceId: 'user-1',
})
戻り値: number。削除された Subscription の数です。
subscriptionCountsubscriptioncountへの直接リンク
この Provider のアクティブな Subscription の総数です。
if (this.subscriptionCount === 0) {
// nothing to poll
}
戻り値: number
PollingPollingへの直接リンク
poll(subscriptions)pollsubscriptionsへの直接リンク
Polling Cycle ごとに、すべてのアクティブな Subscription とともに呼び出されます。外部 Source を確認して通知を送出するにはオーバーライドします。Framework は Polling Cycle の重複を防ぎます。poll() の呼び出しに pollInterval より長い時間がかかった場合、次の Cycle はスキップされます。
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への直接リンク
Polling Timer を開始します。connect() の後に Agent から呼び出されます。冪等であるため、複数回呼び出しても影響はありません。
provider.startPolling()
stopPolling()stoppollingへの直接リンク
Polling Timer を停止します。
provider.stopPolling()
WebhookWebhookへの直接リンク
handleWebhook(request)handlewebhookrequestへの直接リンク
受信した Webhook リクエストを処理します。Payload を解析して Subscription と照合し、Notification Signal を送出するにはオーバーライドします。すぐに使用できる実装については、WebhookSignalProvider を参照してください。
Webhook リクエストを検証した後、アプリケーションで定義した 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への直接リンク
シャットダウン時に呼び出されます。デフォルトの実装では Polling を停止し、すべての Subscription をクリアします。
provider.stop()
通知通知への直接リンク
notify(notification, target)notifynotification-targetへの直接リンク
接続先の Agent に Notification 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:
source:
kind:
"pr-updated"、"new-message")。summary:
priority?:
payload?:
target:
threadId と resourceId を含める必要があります。型型への直接リンク
SignalSubscriptionsignalsubscriptionへの直接リンク
subscribe() が返す Subscription オブジェクトです。
id:
providerId:
threadId:
resourceId:
externalResourceId:
"github:owner/repo#123")。subscribedAt:
metadata:
SignalProviderTargetsignalprovidertargetへの直接リンク
特定の Agent Thread を識別します。
threadId:
resourceId:
agentId?:
型ガード型ガードへの直接リンク
isSignalProvider(obj)issignalproviderobjへの直接リンク
SignalProvider インスタンスかどうかを実行時に確認します。
import { isSignalProvider } from '@mastra/core/signals'
if (isSignalProvider(obj)) {
obj.connect(agent)
}
戻り値: boolean