PubSub
PubSub est la classe de base abstraite du système d'événements de Mastra. Elle définit le contrat implémenté par chaque backend pub/sub, de sorte que le reste de Mastra puisse publier des événements et s'y abonner sans connaître le transport utilisé.
Mastra utilise pub/sub en interne pour le traitement des événements de Workflow, le streaming et la communication entre composants. La plupart des applications utilisent l'implémentation par défaut EventEmitterPubSub et ne construisent jamais directement un PubSub. Implémentez cette classe uniquement si vous avez besoin d'un transport personnalisé.
Pour les implémentations intégrées, consultez EventEmitterPubSub, UnixSocketPubSub, CachingPubSub, RedisStreamsPubSub et GoogleCloudPubSub.
Exemple d'utilisationLien direct vers Exemple d'utilisation
Étendez PubSub et implémentez les quatre méthodes abstraites afin d'ajouter un backend personnalisé.
import { PubSub } from '@mastra/core/events'
import type { Event, EventCallback, SubscribeOptions } from '@mastra/core/events'
export class CustomPubSub extends PubSub {
async publish(topic: string, event: Omit<Event, 'id' | 'createdAt'>): Promise<void> {
// Deliver the event to subscribers of `topic`.
}
async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise<void> {
// Register `cb` to receive events published to `topic`.
}
async unsubscribe(topic: string, cb: EventCallback): Promise<void> {
// Remove a previously registered callback.
}
async flush(): Promise<void> {
// Wait for any in-flight deliveries to settle.
}
}
Transmettez l'instance au constructeur de Mastra :
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'
export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})
Modes de livraisonLien direct vers Modes de livraison
Un PubSub déclare les modes de livraison qu'il prend en charge au moyen de la propriété supportedModes. Mastra lit cette propriété pour déterminer s'il doit exécuter un Worker de longue durée qui récupère les événements.
| Mode | Description |
|---|---|
pull | Les consommateurs lisent activement les données du broker, par exemple avec Redis Streams XREADGROUP. Mastra exécute un Worker d'orchestration pour effectuer la lecture. |
push | Les événements arrivent sans demande du consommateur, dans le processus ou par l'intermédiaire d'un endpoint HTTP. Aucune boucle de lecture n'est nécessaire. |
La valeur par défaut est ['pull'], afin que les implémentations personnalisées conservent le comportement actuel sauf si elles activent la livraison push.
MéthodesLien direct vers Méthodes
Méthodes principalesLien direct vers Méthodes principales
publish(topic, event)Lien direct vers publishtopic-event
Publie un événement dans un topic. Les champs id et createdAt sont attribués par l'implémentation.
await pubsub.publish('my-topic', {
type: 'example',
data: { value: 1 },
runId: 'run-123',
})
subscribe(topic, cb, options?)Lien direct vers subscribetopic-cb-options
Enregistre un callback pour recevoir les événements publiés dans un topic. Lorsque options.group est défini, les abonnés d'un même groupe se disputent les messages et chaque événement est livré à un seul membre. Sans groupe, chaque abonné reçoit chaque événement.
Transmettez options.batch pour activer la livraison par lots. La signature du callback reste inchangée : un lot de N événements est livré sous la forme de N appels consécutifs à cb(event, ack, nack), dans l'ordre de publication. Le traitement par lots est respecté uniquement lorsque la propriété supportsNativeBatching du backend vaut true. Les autres backends ignorent cette option et livrent les événements un par un.
await pubsub.subscribe('my-topic', (event, ack, nack) => {
console.log(event)
})
unsubscribe(topic, cb)Lien direct vers unsubscribetopic-cb
Supprime d'un topic un callback précédemment enregistré.
await pubsub.unsubscribe('my-topic', callback)
flush()Lien direct vers flush
Attend la fin de toutes les livraisons en cours. Appelez cette méthode avant l'arrêt afin d'éviter de perdre des événements.
await pubsub.flush()
clearTopic(topic)Lien direct vers cleartopictopic
Supprime tout l'état conservé pour un topic (historique en cache, entrées de flux persistantes et groupes de consommateurs) dès qu'aucun autre événement ne doit y être publié. Les cycles de vie des Runs de Mastra (Agents durables et moteur de Workflow événementiel) appellent automatiquement cette méthode lorsqu'un Run atteint un état terminal, afin que les topics propres à chaque Run ne s'accumulent pas sur les transports qui conservent les messages.
L'implémentation par défaut ne fait rien : les transports qui ne conservent rien par topic (tels que EventEmitterPubSub) n'ont rien à effacer. Les backends qui rendent les messages persistants, comme RedisStreamsPubSub, la remplacent. Le contrat fonctionne selon le principe du meilleur effort : les implémentations journalisent les échecs au lieu de lever une erreur, car les appelants l'invoquent sans attendre de résultat aux limites du nettoyage.
await pubsub.clearTopic('workflow.events.v2.run-123')
Méthodes de relectureLien direct vers Méthodes de relecture
Ces méthodes permettent de reprendre un flux après une déconnexion. Les implémentations par défaut utilisent un subscribe standard ; les backends qui ne prennent pas en charge l'historique fonctionnent donc uniquement en direct. CachingPubSub remplace ces méthodes afin de relire les événements mis en cache.
getHistory(topic, offset?)Lien direct vers gethistorytopic-offset
Renvoie les événements mis en cache pour un topic à partir de offset. Renvoie un tableau vide lorsque le backend ne possède aucun historique.
const events = await pubsub.getHistory('my-topic', 0)
Renvoie : Promise<Event[]>
subscribeWithReplay(topic, cb)Lien direct vers subscribewithreplaytopic-cb
Relit les événements mis en cache, puis s'abonne aux événements en direct.
await pubsub.subscribeWithReplay('my-topic', event => {
console.log(event)
})
subscribeFromOffset(topic, offset, cb)Lien direct vers subscribefromoffsettopic-offset-cb
Relit les événements mis en cache à partir d'une position connue, puis s'abonne aux événements en direct. Cette méthode est plus efficace qu'une relecture complète lorsque le client connaît sa dernière position.
await pubsub.subscribeFromOffset('my-topic', 42, event => {
console.log(event)
})
PropriétésLien direct vers Propriétés
supportedModes:
["pull"].supportsNativeBatching:
options.batch dans subscribe(). Utilise par défaut false. Les backends qui intègrent le traitement par lots en interne remplacent cette propriété et renvoient true.TypesLien direct vers Types
EventLien direct vers event
type:
id:
data:
runId:
createdAt:
index?:
deliveryAttempt?:
SubscribeOptionsLien direct vers subscribeoptions
group?:
batch?:
supportsNativeBatching vaut true.SubscribeBatchOptionsLien direct vers subscribebatchoptions
Politique de traitement par lots propre à chaque abonnement. La signature du callback ne change pas. Un lot de N événements devient N invocations consécutives du callback dans l'ordre de publication.
maxSize?:
maxWaitMs?:
minIntervalMs?:
maxSize ou maxWaitMs devrait se déclencher, le buffer est conservé jusqu'à ce que cet intervalle se soit écoulé depuis la dernière livraison.isImmediate?:
true pour un événement, le buffer est vidé immédiatement à la publication, sous réserve de minIntervalMs. Il s'agit d'un mécanisme d'échappement propre à chaque événement.coalesce?:
Event constitue une violation du contrat et entraîne la suppression de tout le lot. L'ordre des événements conservés est préservé.maxBufferSize?:
overflow?:
maxBufferSize. coalesce-or-drop-oldest exécute d'abord coalesce, puis supprime les événements les plus anciens si la limite est toujours dépassée.EventCallbackLien direct vers eventcallback
Signature du callback des abonnés : (event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void.