@cqrs-ddd/pipeline-deadletter
Captures a failed operation as a dead letter: its request, error, tenant and correlation
id, sent to a transport (BullMQ, RabbitMQ or PostgreSQL). By default only events are
captured, since the caller of a command or a query already receives the error.
DeadLetterRedriver runs a captured request again later.
Installation
Section titled “Installation”pnpm add @cqrs-ddd/pipeline-deadletter @cqrs-ddd/pipelineEach transport takes a client of its broker or database (bullmq queue, amqplib
confirm channel, pg pool), which the application installs.
import { createPipeline } from '@cqrs-ddd/pipeline';import { DeadLetterBehavior, deadLetter, PostgresDeadLetterTransport,} from '@cqrs-ddd/pipeline-deadletter';
const transport = new PostgresDeadLetterTransport(pool);const pipeline = createPipeline({ behaviors: [new DeadLetterBehavior(transport)] });
export const onUserCreated = pipeline.wrap( { name: 'onUserCreated', kind: 'event' }, deadLetter({ rethrow: false, redactKeys: ['refreshToken'] }),)(async (event: UserCreatedEvent) => mailer.sendWelcome(event.email));The behavior’s constructor takes the transport, optional defaults for every operation and an optional logger.
Transports
Section titled “Transports”| Transport | Notes |
|---|---|
BullMqDeadLetterTransport |
adds each dead letter as a job of a BullMQ queue (job name dead-letter by default) |
RabbitMqDeadLetterTransport |
publishes to RabbitMQ through a confirm channel |
PostgresDeadLetterTransport |
inserts into a table, dead_letters by default; it is also a DeadLetterStore, which the redriver needs |
Another backend implements DeadLetterTransport: send(record).
Options
Section titled “Options”| Option | Meaning | Default |
|---|---|---|
captureKinds |
request kinds that are captured | ['event'] |
rethrow |
rethrow the error after capturing it; false swallows it, on events only |
true |
ignoreErrors |
error classes, or a predicate, that are rethrown without being captured, such as validation errors | none |
redactKeys |
field names masked in the captured payload | none |
redact |
a function that replaces the payload redaction | none |
metadata |
a function that returns extra fields for the record | none |
includeStack |
include the stack trace | true |
ignoreErrors and redactKeys given both as constructor defaults and per operation are
combined; the other options are merged shallowly.
Ordering
Section titled “Ordering”Place the behavior outside retries, so it captures a failure only once the retries are exhausted, and inside validation, so expected validation errors are not captured.
Redriving
Section titled “Redriving”import { DeadLetterRedriver } from '@cqrs-ddd/pipeline-deadletter';
const redriver = new DeadLetterRedriver(transport, { requestTypes: [UserCreatedEvent], dispatch: { // Only the handler that failed, not every subscriber of the event. event: (event, record) => eventHandlers[record.handlerName](event), },});
await redriver.redrive(recordId);The redriver rebuilds the request from the record, matching record.requestName against
the class names in requestTypes, dispatches it, and marks the record resolved. When the
handling fails again, it counts the attempt and rethrows. A record whose payload was
redacted needs a rebuild function that restores the redacted values. During a redrive,
the behavior neither captures the failure again nor swallows it.