Créer un fournisseur de signaux
Dans ce guide, vous allez créer un fournisseur de signaux qui interroge périodiquement un service externe et envoie une notification dans le thread d'un Agent dès qu'une ressource surveillée change. Vous apprendrez à étendre la classe de base SignalProvider, à suivre les abonnements, à émettre des notifications depuis une boucle de polling et à enregistrer le fournisseur auprès d'un Agent.
L'exemple surveille des pipelines de build dans un faux service de CI, mais ce modèle s'applique à toute source interrogée en mode pull : outil de suivi des problèmes, API d'état, file d'attente ou backend personnalisé.
Les fournisseurs de signaux sont en version bêta. Tant que l'API n'est pas stable, des changements incompatibles peuvent intervenir sans changement de version majeure.
PrérequisLien direct vers Prérequis
- Node.js
v22.13.0ou une version ultérieure installé - Une clé d'API provenant d'un fournisseur de modèles pris en charge
- Un projet Mastra existant. Si nécessaire, suivez le guide d'installation.
Ce guide suppose également que vous connaissez les principes généraux des signaux. Pour découvrir l'ensemble de l'API, consultez la référence de SignalProvider.
Ajouter le stockage des notificationsLien direct vers Ajouter le stockage des notifications
Un fournisseur de signaux envoie des signaux de notification dans les threads. Ces notifications nécessitent un adaptateur de stockage compatible avec le domaine des notifications. Configurez le stockage sur votre instance 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 et MongoDB prennent tous en charge les enregistrements de notification. Sans stockage des notifications, les appels à notify() du fournisseur lèvent une erreur à l'exécution.
Créer le client du service externeLien direct vers Créer le client du service externe
Les fournisseurs réels appellent une API externe. Pour que ce guide reste autonome, créez un petit client de CI fictif qui renvoie l'état du build d'un pipeline. Vous pourrez ensuite le remplacer par votre véritable client d'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 }
}
Ce client renvoie un état aléatoire à chaque appel. Lorsque vous brancherez une véritable API, seul ce fichier devra changer.
Créer le fournisseur de signauxLien direct vers Créer le fournisseur de signaux
Étendez SignalProvider, implémentez le champ abstrait id, définissez un pollInterval et surchargez poll(). La classe de base appelle périodiquement poll() avec tous les abonnements actifs. N'émettez une notification que pour les builds qui vous intéressent.
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 },
)
}
}
}
Quelques points à retenir :
subscribe()etunsubscribe()sont protégées dans la classe de base. Encapsulez-les dans vos propres méthodes publiques (watch/unwatch) afin que les appelants puissent gérer les abonnements.externalResourceIdpeut être n'importe quelle chaîne propre au fournisseur. Ici, il s'agit du nom du pipeline ; un fournisseur GitHub pourrait utiliser"github:owner/repo#123".notify()transmet un signal de notification au thread de l'Agent connecté. La méthode lève une erreur si le fournisseur n'a jamais été enregistré auprès d'un Agent.dedupeKeyempêche d'enregistrer deux fois le même échec.
Enregistrer le fournisseur auprès d'un AgentLien direct vers Enregistrer le fournisseur auprès d'un Agent
Transmettez le fournisseur à l'Agent via signals. L'Agent le connecte et démarre automatiquement la boucle de polling.
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],
})
Enregistrez l'Agent auprès de Mastra et ajoutez le stockage configuré lors de la première étape.
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',
}),
})
Abonner un threadLien direct vers Abonner un thread
Un fournisseur interroge uniquement les ressources surveillées par un thread. Abonnez un thread à un pipeline afin que poll() ait une ressource à vérifier.
import { ciSignals } from './agents/dev-agent'
ciSignals.watch({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme/app:main')
Exécutez ce code une fois après l'enregistrement de l'Agent, par exemple depuis un script de configuration ou une route d'API. L'abonnement réside dans le registre en mémoire du fournisseur ; vous devez donc vous réabonner après un redémarrage.
Tester le fournisseur de signauxLien direct vers Tester le fournisseur de signaux
Démarrez le serveur de développement :
- npm
- pnpm
- Yarn
- Bun
npm run dev
pnpm run dev
yarn dev
bun run dev
Vérifiez qu'un thread est abonné, puis observez les logs. Le client fictif renvoie un état aléatoire à chaque interrogation ; après quelques cycles, vous verrez donc un échec de build déclencher une notification pour le thread abonné.
Pour voir l'Agent réagir à la notification, abonnez-vous au thread et diffusez son flux :
const subscription = await devAgent.subscribeToThread({
resourceId: 'user_123',
threadId: 'thread_456',
})
for await (const chunk of subscription.stream) {
console.log(chunk)
}
Lorsqu'un build échoue, le modèle reçoit la notification dans son contexte :
<notification source="ci-signals" type="ci-status" priority="high" status="delivered">Build failed for acme/app:main</notification>
La sortie n'est pas déterministe, car le client fictif choisit l'état au hasard et le modèle formule librement sa réponse. Le libellé exact peut donc varier.
Étapes suivantesLien direct vers Étapes suivantes
Vous pouvez étendre ce fournisseur de signaux pour :
- Remplacer
fetchBuildStatus()par un véritable client d'API. - Conserver les abonnements afin qu'ils survivent à un redémarrage, puis les réhydrater dans
start(). - Ajouter un endpoint de webhook avec
handleWebhook()pour les sources en mode push. - Exposer les Tools
subscribeetunsubscribeavecgetTools()afin que l'Agent puisse gérer ses propres abonnements. - Utiliser
dedupeKeyetcoalesceKeylorsque les notifications doivent être dédupliquées ou regroupées.
Pour en savoir plus :