Aller au contenu principal

LeaseProvider

LeaseProvider est le contrat de bail distribué, distinct de la distribution des événements (PubSub). La couche Signals de Mastra l’utilise pour désigner, entre plusieurs processus — par exemple des invocations serverless — un propriétaire unique d’une ressource, le plus souvent une clé de fil de discussion. Le propriétaire est le processus qui se réveille et exécute le flux de l’Agent ; les autres processus lui acheminent donc les tâches de suivi au lieu de démarrer une exécution concurrente.

La gestion des baux est distincte de pub/sub. Un backend n’implémente LeaseProvider que lorsqu’il peut réellement coordonner un verrou, comme Redis au moyen d’opérations SET/Lua atomiques, ou une map en mémoire pour un processus unique. Les backends incapables de gérer des baux l’omettent ; l’environnement d’exécution de Signals détecte cette capacité et utilise à défaut un fournisseur sans effet, ce qui préserve le comportement à processus unique.

L’implémentation intégrée RedisStreamsPubSub implémente LeaseProvider, ce qui permet aux Signals d’assurer la coordination entre les instances dans les déploiements distribués et serverless.

Exemple d’utilisation
Lien direct vers Exemple d’utilisation

Vous ne construisez pas directement un LeaseProvider. Configurez un backend pub/sub qui l’implémente, tel que RedisStreamsPubSub, dans le constructeur Mastra ; l’environnement d’exécution de Signals l’utilise alors automatiquement pour coordonner les processus.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
// RedisStreamsPubSub implements both PubSub and LeaseProvider
pubsub: new RedisStreamsPubSub({
url: process.env.REDIS_URL,
}),
})

Pour gérer les baux dans un backend personnalisé, implémentez les méthodes ci-dessous. L’environnement d’exécution de Signals détecte cette capacité de manière structurelle — elle fonctionne donc au-delà des limites des packages — et l’utilise uniquement lorsque toutes les méthodes sont présentes.

src/mastra/pubsub.ts
import { PubSub } from '@mastra/core/events'
import type { LeaseProvider } from '@mastra/core/events'

export class CustomPubSub extends PubSub implements LeaseProvider {
async acquireLease(key: string, owner: string, ttlMs: number) {
// Atomically claim the lease, or report the current holder.
return { acquired: true, owner }
}

// ...getLeaseOwner, releaseLease, renewLease, transferLease
}

Méthodes
Lien direct vers Méthodes

Gestion des baux
Lien direct vers Gestion des baux

acquireLease(key, owner, ttlMs)
Lien direct vers acquireleasekey-owner-ttlms

Tente d’acquérir atomiquement un bail sur une clé. Renvoie { acquired: true, owner } si l’appelant a obtenu le bail, ou { acquired: false, owner }, où owner est le détenteur actuel, afin que l’appelant puisse lui acheminer les tâches de suivi. Un même propriétaire peut appeler acquireLease de manière idempotente pour renouveler ou reprendre le bail.

const result = await pubsub.acquireLease('thread:abc', runId, 15000)

if (result.acquired) {
// This process owns the thread, so wake and run the agent.
} else {
// result.owner holds the lease, so route the signal to them.
}

Renvoie : Promise<{ acquired: boolean; owner?: string }>

key:

string
Clé du bail, par exemple une clé de fil de discussion.

owner:

string
Identifiant du propriétaire, tel qu’un runId. Un même propriétaire peut appeler acquireLease de manière idempotente pour renouveler ou libérer le bail.

ttlMs:

number
Durée de vie du bail, en millisecondes.

getLeaseOwner(key)
Lien direct vers getleaseownerkey

Lit le propriétaire actuel d’un bail, ou undefined si aucun bail n’est détenu.

const owner = await pubsub.getLeaseOwner('thread:abc')

Renvoie : Promise<string | undefined>

releaseLease(key, owner)
Lien direct vers releaseleasekey-owner

Libère un bail. Cette opération est sans effet si l’appelant n’est pas le propriétaire actuel : les implémentations vérifient atomiquement la propriété avant la libération, de sorte qu’un renouvellement simultané par un autre propriétaire n’est jamais écrasé.

await pubsub.releaseLease('thread:abc', runId)

Renvoie : Promise<void>

renewLease(key, owner, ttlMs)
Lien direct vers renewleasekey-owner-ttlms

Renouvelle un bail existant détenu par owner et prolonge sa durée de vie. Renvoie true si le renouvellement a réussi et que l’appelant détient toujours le bail, ou false si le bail a été perdu, parce que sa durée de vie a expiré ou qu’un autre propriétaire l’a obtenu.

const stillOwned = await pubsub.renewLease('thread:abc', runId, 15000)

if (!stillOwned) {
// Lost the lease, so stop renewing and let the new owner take over.
}

Renvoie : Promise<boolean>

transferLease(key, fromOwner, toOwner, ttlMs)
Lien direct vers transferleasekey-fromowner-toowner-ttlms

Transfère atomiquement un bail détenu de fromOwner vers toOwner et actualise sa durée de vie sans libérer la clé entre les deux opérations. Cette primitive sans interruption permet à un propriétaire de suivi de reprendre immédiatement la même clé lorsque le propriétaire actuel termine. Par exemple, une exécution de suivi en file d’attente peut prendre le relais à la fin de l’exécution d’un fil de discussion. Une approche naïve consistant à libérer puis à acquérir laisse brièvement la clé vide ; un processus concurrent pourrait alors obtenir le bail libéré et démarrer une exécution concurrente.

Renvoie true si fromOwner détenait toujours le bail et que la propriété a été transférée à toOwner, ou false si le bail était déjà perdu. Dans ce dernier cas, l’appelant doit effectuer un nouvel appel à acquireLease.

const transferred = await pubsub.transferLease('thread:abc', currentRunId, nextRunId, 15000)

if (!transferred) {
// Lease was lost, so acquire fresh instead.
await pubsub.acquireLease('thread:abc', nextRunId, 15000)
}

Renvoie : Promise<boolean>

attention

Les backends qui ne peuvent pas effectuer le transfert de manière atomique doivent tout de même l’implémenter au mieux sous la forme d’un appel à releaseLease(fromOwner) suivi de acquireLease(toOwner), et préciser que l’échange n’est pas atomique, puisqu’un processus concurrent peut obtenir la clé pendant l’intervalle. Le maintien de cette méthode comme obligatoire offre aux appelants un chemin de code unique et fait de l’atomicité une décision explicite propre à chaque backend.

Détection des capacités
Lien direct vers Détection des capacités

L’environnement d’exécution de Signals détecte LeaseProvider de manière structurelle plutôt qu’avec instanceof. La détection fonctionne ainsi même lorsqu’un backend publié séparément résout une autre copie de @mastra/core. Une valeur est considérée comme un LeaseProvider lorsqu’elle expose les cinq méthodes (acquireLease, getLeaseOwner, releaseLease, renewLease, transferLease).

Lorsque le backend pub/sub configuré n’implémente pas LeaseProvider, l’environnement d’exécution utilise à défaut un fournisseur sans effet qui réussit toujours. Chaque appelant remporte sa propre course au bail, tandis que la libération, le renouvellement et le transfert sont inertes, ce qui préserve le comportement attendu à processus unique.

  • PubSub : contrat de distribution des événements, distinct de la gestion des baux
  • RedisStreamsPubSub : backend intégré qui implémente LeaseProvider
  • Signals : environnement d’exécution qui utilise les baux pour coordonner l’exécution des fils de discussion entre les processus
  • Channels : utilise les baux pour coordonner les exécutions d’Agents dans les déploiements serverless et multi-instances