> Discover all available pages from the documentation index: https://mastra.zisheng.pro/fr/llms.txt # 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`](https://mastra.zisheng.pro/fr/reference/pubsub/event-emitter) 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`](https://mastra.zisheng.pro/fr/reference/pubsub/event-emitter), [`UnixSocketPubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/unix-socket-pubsub), [`CachingPubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/caching-pubsub), [`RedisStreamsPubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/redis-streams) et [`GoogleCloudPubSub`](https://mastra.zisheng.pro/fr/reference/pubsub/google-cloud-pubsub). ## Exemple d'utilisation Étendez `PubSub` et implémentez les quatre méthodes abstraites afin d'ajouter un backend personnalisé. ```typescript 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): Promise { // Deliver the event to subscribers of `topic`. } async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise { // Register `cb` to receive events published to `topic`. } async unsubscribe(topic: string, cb: EventCallback): Promise { // Remove a previously registered callback. } async flush(): Promise { // Wait for any in-flight deliveries to settle. } } ``` Transmettez l'instance au constructeur de [Mastra](https://mastra.zisheng.pro/fr/reference/core/mastra-class) : ```typescript import { Mastra } from '@mastra/core' import { CustomPubSub } from './pubsub' export const mastra = new Mastra({ pubsub: new CustomPubSub(), }) ``` ## 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éthodes ### Méthodes principales #### `publish(topic, event)` Publie un événement dans un topic. Les champs `id` et `createdAt` sont attribués par l'implémentation. ```typescript await pubsub.publish('my-topic', { type: 'example', data: { value: 1 }, runId: 'run-123', }) ``` #### `subscribe(topic, 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`](#properties) du backend vaut `true`. Les autres backends ignorent cette option et livrent les événements un par un. ```typescript await pubsub.subscribe('my-topic', (event, ack, nack) => { console.log(event) }) ``` #### `unsubscribe(topic, cb)` Supprime d'un topic un callback précédemment enregistré. ```typescript await pubsub.unsubscribe('my-topic', callback) ``` #### `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. ```typescript await pubsub.flush() ``` #### `clearTopic(topic)` 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`](https://mastra.zisheng.pro/fr/reference/pubsub/redis-streams), 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. ```typescript await pubsub.clearTopic('workflow.events.v2.run-123') ``` ### 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`](https://mastra.zisheng.pro/fr/reference/pubsub/caching-pubsub) remplace ces méthodes afin de relire les événements mis en cache. #### `getHistory(topic, 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. ```typescript const events = await pubsub.getHistory('my-topic', 0) ``` Renvoie : `Promise` #### `subscribeWithReplay(topic, cb)` Relit les événements mis en cache, puis s'abonne aux événements en direct. ```typescript await pubsub.subscribeWithReplay('my-topic', event => { console.log(event) }) ``` #### `subscribeFromOffset(topic, 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. ```typescript await pubsub.subscribeFromOffset('my-topic', 42, event => { console.log(event) }) ``` ## 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 ### `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` **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` 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`): 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. (Default: `256`) **overflow** (`"drop-oldest" | "drop-newest" | "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. (Default: `coalesce-or-drop-oldest`) ### `EventCallback` Signature du callback des abonnés : `(event: Event, ack?: () => Promise, nack?: () => Promise) => void`. **event** (`Event`): Événement livré. **ack** (`() => Promise`): Confirme la réussite du traitement. L'événement est retiré de la file d'attente. **nack** (`() => Promise`): Signale un échec de traitement. L'événement est remis en file d'attente afin d'être livré à nouveau après un délai.