Aller au contenu principal

UnixSocketPubSub

UnixSocketPubSub est une implémentation de PubSub qui transmet des événements entre plusieurs processus sur un même hôte au moyen d’un socket de domaine Unix. Elle élit un processus comme broker, auquel les autres processus se connectent en tant que clients. Si le broker s’arrête, les clients restants élisent automatiquement un nouveau broker.

Utilisez-la lorsque plusieurs processus locaux doivent partager un flux, par exemple pour coordonner des flux de threads dans l’interface de terminal Mastra Code. Pour une diffusion au sein d’un seul processus, utilisez EventEmitterPubSub. Pour une diffusion distribuée entre plusieurs hôtes, utilisez RedisStreamsPubSub ou GoogleCloudPubSub.

UnixSocketPubSub est un transport push : les événements arrivent sans boucle de lecture, Mastra n’exécute donc aucun worker pull pour ce transport.

Exemple d’utilisation
Lien direct vers Exemple d’utilisation

Transmettez un chemin de socket partagé par tous les processus participants.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { UnixSocketPubSub } from '@mastra/core/events'

export const mastra = new Mastra({
pubsub: new UnixSocketPubSub('/tmp/mastra/events.sock'),
})

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

socketPath:

string
Chemin du socket de domaine Unix. Tous les processus qui partagent un flux doivent utiliser le même chemin.

options?:

UnixSocketPubSubOptions
Configuration facultative.
number

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

socketPath:

string
Chemin du socket transmis au constructeur.

supportedModes:

ReadonlyArray<"pull" | "push">
Renvoie ["push"].

isBroker:

boolean
Indique si cette instance joue actuellement le rôle de broker.

remoteClientCount:

number
Nombre de clients distants connectés à ce broker. Toujours égal à 0 lorsque cette instance n’est pas le broker.

Méthodes
Lien direct vers Méthodes

UnixSocketPubSub implémente le contrat PubSub. La méthode ci-dessous est propre à cette implémentation.

close()
Lien direct vers close

Ferme la connexion au socket et, lorsque cette instance est le broker, libère ce rôle. Appelez cette méthode lors d’un arrêt progressif.

await pubsub.close()

Élection du broker
Lien direct vers Élection du broker

Le premier processus qui se lie au socket devient le broker et achemine les événements entre tous les clients connectés. Les autres processus se connectent en tant que clients. Lorsque le broker s’arrête, un fichier de verrouillage exclusif sérialise l’élection suivante. Un seul client devient le nouveau broker. Les clients restants se réabonnent auprès de lui. Ce mécanisme évite un état de split-brain dans lequel deux processus joueraient simultanément le rôle de broker.