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'utilisationLien direct vers Exemple d'utilisation
Provider d'interrogation qui vérifie une API toutes les 30 secondes :
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 :
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 constructeurLien direct vers Paramètres du constructeur
SignalProvider est abstraite. Les sous-classes appellent super() sans argument.
PropriétésLien direct vers Propriétés
id:
name?:
pollInterval?:
poll() à cet intervalle. Conservez undefined ou 0 pour les Providers utilisant uniquement des webhooks.isConnected:
true après l'appel à connect(). Utilisé en interne pour éviter de reconnecter les éléments pendant Agent.__fork().MéthodesLien direct vers Méthodes
ConnexionLien 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 ToolsLien 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 abonnementsLien 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:
threadId et resourceId.externalResourceId:
"github:owner/repo#123").metadata?:
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.
subscriptionCountLien direct vers subscriptioncount
Nombre total d'abonnements actifs de ce Provider.
if (this.subscriptionCount === 0) {
// nothing to poll
}
Renvoie : number
InterrogationLien 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()
WebhooksLien 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 vieLien 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()
NotificationsLien 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:
source:
kind:
"pr-updated", "new-message").summary:
priority?:
payload?:
target:
threadId et resourceId.TypesLien direct vers Types
SignalSubscriptionLien direct vers signalsubscription
Objet d'abonnement renvoyé par subscribe().
id:
providerId:
threadId:
resourceId:
externalResourceId:
"github:owner/repo#123").subscribedAt:
metadata:
SignalProviderTargetLien direct vers signalprovidertarget
Identifie un fil de discussion précis d'un Agent.
threadId:
resourceId:
agentId?:
Garde de typeLien 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