@emergence/queue (0.1.0)
Installation
@emergence:registry=npm install @emergence/queue@0.1.0"@emergence/queue": "0.1.0"About this package
@emergence/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 '@emergence/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.
Côté service
import { createLimiter, Saturated } from '@emergence/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 |