Aller au contenu principal

SignalProvider

Ajouté dans : @mastra/core@1.39.0

Classe de base abstraite permettant de créer des Providers de signal. Un Provider de signal surveille des sources externes (API, webhooks, streams d'événements) et transmet des signaux de notification aux fils de discussion des Agents au moyen d'un registre d'abonnements intégré.

Par défaut, les Providers de signal ne sont pas des Processors. Ceux qui doivent intercepter l'exécution d'un Agent renvoient des Processors depuis getInputProcessors() ou getOutputProcessors(). Ceux qui exposent des Tools appelables par un Agent les renvoient depuis getTools().

Pour un Provider prêt à l'emploi fondé sur les webhooks, consultez WebhookSignalProvider.

Exemple d'utilisation
Lien direct vers Exemple d'utilisation

Provider d'interrogation qui vérifie une API toutes les 30 secondes :

src/signals/slack-signals.ts
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 },
)
}
}
}
}

Enregistrez-le auprès d'un Agent :

src/mastra/index.ts
import { Agent } from '@mastra/core/agent'

const agent = new Agent({
id: 'agent',
signals: [new SlackSignals()],
})

L'Agent appelle connect(this) et enregistre tous les Processors ou Tools renvoyés par le Provider. Il lance ensuite l'interrogation.

Paramètres du constructeur
Lien direct vers Paramètres du constructeur

SignalProvider est abstraite. Les sous-classes appellent super() sans argument.

Propriétés
Lien direct vers Propriétés

id:

TId extends string
Identifiant unique de ce Provider. Les sous-classes doivent l’implémenter comme une propriété readonly.

name?:

string
Nom d'affichage lisible du Provider.

pollInterval?:

number
Intervalle d'interrogation en millisecondes. Lorsqu'il est défini, le framework appelle poll() à cet intervalle. Conservez undefined ou 0 pour les Providers utilisant uniquement des webhooks.

isConnected:

boolean
Indique si ce Provider est connecté à un Agent. Renvoie true après l'appel à connect(). Utilisé en interne pour éviter de reconnecter les éléments pendant Agent.__fork().

Méthodes
Lien direct vers Méthodes

Connexion
Lien direct vers Connexion

connect(agent)
Lien direct vers connectagent

Appelée par le constructeur de l'Agent. Établit le lien bidirectionnel afin que le Provider puisse renvoyer des signaux à l'Agent. Remplacez cette méthode pour exécuter une configuration supplémentaire une fois le lien établi. Appelez toujours 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)
Lien direct vers __registermastramastra

Appelée lorsque l'Agent du Provider est enregistré auprès d'une instance Mastra. Remplacez cette méthode pour accéder au stockage ou à d'autres services Mastra. Appelez toujours super.__registerMastra(mastra).

override __registerMastra(mastra) {
super.__registerMastra(mastra)
// this.mastra is now available
}

Intégration des Processors et des Tools
Lien direct vers Intégration des Processors et des Tools

getInputProcessors()
Lien direct vers getinputprocessors

Renvoie les Processors d'entrée que ce Provider doit enregistrer auprès de l'Agent. Remplacez cette méthode lorsque votre Provider intercepte les étapes d'entrée de l'Agent, par exemple pour injecter des indications de contexte ou détecter des appels de Tool.

getInputProcessors() {
return [this]
}

Renvoie : InputProcessorOrWorkflow[]

getOutputProcessors()
Lien direct vers getoutputprocessors

Renvoie les Processors de sortie que ce Provider doit enregistrer auprès de l'Agent. Remplacez cette méthode lorsque votre Provider intercepte les étapes de sortie de l'Agent.

getOutputProcessors() {
return [this]
}

Renvoie : OutputProcessorOrWorkflow[]

getTools()
Lien direct vers gettools

Renvoie les Tools que ce Provider expose à l'Agent. Remplacez cette méthode lorsque votre Provider ajoute des Tools appelables par l'Agent, tels que des commandes d'abonnement ou de désabonnement.

getTools() {
return {
subscribe_pr: createTool({ /* ... */ }),
unsubscribe_pr: createTool({ /* ... */ }),
}
}

Renvoie : Record<string, unknown>

Suivi des abonnements
Lien direct vers Suivi des abonnements

subscribe(target, externalResourceId, metadata?)
Lien direct vers subscribetarget-externalresourceid-metadata

Abonne un fil de discussion à une ressource externe. Il s'agit d'une méthode protégée : appelez-la depuis l'implémentation de votre Provider.

const sub = this.subscribe(
{ threadId: 'thread-1', resourceId: 'user-1' },
'github:mastra-ai/mastra#123',
{ pr: 123 },
)

Renvoie : SignalSubscription : l'abonnement créé ou l'abonnement existant avec les métadonnées fusionnées.

target:

SignalProviderTarget
Fil de discussion à abonner. Doit inclure threadId et resourceId.

externalResourceId:

string
Identifiant de la ressource externe propre au Provider (par exemple, "github:owner/repo#123").

metadata?:

Record<string, unknown>
Données supplémentaires à stocker avec l'abonnement. Fusionnées avec les métadonnées existantes en cas d'abonnement en double.

unsubscribe(target, externalResourceId)
Lien direct vers unsubscribetarget-externalresourceid

Supprime un abonnement.

const removed = this.unsubscribe(
{ threadId: 'thread-1', resourceId: 'user-1' },
'github:mastra-ai/mastra#123',
)

