Aller au contenu principal

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é.

beta

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érequis
Lien direct vers Prérequis

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 notifications
Lien 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.

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 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 externe
Lien 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.

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

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 signaux
Lien 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.

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

Quelques points à retenir :

  • subscribe() et unsubscribe() 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.
  • externalResourceId peut ê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.
  • dedupeKey empêche d'enregistrer deux fois le même échec.

Enregistrer le fournisseur auprès d'un Agent
Lien 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.

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

Enregistrez l'Agent auprès de Mastra et ajoutez le stockage configuré lors de la première étape.

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

Abonner un thread
Lien 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.

src/mastra/subscribe.ts
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 signaux
Lien direct vers Tester le fournisseur de signaux

Démarrez le serveur de développement :

npm 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 :

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

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 suivantes
Lien 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 subscribe et unsubscribe avec getTools() afin que l'Agent puisse gérer ses propres abonnements.
  • Utiliser dedupeKey et coalesceKey lorsque les notifications doivent être dédupliquées ou regroupées.

Pour en savoir plus :