Samples
00-showcase
sample/00-showcase
is a runnable, end-to-end demonstration of the whole pattern on SQLite and
the in-process transport — the no-broker default profile — with no migration
step (it creates the tables inline so it runs as a single script):
schema.ts— the library'soutboxEvents/inboxEventsfactories combined with the businessorderstable and anorder_audittable (the consumer's observable side effect).app.module.ts— a global Drizzle module,ClsModule.forRootwith the Drizzle transactional adapter (enableTransactionProxy: true), andMessagingModule.forRootwithSqliteOutboxStore,SqliteInboxStore, and anInProcessOutboxTransportover a sharedOutboxRegistry.order.service.ts—placeOrderinserts the order row andenqueues theorder.placedevent in the same@Transactional()method. The payload is a plain interface —enqueueaccepts it without casts.order-placed.handler.ts— the consumer: it registers itself fororder.placedon module init and pairs withInboxService.runOnce(keyed onidempotencyKey ?? id) so the audit row is written exactly once.scripts/smoke.ts— drives the flow and asserts each guarantee.
What it proves
The smoke script asserts the properties that make the pattern correct:
- Atomic outbox — after
placeOrder, both theordersrow and exactly oneoutbox_eventsrow exist. They committed in the same transaction, so there is no event without the work and no work without the event. - Claim and dispatch — one
OutboxClaimer.tick()delivers the committed event to the registered handler and marks the row completed (report.completed === 1). In production with a broker, the Kafka transport publishes instead — same seam, different transport. - Exactly-once inbox — the handler's first delivery writes one
order_auditrow; replaying the same logical event (delivery is at-least-once) is deduplicated by the inbox and writes no second row. - Unroutable events fail fast — an event on a topic with no registered
handler maps to
PermanentError, so the claimer fails the row immediately instead of retrying forever.
On success it prints:
Showcase smoke passed: atomic outbox → in-process dispatch → exactly-once inbox.
Running it
From the repository root:
npm install
npm run showcase
The showcase runs the in-process default profile — no Docker, no Kafka. Rebind
KafkaOutboxTransport and a thin @KafkaConsumer (see the
Quick Start) to take the same flow to production.
01-kafka
sample/01-kafka
takes the showcase one step further: it drives the whole pair over the real
Kafka transport — KafkaOutboxTransport on the producer side and an actual
@KafkaConsumer delegating to KafkaInboxConsumer on the consumer side — using
@nest-native/kafka's
in-memory broker (KafkaTestModule), so it still runs with no cluster.
app.module.ts— wiresKafkaTestModule.forRoot()andMessagingModule.forRootAsync({ ..., useTransport: (producer) => new KafkaOutboxTransport(producer) }).order.consumer.ts— a thin@KafkaConsumer('order.placed')that delegates to the library'sKafkaInboxConsumer.consume(...), supplying the payload validator and the exactly-once side effect.scripts/smoke.ts— places an order, runsOutboxClaimer.tick()(which publishes through Kafka to the consumer), asserts one audit row, then re-emits the same message to prove the inbox deduplicates the redelivery.
Run it from the repository root with npm run sample:focused. This is the closest
you can get to the production path without a broker; point KafkaTestModule at a
real cluster (or use KafkaModule) and the same code runs unchanged.
02-rabbitmq
sample/02-rabbitmq
runs the pair over a real RabbitMQ broker: RabbitOutboxTransport on the
producer side and RabbitInboxConsumer on billing's queue. It is the single
service to read first — every path the adapter promises is asserted against
the broker, not simulated.
rabbitmq.ts— the application owns itsamqplibconnection (connect(url, { recovery: true })) and declares its topology at startup; the transport and the inbox only open channels on it.topology.ts— a topic exchange for events, a quorum work queue with a dead-letter exchange and a dead-letter routing key (without it, a message over the delivery limit is dead-lettered with its original key, misses the dead-letter queue's binding, and is dropped), and the dead-letter queue.billing.consumer.ts— subscribes with manual acks and a prefetch, hands every delivery toRabbitInboxConsumer, and dead-letters poison on a confirm channel. When a channel closes under it — a dropped connection, a channel the broker closed, a cancelled consumer — it subscribes again on new ones.scripts/smoke.ts— places an order and asserts:- Confirmed publish — the claimer's row completes only after the broker acked the message and did not return it; billing issues one invoice.
- Duplicate — the same event published again (what a redelivery after a lost ack looks like) is acked as a duplicate and bills nothing.
- Poison — an invalid payload lands in the dead-letter queue with its
reason in an
x-errorheader, instead of being requeued forever. - Unroutable —
order.refunded, which no queue binds yet, comes back to the transport; the attempt fails and the event stays in the outbox for a retry rather than being acked into the void. - Dropped connection — the broker closes the application's connection (through the management API, as a restart would). amqplib reconnects, the transport opens a new confirm channel, billing subscribes again on its own, and the next order is invoiced once.
03-rabbitmq-services
sample/03-rabbitmq-services
is the pattern doing the job it exists for: two services, orders and
shipping, each a separate process with its own database file, choreographed
over one broker.
orderstakes an order and publishesorder.placedthrough its outbox.shippingconsumes it through its inbox and, in the same transaction, books a shipment and enqueuesshipment.scheduledin its own outbox — the dedup row, the shipment and the outgoing event commit together, so consuming once and publishing once is one step.ordersconsumesshipment.scheduledand marks the order scheduled.- Each service owns its queue (
shared/topology.ts) and runs its own outbox relay (runWorkerLoop), stopped inbeforeApplicationShutdownso it never publishes on a connection that is closing. shared/inbox-subscription.tskeeps each consumer subscribed: whenever its channel closes, it subscribes again on a new one, a second later.
The smoke script forks both services and asserts: three orders cross both
services; while shipping is down, orders keeps accepting and publishing
orders and RabbitMQ keeps them in shipping's durable queue; a new shipping
process on the same database works through that backlog; the broker then
drops both services' connections, and without a restart both reconnect,
subscribe again and carry the next orders through; an event delivered
again books no second shipment and publishes no second
shipment.scheduled; and both outboxes drain, one event per order.
Why one process per service: @nestjs-cls/transactional keeps its transaction
host in process-global state keyed by connection name, and the outbox and inbox
use the default connection. Two Nest applications started in one process —
in a test, say — share one transaction host, so one of them writes through the
other's database. Give each application its own process.
Running the RabbitMQ samples
Both need a broker, and its management API for the dropped-connection step. From the repository root:
npm run infra:up
RABBITMQ_URL=amqp://messaging:messaging@127.0.0.1:56720 \
RABBITMQ_MANAGEMENT_URL=http://messaging:messaging@127.0.0.1:15670 \
npm run sample:focused
Without RABBITMQ_URL they skip locally with a notice, and without
RABBITMQ_MANAGEMENT_URL they skip only the dropped-connection step. In CI,
the sample jobs run a RabbitMQ service container, and a missing URL fails the
job instead.