Skip to Content
Jumentix DocsPackages@jumentix/dead-letter-queue

@jumentix/dead-letter-queue

Private queue for transactions the mutex refused.

A queued write has not happened. Everything in this package follows from that one sentence. The service still throws ResourceLockedError; the queue only makes the attempt recoverable.


1. The problem

UserService guards every mutation with the mutex. When the resource is already locked, it refuses:

const { result: { previouslyLocked } } = await this.mutexService.lock(this.entityName, id);
if (previouslyLocked) await this.rejectLocked('update', id, data);

There are twelve such sites. Before this package, the caller got an error and the intended transaction was discarded. Nothing recorded that it had been attempted, so a write lost to contention was indistinguishable from a write nobody ever made.

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

The client is refused in both paths. That is the point: the write has not happened, and a client told otherwise would act on a lie.


2. The pieces

apps/backend-template@jumentix/dead-letter-queueUserService12 rejection sitescomposeUserDeadLetterReplayhandlers + wiringDeadLetterQueueenqueue · pending · replayDeadLetterReplayWorkerinterval drainInMemoryStoreKeyValueStoreRedisrejectLocked()handlers call backget · set · del
PieceResponsibility
DeadLetterQueueRecords, statuses, attempt bound, replay ordering
DeadLetterReplayWorkerCalls replay on an interval, without overlapping
InMemoryDeadLetterStoreProcess-local; tests and single-process runtimes
KeyValueDeadLetterStoreRedis, through the client the mutex already uses
composeUserDeadLetterReplayMaps operation names back to UserService methods

3. The record lifecycle

pendingsucceededabandonedenqueue()handler returnedbound reachedhandler threw, under the boundTerminal records are kept, not deleted: “did that write ever happen” must stay answerable.
StatusMeaningPicked by replay?
pendingRefused, replayableyes
succeededA handler performed the writeno
abandonedAttempt bound reachedno, ever

4. Replay goes through the service

This is the decision most likely to be got wrong, so it is worth being explicit.

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

Two things this diagram is defending:

Through the service, not the repository. The service re-acquires the mutex, so a record whose lock has not cleared is refused again and stays queued. A repository-level replay would write past the lock, which is exactly the corruption the mutex exists to prevent.

orThrow is not decoration. UserService reports failure in response.error and does not throw. A handler that ignored that would return normally for a still-locked record, the queue would mark it succeeded, and the write would be lost while the report said it landed. Every handler goes through orThrow.


5. API

DeadLetterQueue

new DeadLetterQueue({ store?, maxAttempts?, now?, newId? })
OptionDefaultNotes
storeInMemoryDeadLetterStoreUse KeyValueDeadLetterStore in a real runtime
maxAttempts5Must be at least 1; a queue that never attempts is refused at construction
now() => new Date()Injected for deterministic tests
newIdtimestamp + counterThe counter is why two records in the same millisecond do not collide
MethodReturns
enqueue(input)The stored pending record
pending()Replayable records, in rejection order
find(id)One record, whatever its status
replay(handlers){ replayed, retried, abandoned, skipped } — ids, not counts

enqueue refuses input that could not be replayed: entityName, resourceId and operation are all required.

DeadLetterReplayWorker

new DeadLetterReplayWorker({ queue, handlers, intervalMs?, onReport?, onError? })
MemberBehaviour
start()Begins the interval (default 30_000 ms). The timer is unrefed
stop()Ends it; no tick runs afterwards
tick()Forces one drain — on shutdown, or on the unlock of a resource
runningWhether the interval is active

Three properties, each a way a naive timer loop goes wrong:

  • No overlap. A drain slower than the interval is not started again while the previous one runs, or the same record is replayed twice concurrently.
  • A failing drain does not stop the worker. Redis down is a bad tick, not the end. A worker that dies on the first error is indistinguishable from one that was never started.
  • Stopping is complete, including a tick already scheduled.

Stores

StoreUse
InMemoryDeadLetterStoreTests, single-process runtimes
KeyValueDeadLetterStore(client, { prefix })Redis, prefix defaults to dlq

The Redis client this repository uses exposes get, set and del — no SCAN, no KEYS. The store therefore keeps an explicit index of record ids under one key, which is also what preserves replay order.

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

6. Wiring

Pass the queue as a service. Without it, behaviour is exactly as before.

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);

composeUsersAuthServices already does this and returns both. It builds the worker but does not start it: a background timer is the runtime’s to start and, more importantly, to stop.

deadLetterWorker?.start();
process.on('SIGTERM', () => deadLetterWorker?.stop());

Without a keyValueStorageClient both are undefined, on purpose: a process-local queue would be lost on restart while looking like durability.


7. Try it

A runnable script, no Redis required:

bun run --filter @jumentix/dead-letter-queue example

It is examples/replay.ts in this package — edit it and re-run. It walks the whole lifecycle: a refused write, a replay that fails because the lock still holds, a replay that succeeds, and a record that reaches the attempt bound.

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. Operating it

SymptomWhat it meansWhat to do
pending() growsLocks are not clearing, or the worker is not startedCheck worker.running and the mutex TTL
Records reach abandonedA resource stayed locked for maxAttempts drainsRead lastError; the write is lost and the record is the evidence
skipped in a reportAn operation has no registered handlerA wiring mistake — the record is kept, not discarded
onError firesThe drain itself failed, for example Redis downThe worker keeps ticking; fix the store

The report returns ids, not counts, so an operator can look a record up rather than infer from a total.


9. Validation

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 # against a real Redis container

The integration suite skips itself without RUN_REDIS_INTEGRATION=1 rather than passing. A suite that reports success without its dependency is the false green this repository keeps finding.

What only the real server can answer, and what that suite therefore asserts: values come back as strings and parse, a second queue over the same server sees what the first wrote, the index holds rejection order, the worker drains, and an abandoned record is still abandoned after a reopen.

Last updated on