Aller au contenu principal

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'utilisation
Lien direct vers Exemple d'utilisation

Étendez PubSub et implémentez les quatre méthodes abstraites afin d'ajouter un backend personnalisé.

src/mastra/pubsub.ts
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 :

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'

export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})

Modes de livraison
Lien 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.

ModeDescription
pullLes 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.
pushLes é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éthodes
Lien direct vers Méthodes

Méthodes principales
Lien 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 relecture
Lien 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és
Lien direct vers Propriétés

supportedModes:

ReadonlyArray<"pull" | "push">
Modes de livraison pris en charge par l'implémentation. Utilise par défaut ["pull"].

supportsNativeBatching:

boolean
Indique si l'implémentation respecte 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.

Types
Lien direct vers Types

Event
Lien direct vers event

type:

string
Identifiant du type d'événement.

id:

string
Identifiant unique de l'événement, attribué par l'implémentation lors de la publication.

data:

any
Payload de l'événement.

runId:

string
Run auquel appartient l'événement.

createdAt:

Date
Horodatage attribué par l'implémentation lors de la publication.

index?:

number
Position séquentielle utilisée pour reprendre à partir d’un offset précis.

deliveryAttempt?:

number
Nombre de fois où l'événement a été livré. Commence à 1. Utilise par défaut 1 lorsque le backend ne suit pas les nouvelles livraisons.

SubscribeOptions
Lien direct vers subscribeoptions

group?:

string
Lorsque cette option est définie, les abonnés d'un même groupe se disputent les messages et chaque événement est livré à un seul membre. Lorsqu'elle est omise, chaque abonné reçoit chaque événement.

batch?:

SubscribeBatchOptions
Active la livraison par lots pour cet abonnement. Lorsque cette option est omise, les événements sont livrés un par un. Elle est respectée uniquement par les backends dont supportsNativeBatching vaut true.

SubscribeBatchOptions
Lien 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?:

number
Nombre maximal d'événements conservés avant de forcer un flush.

maxWaitMs?:

number
Durée maximale, en millisecondes, pendant laquelle l'événement le plus ancien peut rester dans le buffer. Le minuteur démarre lorsque le buffer passe de vide à non vide.

minIntervalMs?:

number
Durée minimale, en millisecondes, entre deux livraisons de lots consécutives. Même lorsque 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?:

(event: Event) => boolean
Lorsqu'elle renvoie 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?:

(events: Event[]) => Event[]
Appliquée au lot en file d'attente avant la livraison pour supprimer les événements remplacés. Doit renvoyer un sous-ensemble de son entrée en conservant l'identité des références ; renvoyer de nouveaux objets 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?:

number
= 256
Nombre maximal d'événements que le buffer peut contenir avant le déclenchement de la gestion du dépassement. Les événements marqués comme immédiats ne sont jamais supprimés lors d'un dépassement.

overflow?:

"drop-oldest" | "drop-newest" | "coalesce-or-drop-oldest"
= coalesce-or-drop-oldest
Stratégie de dépassement lorsque le buffer dépasse 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.

EventCallback
Lien direct vers eventcallback

Signature du callback des abonnés : (event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void.

event:

Event
Événement livré.

ack?:

() => Promise<void>
Confirme la réussite du traitement. L'événement est retiré de la file d'attente.

nack?:

() => Promise<void>
Signale un échec de traitement. L'événement est remis en file d'attente afin d'être livré à nouveau après un délai.