@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.
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
| Piece | Responsibility |
|---|---|
DeadLetterQueue | Records, statuses, attempt bound, replay ordering |
DeadLetterReplayWorker | Calls replay on an interval, without overlapping |
InMemoryDeadLetterStore | Process-local; tests and single-process runtimes |
KeyValueDeadLetterStore | Redis, through the client the mutex already uses |
composeUserDeadLetterReplay | Maps operation names back to UserService methods |
3. The record lifecycle
| Status | Meaning | Picked by replay? |
|---|---|---|
pending | Refused, replayable | yes |
succeeded | A handler performed the write | no |
abandoned | Attempt bound reached | no, ever |
4. Replay goes through the service
This is the decision most likely to be got wrong, so it is worth being explicit.
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? })| Option | Default | Notes |
|---|---|---|
store | InMemoryDeadLetterStore | Use KeyValueDeadLetterStore in a real runtime |
maxAttempts | 5 | Must be at least 1; a queue that never attempts is refused at construction |
now | () => new Date() | Injected for deterministic tests |
newId | timestamp + counter | The counter is why two records in the same millisecond do not collide |
| Method | Returns |
|---|---|
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? })| Member | Behaviour |
|---|---|
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 |
running | Whether 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
| Store | Use |
|---|---|
InMemoryDeadLetterStore | Tests, 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.
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 exampleIt 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
| Symptom | What it means | What to do |
|---|---|---|
pending() grows | Locks are not clearing, or the worker is not started | Check worker.running and the mutex TTL |
Records reach abandoned | A resource stayed locked for maxAttempts drains | Read lastError; the write is lost and the record is the evidence |
skipped in a report | An operation has no registered handler | A wiring mistake — the record is kept, not discarded |
onError fires | The drain itself failed, for example Redis down | The 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 containerThe 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.