Skip to Content
PortuguêsDocumentação JumentixPacotes@jumentix/dead-letter-queue

@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.

ClientUserServiceMutexServiceDeadLetterQueueupdate(user-1, payload)lock(User, user-1)previouslyLocked = trueResourceLockedErrorenqueue(User, user-1, update, payload)ResourceLockedErrorwithout the queue — the attempt disappearedwith the queue — recorded, and still refused

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

apps/backend-template@jumentix/dead-letter-queueUserService12 rejection sitescomposeUserDeadLetterReplayhandlers + wiringDeadLetterQueueenqueue · pending · replayDeadLetterReplayWorkerinterval drainInMemoryStoreKeyValueStoreRedisrejectLocked()handlers call backget · set · del
PeçaResponsabilidade
DeadLetterQueueRegistos, estados, limite de tentativas, ordem de replay
DeadLetterReplayWorkerChama replay num intervalo, sem sobreposição
InMemoryDeadLetterStoreLocal ao processo; testes e runtimes de processo único
KeyValueDeadLetterStoreRedis, através do cliente que o mutex já usa
composeUserDeadLetterReplayMapeia nomes de operação de volta para métodos do UserService

3. O ciclo de vida do registo

pendingsucceededabandonedenqueue()handler returnedbound reachedhandler threw, under the boundTerminal records are kept, not deleted: “did that write ever happen” must stay answerable.
EstadoSignificadoEscolhido pelo replay?
pendingRecusada, reexecutávelsim
succeededUm handler executou a escritanão
abandonedLimite de tentativas atingidonunca

4. O replay passa pelo serviço

Esta é a decisão com maior probabilidade de ser mal feita, por isso vale ser explícito.

ReplayWorkerhandlerorThrowUserServicere-acquires the mutexlock clearedstatus = succeededstill lockedattempts + 1, stays pendingThrough the service, never the repository: a repository-level replay would write past the lock.

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çãoOmissãoNotas
storeInMemoryDeadLetterStoreUsar KeyValueDeadLetterStore num runtime real
maxAttempts5Mínimo 1; uma fila que nunca tenta é recusada na construção
now() => new Date()Injetado para testes determinísticos
newIdtimestamp + contadorO contador é o que evita colisão entre dois registos no mesmo milissegundo
MétodoDevolve
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? })
MembroComportamento
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
runningSe 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

StoreUso
InMemoryDeadLetterStoreTestes, 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.

dlq:indexid-1, id-2, id-3dlq:record:id-1dlq:record:id-2dlq:record:id-31st2nd3rd

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

SintomaO que significaO que fazer
pending() cresceLocks não libertam, ou o worker não arrancouVerificar worker.running e o TTL do mutex
Registos chegam a abandonedUm recurso ficou bloqueado durante maxAttempts drenagensLer lastError; a escrita perdeu-se e o registo é a prova
skipped num relatórioUma operação não tem handler registadoErro de ligação — o registo é mantido, não descartado
onError disparaA própria drenagem falhou, por exemplo Redis em baixoO 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 real

A 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.

Last updated on