Signal provider の構築
このガイドでは、一定間隔で外部サービスをポーリングし、監視対象のリソースが変化するたびに Agent スレッドへ通知を送る signal provider を構築します。SignalProvider 基底クラスの拡張、サブスクリプションの追跡、ポーリングループからの通知送信、Agent への provider 登録について説明します。
この例では架空の CI サービスのビルドパイプラインを監視しますが、このパターンは課題管理システム、ステータス API、キュー、独自バックエンドなど、任意の pull 型ソースに適用できます。
Signal provider は beta です。API が安定するまでは、major version を上げずに破壊的変更が行われる可能性があります。
前提条件前提条件への直接リンク
- Node.js
v22.13.0以降がインストールされていること - サポートされている Model Provider の API キー
- 既存の Mastra プロジェクト。必要に応じてインストールガイドに従ってください。
また、このガイドでは signal の概要を理解していることを前提とします。API の全体像については、SignalProvider referenceを参照してください。
Notification storage の追加Notification storage の追加への直接リンク
signal provider はスレッドに 通知 signal を送り、通知には通知 domain をサポートする storage アダプターが必要です。Mastra インスタンスに storage を設定します。
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 はいずれも通知 record をサポートします。通知 storage がない場合、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 を接続する際に変更するのは、このファイルだけです。
Signal provider の構築Signal provider の構築への直接リンク
SignalProvider を拡張し、abstract な id フィールドを実装して pollInterval を設定し、poll() を override します。基底クラスは、一定間隔で有効なすべてのサブスクリプションとともに 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 です。呼び出し側がサブスクリプションを管理できるよう、独自の public method(watch/unwatch)でラップします。externalResourceIdには provider 固有の任意の string を指定できます。ここではパイプライン名です。GitHub provider なら"github:owner/repo#123"のような値を使用できます。notify()は、接続された Agent のスレッドに通知 signal を転送します。provider が一度も Agent に登録されていない場合は例外を送出します。dedupeKeyは、同じ failure が重複して保存されるのを防ぎます。
Agent への provider の登録Agent への provider の登録への直接リンク
signals を介して Agent に provider を渡します。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],
})
Agent を Mastra に登録し、最初のステップで使用した storage を追加します。
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 の subscribeThread の subscribeへの直接リンク
provider は、スレッドが監視しているリソースだけをポーリングします。poll() に確認対象を与えるため、スレッドをパイプラインに subscribe します。
import { ciSignals } from './agents/dev-agent'
ciSignals.watch({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme/app:main')
Agent の登録後に、setup スクリプトや API route などから一度実行します。サブスクリプションは provider の in-memory registry に保存されるため、再起動後は再度 subscribe してください。
Signal provider のテストSignal provider のテストへの直接リンク
開発サーバーを起動します。
- npm
- pnpm
- Yarn
- Bun
npm run dev
pnpm run dev
yarn dev
bun run dev
スレッドが subscribe されていることを確認し、ログを監視します。架空のクライアントは poll ごとにランダムなステータスを返すため、数回以内に failed ビルドが発生し、subscribe されたスレッドに通知が送られます。
Agent が通知に反応する様子を確認するには、スレッドを subscribe してストリームします。
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>
架空のクライアントはステータスをランダムに生成し、モデルも自由に応答を表現するため、出力は非決定的で、実際の文言は異なります。
次のステップ次のステップへの直接リンク
この signal provider は次のように拡張できます。
fetchBuildStatus()を実際の API クライアントに置き換える。- 再起動後も保持されるようサブスクリプションを永続化し、
start()で再構築する。 - push 型ソース向けに、webhook エントリーポイントを
handleWebhook()で追加する。 subscribeとunsubscribeTool をgetTools()で公開し、Agent 自身がサブスクリプションを管理できるようにする。- 通知の重複排除やバッチ化が必要な場合に、
dedupeKeyとcoalesceKeyを使用する。
関連情報: