• Joined on 2026-08-12

@synvectis/queue (0.2.2)

Published 2026-09-06 13:09:48 +02:00 by neigel

Installation

@synvectis:registry=
npm install @synvectis/queue@0.2.2
"@synvectis/queue": "0.2.2"

About this package

@synvectis/queue

La file de travaux différés d'une application Synvectis.

Deux mécanismes, deux métiers

⚠ Il ne faut pas les confondre, et c'est tout le dessin de ce paquet.

createQueue createLimiter
Chez qui L'émetteur — celui qui a quelque chose à perdre Le service qui calcule
Ce qu'il fait Persiste, rejoue, ordonne, trace Compte ce qui est en vol, refuse au-delà
Ce qu'il détient Des travaux, dans Mongo Rien
Survit à un redémarrage Oui — il vit dans Mongo, pas dans le processus Sans objet
Ce qu'il sait Qu'un travail ne doit pas se perdre Que le GPU tient un modèle à la fois

Les mettre au même endroit, c'est soit donner au service un état qu'il n'a pas à tenir — et il cesse d'être remplaçable —, soit demander à l'émetteur de deviner une capacité matérielle qu'il ignore.

Et ils s'emboîtent : un service saturé répond 429, que la file d'en face traite comme un échec passager, donc rejouable avec son recul. Rien à inventer de part et d'autre.

Côté émetteur

import { createQueue, register, PermanentFailure } from '@synvectis/queue'

register('envoi-courriel', {
  queue: 'io',
  maxAttempts: 5,
  handler: async ({ payload }) => envoyer(payload),
})

const queue = createQueue({
  connection: mongoose.connection,
  redis,
  queues: { io: { concurrency: 20 }, compute: { concurrency: 2 } },
  onEvent: (event, job) => socket.emit(event, job),   // facultatif
})
queue.start()

// Un battement : le même type rejoué, sans document et sans doublon même à
// plusieurs répliques.
await queue.schedule('training-poll', 5_000)

await queue.submit('envoi-courriel', {
  subject: { tenant: orgId, actor: userId, scope: projetId },
  payload: { destinataire, gabarit },
  dedupKey: `courriel:${destinataire}:${gabarit}`,
  priority: 100,          // ⚠ PLUS PETIT = PLUS PRIORITAIRE. Défaut : 400
  source: 'facturation',  // pour la trace, jamais pour décider
})

⚠ La charge (payload) est une RÉFÉRENCE, jamais un contenu. Elle est persistée, listable et recopiée dans les sauvegardes : un fichier, un secret ou une donnée personnelle n'y ont pas leur place. Ce qu'il faut pour agir se résout à l'exécution.

⚠ Le sujet est trois chaînes opaques, sans aucune ref Mongoose. C'est ce qui rend ce paquet utilisable par une application qui n'a ni les mêmes collections ni la même notion de projet. Une chaîne ne se joint pas.

⚠ onEvent ne peut pas faire échouer un travail. L'appel est enveloppé : un abonné qui lève — socket fermée, sérialisation impossible — ferait sinon rater un travail qui a réussi, et le ferait rejouer. Un observateur ne doit pas pouvoir provoquer ce qu'il observe.

⚠ Un battement n'a pas de document. schedule inscrit une horloge, pas un travail : rien à lister, à rejouer ou à annuler. Son jobId est stable, donc trois répliques qui démarrent n'accélèrent pas la cadence — ce qu'un setInterval fait, lui, par construction.

⚠ PermanentFailure ne se rejoue pas. Un droit qui manque ne réapparaîtra pas en réessayant ; le rejouer noierait la vraie cause sous des tentatives.

Sans Redis

redis est facultatif. Sans lui, le déclenchement se fait par sondage Mongo — aucune infrastructure à ajouter :

const queue = createQueue({
  connection: mongoose.connection,
  queues: { default: { concurrency: 2 } },
  pollMs: 1_000,   // le pas de sondage
  lockMs: 30_000,  // la durée d'une réservation
})

⚠ Le paquet n'embarque pas Redis, il s'en passe. Redis est un serveur, pas une bibliothèque : aucun paquet npm n'en glisse un dans un processus, et les imitations en mémoire ne sont partagées par aucune autre instance — deux répliques auraient chacune sa file, et le même travail partirait deux fois. Or la durabilité n'a jamais été dans Redis : elle est dans Mongo.

⚠ Ce que Redis apporte vraiment, et c'est le seul critère du choix :

Avec Redis Par sondage Mongo
Départ d'un travail quelques millisecondes jusqu'à pollMs
Coût au repos nul (attente bloquante) une requête indexée par classe et par tour
Durabilité la même — elle est dans Mongo la même
Sûreté en multi-instance la même la même (prise atomique)
Infrastructure un serveur de plus aucune

Pour vingt travaux par jour, le sondage ne se voit pas. Pour vingt par seconde, Redis redevient le bon outil.

⚠ La prise de travail est un findOneAndUpdate atomique, et la réservation est renouvelée tant que le travail tourne. Sans ce renouvellement, tout travail plus long que lockMs serait repris par un autre ouvrier — donc exécuté deux fois.

Côté service

import { createLimiter, Saturated } from '@synvectis/queue'

const limiter = createLimiter(Number(process.env.CALCULS_SIMULTANES) || 2)

app.post('/api/calcul/:methode', async (req, res, next) => {
  try {
    res.json(await limiter.run(() => calculer(req.params.methode, req.body)))
  } catch (err) {
    if (err instanceof Saturated) return res.status(err.status).json({ error: err.message })
    next(err)
  }
})

⚠ Pas de salle d'attente, délibérément. Faire patienter à l'intérieur ferait calculer pour un appelant qui a peut-être déjà abandonné sur son propre délai. On refuse tout de suite ; le rejeu appartient à celui qui détient la durabilité.

Conventions

Code en anglais, commentaires et messages en français. Les décisions qui tiennent pour une raison précise portent un bloc ⚠ qui dit laquelle — souvent le bug qui l'a imposée.

Dependencies

Development Dependencies

ID Version
bullmq ^5.0.0
ioredis ^5.0.0
mongodb-memory-server ^10.1.4
mongoose 9.2.1
typescript ^5.9.3
vitest ^4.1.0

Peer Dependencies

ID Version
bullmq >=5
ioredis >=5
mongoose >=8
Details
npm
2026-09-06 13:09:48 +02:00
70
latest
31 KiB
Assets (1)
Versions (1) View all
0.2.2 2026-09-06