Renvoie : boolean : true en cas de suppression, false si aucun abonnement correspondant n'existait.

getSubscriptions()
Lien direct vers getsubscriptions

Renvoie tous les abonnements actifs de ce Provider.

const allSubs = this.getSubscriptions()

Renvoie : SignalSubscription[]

getSubscriptionsForResource(externalResourceId)
Lien direct vers getsubscriptionsforresourceexternalresourceid

Renvoie tous les abonnements associés à une ressource externe précise.

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

Renvoie : SignalSubscription[]

getSubscriptionsForThread(target)
Lien direct vers getsubscriptionsforthreadtarget

Renvoie tous les abonnements associés à un fil de discussion précis.

const subs = this.getSubscriptionsForThread({
threadId: 'thread-1',
resourceId: 'user-1',
})

Renvoie : SignalSubscription[]

hasSubscription(target, externalResourceId)
Lien direct vers hassubscriptiontarget-externalresourceid

Vérifie si un abonnement existe.

if (this.hasSubscription(target, 'github:mastra-ai/mastra#123')) {
// already subscribed
}

Renvoie : boolean

unsubscribeAll(target)
Lien direct vers unsubscribealltarget

Supprime tous les abonnements associés à un fil de discussion.

const removed = this.unsubscribeAll({
threadId: 'thread-1',
resourceId: 'user-1',
})

Renvoie : number : nombre d'abonnements supprimés.

subscriptionCount
Lien direct vers subscriptioncount

Nombre total d'abonnements actifs de ce Provider.

if (this.subscriptionCount === 0) {
// nothing to poll
}

Renvoie : number

Interrogation
Lien direct vers Interrogation

poll(subscriptions)
Lien direct vers pollsubscriptions

Appelée à chaque cycle d'interrogation avec tous les abonnements actifs. Remplacez cette méthode pour vérifier les sources externes et émettre des notifications. Le framework empêche le chevauchement des cycles : si un appel à poll() dure plus longtemps que pollInterval, le cycle suivant est ignoré.

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()
Lien direct vers startpolling

Démarre le minuteur d'interrogation. Appelée par l'Agent après connect(). Cette méthode est idempotente : plusieurs appels restent sans effet supplémentaire.

provider.startPolling()

stopPolling()
Lien direct vers stoppolling

Arrête le minuteur d'interrogation.

provider.stopPolling()

Webhooks
Lien direct vers Webhooks

handleWebhook(request)
Lien direct vers handlewebhookrequest

Traite une requête webhook entrante. Remplacez cette méthode pour analyser le payload et le faire correspondre aux abonnements, puis émettre des signaux de notification. Consultez WebhookSignalProvider pour une implémentation prête à l'emploi.

Appelez cette méthode depuis un endpoint HTTP défini par l'application après avoir vérifié la requête webhook.

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

Renvoie : Promise<{ status?: number; body?: unknown }>

Cycle de vie
Lien direct vers Cycle de vie

start()
Lien direct vers start

Appelée après connect() pour exécuter l'initialisation asynchrone. Remplacez cette méthode lorsque la configuration nécessite que l'Agent ou l'instance Mastra soit disponible.

async start() {
await this.loadInitialState()
}

stop()
Lien direct vers stop

Appelée lors de l'arrêt. L'implémentation par défaut interrompt l'interrogation et efface tous les abonnements.

provider.stop()

Notifications
Lien direct vers Notifications

notify(notification, target)
Lien direct vers notifynotification-target

Envoie un signal de notification à l'Agent connecté. Il s'agit d'un wrapper protégé pratique autour de agent.sendNotificationSignal().

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:

object
Payload de la notification.
object

source:

string
Identifiant de la source de la notification.

kind:

string
Type d'événement (par exemple, "pr-updated", "new-message").

summary:

string
Résumé lisible de l'événement.

priority?:

"high" | "medium" | "low"
Priorité de la notification.

payload?:

unknown
Données arbitraires jointes à la notification.

target:

SignalProviderTarget
Fil de discussion à notifier. Doit inclure threadId et resourceId.

Types
Lien direct vers Types

SignalSubscription
Lien direct vers signalsubscription

Objet d'abonnement renvoyé par subscribe().

id:

string
Identifiant unique de l'abonnement.

providerId:

string
Provider auquel appartient cet abonnement.

threadId:

string
Fil de discussion qui reçoit les signaux.

resourceId:

string
Ressource à laquelle appartient le fil de discussion.

externalResourceId:

string
Identifiant de la ressource externe propre au Provider (par exemple, "github:owner/repo#123").

subscribedAt:

Date
Date de création de l'abonnement.

metadata:

Record<string, unknown>
Métadonnées propres au Provider stockées avec l'abonnement.

SignalProviderTarget
Lien direct vers signalprovidertarget

Identifie un fil de discussion précis d'un Agent.

threadId:

string
Fil de discussion à cibler.

resourceId:

string
Ressource à laquelle appartient le fil de discussion.

agentId?:

string
Identifiant de l'Agent.

Garde de type
Lien direct vers Garde de type

isSignalProvider(obj)
Lien direct vers issignalproviderobj

Vérification à l'exécution des instances de SignalProvider.

import { isSignalProvider } from '@mastra/core/signals'

if (isSignalProvider(obj)) {
obj.connect(agent)
}

Renvoie : boolean