@jumentix/dead-letter-queue
Fila privada para transações recusadas pelo mutex.
Uma escrita em fila não aconteceu. Tudo neste pacote decorre dessa frase. O serviço continua a lançar
ResourceLockedError; a fila apenas torna a tentativa recuperável.
1. O problema
O UserService protege cada mutação com o mutex. Quando o recurso já está
bloqueado, recusa:
const { result: { previouslyLocked } } = await this.mutexService.lock(this.entityName, id);
if (previouslyLocked) await this.rejectLocked('update', id, data);São doze pontos assim. Antes deste pacote, o chamador recebia um erro e a transação pretendida era descartada. Nada registava que tinha sido tentada, por isso uma escrita perdida por contenção era indistinguível de uma escrita que ninguém chegou a fazer.
O cliente é recusado nos dois caminhos. É esse o ponto: a escrita não aconteceu, e um cliente informado do contrário agiria sobre uma mentira.
2. As peças
| Peça | Responsabilidade |
|---|---|
DeadLetterQueue | Registos, estados, limite de tentativas, ordem de replay |
DeadLetterReplayWorker | Chama replay num intervalo, sem sobreposição |
InMemoryDeadLetterStore | Local ao processo; testes e runtimes de processo único |
KeyValueDeadLetterStore | Redis, através do cliente que o mutex já usa |
composeUserDeadLetterReplay | Mapeia nomes de operação de volta para métodos do UserService |
3. O ciclo de vida do registo
| Estado | Significado | Escolhido pelo replay? |
|---|---|---|
pending | Recusada, reexecutável | sim |
succeeded | Um handler executou a escrita | não |
abandoned | Limite de tentativas atingido | nunca |
4. O replay passa pelo serviço
Esta é a decisão com maior probabilidade de ser mal feita, por isso vale ser explícito.
Duas coisas que este diagrama defende:
Pelo serviço, não pelo repositório. O serviço readquire o mutex, por isso um registo cujo lock não libertou é recusado outra vez e continua em fila. Um replay ao nível do repositório escreveria por cima do lock, que é exatamente a corrupção que o mutex existe para prevenir.
orThrow não é decoração. O UserService reporta falha em
response.error e não lança. Um handler que ignorasse isso retornaria
normalmente para um registo ainda bloqueado, a fila marcaria succeeded, e a
escrita perder-se-ia com o relatório a dizer que passou. Todos os handlers passam
por orThrow.
5. API
DeadLetterQueue
new DeadLetterQueue({ store?, maxAttempts?, now?, newId? })| Opção | Omissão | Notas |
|---|---|---|
store | InMemoryDeadLetterStore | Usar KeyValueDeadLetterStore num runtime real |
maxAttempts | 5 | Mínimo 1; uma fila que nunca tenta é recusada na construção |
now | () => new Date() | Injetado para testes determinísticos |
newId | timestamp + contador | O contador é o que evita colisão entre dois registos no mesmo milissegundo |
| Método | Devolve |
|---|---|
enqueue(input) | O registo pending guardado |
pending() | Registos reexecutáveis, por ordem de recusa |
find(id) | Um registo, qualquer que seja o estado |
replay(handlers) | { replayed, retried, abandoned, skipped } — ids, não contagens |
O enqueue recusa entrada que não poderia ser reexecutada: entityName,
resourceId e operation são obrigatórios.
DeadLetterReplayWorker
new DeadLetterReplayWorker({ queue, handlers, intervalMs?, onReport?, onError? })| Membro | Comportamento |
|---|---|
start() | Inicia o intervalo (omissão 30_000 ms). O timer leva unref |
stop() | Termina; nenhum tick corre depois |
tick() | Força uma drenagem — no shutdown, ou ao libertar um recurso |
running | Se o intervalo está ativo |
Três propriedades, cada uma um modo de um ciclo de timer ingénuo correr mal:
- Sem sobreposição. Uma drenagem mais lenta que o intervalo não recomeça enquanto a anterior corre, ou o mesmo registo é reexecutado duas vezes em simultâneo.
- Uma drenagem que falha não mata o worker. Redis em baixo é um tick mau, não o fim. Um worker que morre ao primeiro erro é indistinguível de um que nunca arrancou.
- Parar é completo, incluindo um tick já agendado.
Stores
| Store | Uso |
|---|---|
InMemoryDeadLetterStore | Testes, runtimes de processo único |
KeyValueDeadLetterStore(client, { prefix }) | Redis, prefixo por omissão dlq |
O cliente Redis que este repositório usa expõe get, set e del — sem
SCAN, sem KEYS. Por isso a store mantém um índice explícito dos ids dos
registos sob uma chave, que é também o que preserva a ordem de replay.
6. Ligação
Passar a fila como serviço. Sem ela, o comportamento é exatamente o anterior.
import { composeUserDeadLetterQueue, composeUserDeadLetterWorker }
from '@src/modules/Users/composition/composeUserDeadLetterReplay';
const deadLetterQueue = composeUserDeadLetterQueue(keyValueStorageClient);
const userService = UserService.compile({
dataRepository,
services: { mutexService, passwordCryptoService, deadLetterQueue }
});
const deadLetterWorker = composeUserDeadLetterWorker(deadLetterQueue, userService);O composeUsersAuthServices já faz isto e devolve ambos. Constrói o worker mas
não o arranca: um timer em fundo é do runtime arrancar e, mais importante,
parar.
deadLetterWorker?.start();
process.on('SIGTERM', () => deadLetterWorker?.stop());Sem keyValueStorageClient ambos ficam undefined, de propósito: uma fila local
ao processo perder-se-ia no restart parecendo durável.
7. Experimentar
Um script executável, sem Redis:
bun run --filter @jumentix/dead-letter-queue exampleÉ o examples/replay.ts deste pacote — edita e volta a correr. Percorre o ciclo
todo: uma escrita recusada, um replay que falha porque o lock ainda se mantém, um
replay que passa, e um registo que atinge o limite de tentativas.
import { DeadLetterQueue } from '@jumentix/dead-letter-queue';
const queue = new DeadLetterQueue({ maxAttempts: 2 });
await queue.enqueue({
entityName: 'User', resourceId: 'user-1', operation: 'update', payload: { firstName: 'Ada' }
});
let locked = true;
const handlers = {
update: async () => { if (locked) throw new Error('User user-1 is locked'); }
};
await queue.replay(handlers); // retried: [ 'dlq-...' ]
locked = false;
await queue.replay(handlers); // replayed: [ 'dlq-...' ]8. Operação
| Sintoma | O que significa | O que fazer |
|---|---|---|
pending() cresce | Locks não libertam, ou o worker não arrancou | Verificar worker.running e o TTL do mutex |
Registos chegam a abandoned | Um recurso ficou bloqueado durante maxAttempts drenagens | Ler lastError; a escrita perdeu-se e o registo é a prova |
skipped num relatório | Uma operação não tem handler registado | Erro de ligação — o registo é mantido, não descartado |
onError dispara | A própria drenagem falhou, por exemplo Redis em baixo | O worker continua a tiquetaquear; corrigir a store |
O relatório devolve ids, não contagens, para que um operador possa procurar o registo em vez de inferir de um total.
9. Validação
bun run --filter @jumentix/dead-letter-queue build
bun run --filter @jumentix/dead-letter-queue typecheck
bun run --filter @jumentix/dead-letter-queue lint
bun run --filter @jumentix/dead-letter-queue test
bun run smoke:dead-letter:redis # contra um contentor Redis realA suite de integração salta-se a si própria sem RUN_REDIS_INTEGRATION=1 em
vez de passar. Uma suite que reporta sucesso sem a sua dependência é o
falso-verde que este repositório continua a encontrar.
O que só o servidor real responde, e que essa suite portanto verifica: valores voltam como strings e fazem parse, uma segunda fila sobre o mesmo servidor vê o que a primeira escreveu, o índice mantém a ordem de recusa, o worker drena, e um registo abandonado continua abandonado depois de reabrir.