Fournisseurs de signaux
Ajouté dans : @mastra/core@1.39.0
Cette fonctionnalité est en bêta. Des changements incompatibles peuvent survenir sans hausse de version majeure tant que l’API n’est pas stable.
Un fournisseur de signaux surveille une source externe, comme GitHub, Slack, un système d’intégration continue (CI) ou votre propre API, et envoie des signaux de notification aux fils de discussion d’agents abonnés.
Quand utiliser un fournisseur de signauxLien direct vers Quand utiliser un fournisseur de signaux
Utilisez un fournisseur de signaux lorsqu’un système externe produit des événements auxquels un agent doit réagir et que vous souhaitez confier à Mastra la gestion des abonnements.
- La source émet des événements liés à une ressource qui intéresse un fil de discussion, par exemple une demande de fusion, un canal ou une compilation.
- Vous souhaitez centraliser le suivi des fils de discussion qui surveillent chaque ressource externe.
- Vous souhaitez recevoir les événements par interrogation périodique, par webhook ou par les deux moyens.
Si vous devez seulement envoyer un événement ponctuel à un fil de discussion, appelez plutôt agent.sendNotificationSignal() directement.
Fonctionnement des fournisseurs de signauxLien direct vers Fonctionnement des fournisseurs de signaux
Un fournisseur de signaux constitue le côté producteur du système de signaux. Il transmet les événements externes à un fil de discussion, tandis que les API de signaux déterminent la façon dont celui-ci les consomme.
Un fournisseur de signaux réunit trois fonctionnalités :
- Suivi des abonnements : la classe de base
SignalProviderconserve un registre en mémoire qui associe chaque fil de discussion d’agent aux ressources externes qu’il surveille. - Ingestion : remplacez
poll()pour les sources interrogées par l’application, ouhandleWebhook()pour les sources qui lui envoient des données. - Distribution : lorsqu’un événement correspond à un abonnement, appelez la méthode utilitaire protégée
notify()pour transmettre un signal de notification au fil de discussion de l’agent connecté.
Pour enregistrer un fournisseur, transmettez-le à un agent.
L’agent connecte le fournisseur et lance l’interrogation périodique lorsqu’un pollInterval est défini. Il intègre également les processeurs et outils exposés par le fournisseur.
import { Agent } from '@mastra/core/agent'
import { CiSignals } from '../signals/ci-signals'
export const supportAgent = new Agent({
id: 'support-agent',
name: 'Support Agent',
instructions: 'Help the user triage updates.',
model: 'openai/gpt-5.6-sol',
signals: [new CiSignals()],
})
La distribution des notifications nécessite un adaptateur de stockage prenant en charge les notifications, comme libSQL, PostgreSQL ou MongoDB. Configurez le stockage sur l’instance Mastra afin que notify() puisse enregistrer les notifications.
Démarrage rapideLien direct vers Démarrage rapide
L’exemple suivant présente un fournisseur reposant sur l’interrogation périodique qui surveille des pipelines CI et émet une notification lorsqu’un pipeline suivi échoue.
import { SignalProvider } from '@mastra/core/signals'
import type { SignalProviderTarget, SignalSubscription } from '@mastra/core/signals'
type BuildStatus = {
id: string
status: 'passed' | 'failed'
}
const builds = new Map<string, BuildStatus>([
['acme-app-main', { id: 'build_123', status: 'failed' }],
])
async function fetchBuildStatus(pipeline: string): Promise<BuildStatus> {
return builds.get(pipeline) ?? { id: 'build_unknown', status: 'passed' }
}
export class CiSignals extends SignalProvider<'ci-signals'> {
readonly id = 'ci-signals' as const
readonly pollInterval = 30_000
watch(target: SignalProviderTarget, pipeline: string) {
return this.subscribe(target, pipeline)
}
unwatch(target: SignalProviderTarget, pipeline: string) {
return this.unsubscribe(target, pipeline)
}
async poll(subscriptions: SignalSubscription[]) {
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 },
)
}
}
}
Enregistrez le fournisseur auprès d’un agent, puis abonnez un fil de discussion au pipeline que vous souhaitez surveiller.
import { Agent } from '@mastra/core/agent'
import { CiSignals } from '../signals/ci-signals'
export const ciSignals = new CiSignals()
export const supportAgent = new Agent({
id: 'support-agent',
name: 'Support Agent',
instructions: 'Help the user triage CI updates.',
model: 'openai/gpt-5.6-sol',
signals: [ciSignals],
})
ciSignals.watch({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme-app-main')
Mastra appelle poll() selon le pollInterval en lui transmettant tous les abonnements actifs. En l’absence d’abonnement, Mastra ignore le cycle et ne fait pas chevaucher les cycles. Ainsi, une exécution lente de poll() ne s’exécute pas en parallèle avec elle-même.
Pour créer de bout en bout un fournisseur reposant sur l’interrogation périodique, avec stockage des notifications, enregistrement de l’agent, abonnement du fil de discussion et tests, consultez Créer un fournisseur de signaux.
Fournisseurs par interrogation périodique et par webhookLien direct vers Fournisseurs par interrogation périodique et par webhook
Utilisez l’interrogation périodique lorsque la source externe n’envoie pas d’événements à votre application. Définissez pollInterval et remplacez poll(subscriptions). Chaque abonnement comprend le fil de discussion cible et l’identifiant de la ressource externe à examiner.
Utilisez des webhooks lorsque la source externe peut appeler votre application. Remplacez handleWebhook(request), analysez la charge utile, recherchez les abonnements correspondants, puis appelez notify() pour chacun d’eux.
import { SignalProvider } from '@mastra/core/signals'
import type { SignalProviderWebhookRequest } from '@mastra/core/signals'
export class CiSignals extends SignalProvider<'ci-signals'> {
readonly id = 'ci-signals' as const
async handleWebhook(request: SignalProviderWebhookRequest) {
const payload = request.body as { pipeline: string; status: string }
const subscriptions = this.getSubscriptionsForResource(payload.pipeline)
for (const sub of subscriptions) {
await this.notify(
{
source: this.id,
kind: 'ci-status',
priority: 'high',
summary: `Build ${payload.status} for ${payload.pipeline}`,
payload,
},
{ resourceId: sub.resourceId, threadId: sub.threadId },
)
}
return { status: 200, body: { matched: subscriptions.length } }
}
}
handleWebhook() est une méthode du fournisseur, et non une route HTTP montée automatiquement. Appelez-la depuis votre propre point de terminaison en lui transmettant le corps et les en-têtes de la requête, ainsi que les éventuels paramètres de route. Consultez la référence de SignalProvider pour en savoir plus sur les abonnements, l’interrogation périodique, le cycle de vie et notify(). Pour connaître la structure complète de la charge utile d’une notification, notamment les champs de déduplication et de regroupement, consultez la référence de Agent.sendNotificationSignal().
Fournisseur de webhooks intégréLien direct vers Fournisseur de webhooks intégré
Pour les sources de webhooks génériques, utilisez WebhookSignalProvider au lieu d’écrire une sous-classe. Configurez-le avec une fonction qui extrait un identifiant de ressource de la charge utile et, éventuellement, une fonction qui construit la notification.
import { Agent } from '@mastra/core/agent'
import { WebhookSignalProvider } from '@mastra/core/signals'
const webhooks = new WebhookSignalProvider({
extractResourceId: payload => (payload as { repository: string }).repository,
buildNotification: (payload, sub) => ({
source: 'ci',
kind: 'build-status',
priority: 'medium',
summary: `Build ${(payload as { status: string }).status} for ${sub.externalResourceId}`,
}),
})
export const supportAgent = new Agent({
id: 'support-agent',
name: 'Support Agent',
instructions: 'Help the user triage updates.',
model: 'openai/gpt-5.6-sol',
signals: [webhooks],
})
webhooks.subscribeThread({ resourceId: 'user_123', threadId: 'thread_456' }, 'acme/app')
Lorsqu’un webhook arrive, appelez webhooks.handleWebhook({ body, headers }) depuis votre route. Le fournisseur compare l’identifiant de ressource extrait à ses abonnements et notifie chaque fil de discussion correspondant.
Fonctionnalités avancées des fournisseursLien direct vers Fonctionnalités avancées des fournisseurs
Un fournisseur peut prendre en charge d’autres fonctionnalités que l’ingestion d’événements. Ajoutez uniquement celles dont votre source a besoin.
- Abonnements durables : le registre de base est conservé en mémoire et propre à chaque processus. Gérez vous-même la persistance des abonnements lorsqu’ils doivent survivre à un redémarrage, puis restaurez-les dans
start(). - Fonctions de cycle de vie : remplacez
start()pour la configuration asynchrone etstop()pour le nettoyage. Appelezsuper.stop()lorsque vous remplacezstop()afin que le fournisseur de base puisse arrêter l’interrogation périodique et vider son registre. - Processeurs et outils : renvoyez les processeurs depuis
getInputProcessors()ougetOutputProcessors(), et les outils que l’agent peut appeler depuisgetTools().
Le package @mastra/github-signals est un fournisseur de signaux destiné à la production. Il surveille les demandes de fusion GitHub et notifie les fils de discussion des commentaires, de l’état des revues, de l’état de l’intégration continue et des fusions. Utilisez-le comme référence pour l’interrogation périodique, les abonnements durables, les outils, les processeurs et les fonctions de cycle de vie.
import { Agent } from '@mastra/core/agent'
import { GithubSignals } from '@mastra/github-signals'
export const devAgent = new Agent({
id: 'dev-agent',
name: 'Dev Agent',
instructions: 'Help triage pull request activity.',
model: 'openai/gpt-5.6-sol',
signals: [new GithubSignals()],
})