Djaouad

Reliable Domain Events and the Transactional Outbox: a Hybrid Event-Driven Pipeline in a DDD + Onion Architecture Backend

October 9, 202656 min read
typescriptevent-driventransactional-outboxdomain-eventsbullmqpostgresidempotencydddonion-architecturebackendnodejs

One outbox table, two polling workers, BullMQ queues, registries and idempotency keys, all of it built as ports and adapters on top of the same DDD + Onion Architecture codebase.

This post builds on Dependency Injection from Scratch, the container, scopes and composition roots described there are used here without re-explaining them. It's also the machinery behind the event driven RAG pipeline of Building a Shopping Assistant AI Agent from Scratch. also, the full application is in this repository


when an order gets created, saving it is only half of the job, the system also has to send a confirmation email, create the shipment in the shipping provider's API, update the embeddings knowledge base,... etc. the naive way is to do those right after the save() call, and that's where a reliable system quietly turns into an unreliable one: if the process crashes between the commit and the publish, the order exists but nobody was ever told, and if we call the shipping API inside the transaction, the remote call can't be rolled back when the commit fails.

This post is about closing that gap with one durable pipeline, and about how a single outbox table can carry two different kinds of side effects.


Part I, The Problem

1.1 Two things happen when an order is created

  • the 1st is a state change in the database (new order created, inventory updated,... etc.), this can be saved atomically, many entities in one transaction, if one write fails they all fail, which keeps our system consistent.
  • the 2nd is side effects (send a confirmation email, create the order in the shipping API, update the embeddings knowledge base,... etc.), these are risky to do inside the same transaction because they fail and succeed independently of the DB write, and they could be heavy tasks that need async processing (publishing domain events to an external queue, for example).

The side effects introduce a new potential inconsistency point: we can't atomically write to our database and call external systems. Imagine a use case that does:

  1. save the new state to the DB
  2. create the order in the shipping API
  3. publish the domain events to a queue
await this.db.transaction(async (tx) => {
  await this.orderRepository.save(order, tx);
  await this.cartRepository.save(cart, tx);

  await Promise.all(products.map((p) => this.productRepository.save(p, tx)));
});

await this.shippingProviderGateway.createShipment(order);

await this.eventPublisher.publish(events);
Crash pointResult
after 1, before 2- order is saved in DB
- our system is inconsistent with the shipping API's database
- all domain events are lost (order created, stock updated,... etc.)
after 2, before 3- order is saved in DB
- our system is consistent with the shipping API's database
- all domain events are lost (order created, stock updated,... etc.)

And moving the shipping call inside the transaction doesn't fix it, it just creates a third failure mode: the remote call succeeds, then the commit fails (a deadlock, a constraint, a dropped connection), and now the shipping provider knows about an order that doesn't exist in our database. A remote call can't be rolled back. On top of that, the transaction holds a pooled DB connection open for the whole duration of an HTTP call we don't control.

1.2 The thesis

We have two kinds of side effects:

  1. actions that must be taken for the use case to be considered successful (like syncing the order state between our database and the shipping API), these represent "to do" items for our system, and each one requires executing exactly one specific workflow (one external API request).
  2. events that must reach the right queues to notify other parts of the system that something significant to the business happened (order created, user registered, inventory updated,... etc.), these could trigger many different workflows throughout the system (one order.created event could trigger a confirmation email, an inventory update, an analytics update,... etc.)

Different in meaning, but identical in what they need from the infrastructure: be recorded durably in the same transaction as the state change, then be delivered later, at least once, to whoever needs them. This post builds one pipeline for both.


Part II, First Principles

2.1 Domain events as internal system notifications

Domain events are past tense, immutable objects collected by aggregates in-memory at runtime. they represent something that already happened, an event significant to the business, that could trigger different side effects. domain events are usually emitted whenever an aggregate's state is mutated through business logic.

export class Order {
  private _events: DomainEvent[] = [];
  // rest of the aggregate attributes, methods,... etc.

  confirm(): void {
    // Business Logic: update order status to confirmed
    // ...

    // record the event
    this.recordThat(
      new OrderConfirmed(
        this.id.value,
        this.userId.value,
        this._orderItems.length,
        this.getTotalOrderPrice().amount,
        this.getTotalOrderPrice().currency,
        this.getSelectedShippingProvider(),
      ),
    );
  }

  private recordThat(event: DomainEvent): void {
    this._events.push(event);
  }

  pullEvents(): DomainEvent[] {
    const events = [...this._events];
    this._events = [];
    return events;
  }
}

they're drained from the aggregate at the end of each use case to be published.

export class ConfirmOrderService {
  // ...

  async execute(command: ConfirmOrderCommand): Promise<void> {
    const orderId = OrderId.of(command.orderId);

    const order = await this.orderRepository.find(orderId);

    if (!order) throw new NotFoundError("order", orderId.value);

    order.confirm(); // event is recorded here

    const events = order.pullEvents(); // returns all collected events

    // save updated order aggregate to db

    // publish events to queue
  }
}

The last two comments are the whole post: where and how do those two steps happen reliably?

2.2 Delivery

There are two places a domain event can be delivered to its consumers:

  • in-process: a synchronous dispatcher (an in-memory event bus) calls the handlers inside the same process, usually inside the use case. Simple, no infrastructure, but the handlers share the request's fate: a slow handler slows the request, a failing handler fails (or worse, is swallowed by) the use case, and a crash loses everything not yet handled.
  • out-of-process: the event goes to a broker/queue (BullMQ over Redis here), and independent workers consume it. The use case returns fast, consumers retry on their own, and they scale and fail independently.

We want the second, and any out-of-process delivery has to pick a guarantee:

GuaranteeMeaningCost
at-most-oncefire and forget, never retriedevents can be silently lost
at-least-onceretried until acknowledgedthe same event can be delivered more than once
exactly-oncedelivered once and only oncenot achievable end to end over a network, you simulate it

We choose at-least-once delivery + idempotent consumers. delivery never loses an event, and each consumer makes processing the same event twice harmless, which is the closest thing to "exactly-once" that actually exists in practice. Keep this sentence in mind, because every design decision in the rest of the post is a consequence of it.

2.3 External calls inside a DB transaction

Back to the three failure modes from Part I, now stated as a rule: a database transaction can only protect things that live inside that database.

  • publishing to a queue/calling an API after the commit can be lost if the process dies in between.
  • publishing/calling before the commit can't be taken back if the commit fails.
  • doing it inside the transaction body is just "before the commit" with extra steps, plus a held connection while we wait on a network we don't control (pool exhaustion under load is the usual symptom).

So we need a way to make "I must tell the outside world about this" part of the transaction itself.

2.4 How the transactional outbox solves it

The transactional outbox pattern, in one sentence: don't send the message, write it to a table in the same transaction as the state change, and let another process send it later.

The intent to call the outside world becomes ordinary data, and ordinary data in the same transaction is atomic with the rest of the state, either the order and its outbox rows are all committed, or none of them are.

Rendering diagram…

The relay is allowed to fail, crash and retry, because the rows are durable. What it can never do is lose a row, and that's the guarantee we were missing.

2.5 Events vs actions, on one table

Domain eventOutbox action
Meaninga fact: "something happened"a command: "do this"
Exampleorder.created, product.updatedcreate_order_in_shipping_api
Namingpast tenseimperative
Consumerszero, one or many (fan-out to queues)exactly one handler
Typicallynotifies other parts of the systemwrites to an external system
Created byaggregates, drained with pullEvents()the application service explicitly
category columndomain-eventoutbox-job

(in the code the second category is called outbox-job and the entries OutboxJobEntry, in prose i'll call them outbox actions, because that's what they are: actions our system owes the outside world.)


Part III, The Hybrid Pipeline

3.1 Why one shared pipeline works

Look at the lifecycle of both kinds of side effects:

  1. recorded as a row, in the same transaction as the state change.
  2. found by a polling worker, claimed so nobody else processes it.
  3. published to a BullMQ queue.
  4. validated and handled by a worker, idempotently.
  5. cleaned up/recovered by helper workers if something goes wrong.

Every step is the same for events and actions. The only differences are:

  • the category of the row, so each processor polls its own slice.
  • the publish target: an action goes to exactly one queue (outbox-queue), an event is mapped to zero or more queues (email-queue, embedding-queue,... etc.)
  • the handlers: an action has one handler service, an event has one handler service per queue that cares about it.

So instead of two systems we build one table, one repository, one state machine, one recovery story, and two thin processors on top.

3.2 Overview diagram

Rendering diagram…

Everything left of the queues is the relay side (Part V and VI), everything right of them is the consumer side (Part VII and VIII).

3.3 The outbox table

export const outbox = pgTable(
  "outbox",
  {
    id: varchar("id", { length: 40 }).notNull().primaryKey(),
    category: outboxCategoryEnum("category").notNull(),
    // "order.created" for a "domain-event" category
    // "create_order_in_shipping_api" for an "outbox-job" category
    event_type: varchar("event_type", { length: 100 }).notNull(),
    // optional, useful to query all domain events of an aggregate ordered by created_at for debugging
    aggregate_id: varchar("aggregate_id", { length: 80 }),
    payload: jsonb("payload").notNull(),
    status: outboxStatusEnum("status").notNull().default("PENDING"),
    attempts: integer("attempts").notNull().default(0),
    scheduledAt: timestamp("scheduled_at", { withTimezone: true })
      .notNull()
      .defaultNow(),
    processed_at: timestamp("processed_at", { withTimezone: true }),
    error_message: text("error_message"),
    created_at: timestamp("created_at", { withTimezone: true })
      .notNull()
      .defaultNow(),
    locked_at: timestamp("locked_at", { withTimezone: true }),
  },
  (t) => [
    // each worker polls its own slice efficiently
    index("outbox_category_status_scheduled_idx").on(
      t.category,
      t.status,
      t.scheduledAt,
    ),
    index("outbox_aggregate_id_idx").on(t.aggregate_id),
    index("outbox_event_type_idx").on(t.event_type),
  ],
);

export const outboxCategoryEnum = pgEnum("outbox_category", [
  "outbox-job",
  "domain-event",
]);

export const outboxStatusEnum = pgEnum("outbox_status", [
  "PENDING",
  "PROCESSING",
  "COMPLETED",
  "FAILED",
]);
  • category is what makes the table hybrid, both processors read the same table, each one filtering on its own category. This is why the composite index starts with (category, status, scheduled_at): each worker's polling query is category = X and status = 'PENDING' and scheduled_at <= now(), exactly the index's column order.
  • event_type holds the event code or the action name, and payload holds the event itself (for events) or the action's arguments (for actions), as jsonb.
  • aggregate_id is only filled for events, it's not used by the pipeline, it exists so you can answer "what happened to this order?" with one query when debugging.
  • scheduled_at does two jobs: it's the retry backoff (a failed publication is rescheduled into the future) and it lets an action be delayed on purpose (saveJob accepts an optional scheduledAt).
  • attempts counts publication attempts, locked_at records when a row was claimed (the stuck-rows worker uses it), processed_at and error_message are for post-mortems.

One table or two? Two tables would be cleaner to type, but you'd duplicate the state machine, the repository, the stuck-row recovery and the cleanup. One table with a discriminator trades a little type safety for a lot less machinery, and we win the type safety back at the repository boundary (3.4).

And the lifecycle of a row is a small state machine:

Rendering diagram…

3.4 The OutboxRepository port

export const OutboxAction = {
  CREATE_ORDER_IN_SHIPPING_API: "create_order_in_shipping_api",
  DELETE_ORDER_IN_SHIPPING_API: "delete_order_in_shipping_api",
  UPDATE_ORDER_IN_SHIPPING_API: "update_order_in_shipping_api",
  CREATE_SHIPMENT_IN_SHIPPING_API: "create_shipment_in_shipping_api",
} as const;

export type OutboxAction = (typeof OutboxAction)[keyof typeof OutboxAction];

export const OutboxCategory = {
  OUTBOX_JOB: "outbox-job",
  DOMAIN_EVENT: "domain-event",
} as const;

export const OutboxStatus = {
  PENDING: "PENDING",
  PROCESSING: "PROCESSING",
  COMPLETED: "COMPLETED",
  FAILED: "FAILED",
} as const;
// (+ the matching type aliases)

export type OutboxJobEntry = {
  id: string;
  category: typeof OutboxCategory.OUTBOX_JOB;
  eventType: OutboxAction;
  payload: unknown;
  status: OutboxStatus;
  attempts: number;
  scheduledAt: Date;
  processedAt: Date | null;
  errorMessage: string | null;
  createdAt: Date;
  lockedAt: Date | null;
};

export type OutboxDomainEventEntry = {
  id: string;
  category: typeof OutboxCategory.DOMAIN_EVENT;
  eventType: DomainEventCode;
  aggregateId: string | null;
  // ... same bookkeeping fields
};

export type OutboxRepository = {
  // called by application services inside the same transaction as aggregate saves
  saveJob(
    params: {
      action: OutboxAction;
      payload: Record<string, unknown>;
      scheduledAt?: Date; // optional, defaults to now
    },
    tx: TransactionClient,
  ): Promise<void>;

  saveEvents(events: DomainEvent[], tx: TransactionClient): Promise<void>;

  // used by the processor workers
  getPendingJobs(limit: number): Promise<OutboxJobEntry[]>;
  getPendingEvents(limit: number): Promise<OutboxDomainEventEntry[]>;
  updateRowToProcessing(params: UpdateRowToProcessingParams): Promise<boolean>;
  updateRowToPending(params: UpdateRowToPendingParams): Promise<void>;
  updateRowToCompleted(params: UpdateRowToCompletedParams): Promise<void>;
  updateRowToFailed(params: UpdateRowToFailedParams): Promise<void>;

  // used by the helper workers
  deleteCompletedRows(olderThan: Date, tx: TransactionClient): Promise<void>;
  getStuckRows(
    batchSize: number,
    stuckBefore: Date,
    tx?: TransactionClient,
  ): Promise<(OutboxJobEntry | OutboxDomainEventEntry)[]>;
};
  • the single table splits back into two discriminated entry types at the repository boundary: getPendingJobs can only return OutboxJobEntry (with eventType: OutboxAction) and getPendingEvents only OutboxDomainEventEntry (with eventType: DomainEventCode), the varchar column becomes a precise union type, so the processors never deal with an untyped string. (the one cast lives inside the Postgres adapter, where we know the category, and is commented as such.)
  • the tx parameter is required on saveJob/saveEvents, you can't write to the outbox outside a transaction by accident, which is the entire point of the pattern.
  • this port is defined in the application layer instead of the domain layer: if an interface accepts or returns domain entities/aggregates it belongs to the domain, but the outbox exists only to support an application workflow/technical pattern, domain code never sees it. Same reasoning as the idempotency keys repository port.

Part IV, Writing to the Outbox

4.1 Saving domain events

The write side is the simplest part of the whole system, which is the idea. Here's CreateOrderService, trimmed to the relevant part:

cart.clear();

const orderEvents = order.pullEvents();
const productEvents = products.flatMap((p) => p.pullEvents());
const cartEvents = cart.pullEvents();

await this.db.transaction(async (tx) => {
  await this.idempotencyKeysRepository.create(
    idempotencyKey,
    "CreateOrderService",
    tx,
    { orderId: order.id.value },
  );

  await Promise.all([
    this.orderRepository.save(order, tx),
    this.cartRepository.save(cart, tx),
    this.outboxRepository.saveEvents(
      [...orderEvents, ...productEvents, ...cartEvents],
      tx,
    ),
  ]);

  await Promise.all(products.map((p) => this.productRepository.save(p, tx)));
});
  • one use case touches three aggregates (order, cart, products), each one recorded its own events (order.created, cart.cleared, stock.reserved,... etc.), the service drains all of them with pullEvents() and saves them as a flat list.
  • the events are drained after all the business logic ran and before the transaction opens, so the transaction only contains writes.
  • saveEvents maps every event to a row (category: domain-event, status: PENDING, event_type: event.eventType, aggregate_id, and the event object itself as the payload) and inserts them in one statement.
  • anything this service reads from an external gateway is read-only (it validates the shipping price and the postal code against the provider before touching state), reads have no side effects to compensate, so they're fine to do outside of the transaction. only external writes need the outbox.

4.2 Saving outbox actions

Cancelling an order is the interesting case, if the order was already created in the shipping provider's system (it has a tracking number), that shipment has to be deleted there too:

await this.db.transaction(async (tx) => {
  await Promise.all([
    this.orderRepository.save(order, tx),
    this.outboxRepository.saveEvents([...orderEvents, ...productEvents], tx),
  ]);

  await Promise.all(products.map((p) => this.productRepository.save(p, tx)));

  // if the order has a tracking number, it means this order was created in the shipping api,
  // so we schedule a transactional outbox action to delete it, since we don't wanna make
  // external (gateway) write requests inside of a transaction
  if (trackingNumber) {
    const payload: OutboxJobPayloadType<
      typeof OutboxAction.DELETE_ORDER_IN_SHIPPING_API
    > = {
      shippingProvider: order.getSelectedShippingProvider(),
      trackingNumber,
    };

    await this.outboxRepository.saveJob(
      {
        action: OutboxAction.DELETE_ORDER_IN_SHIPPING_API,
        payload,
      },
      tx,
    );
  }
});
  • OutboxJobPayloadType<typeof OutboxAction.DELETE_ORDER_IN_SHIPPING_API> types the payload per action, you can't schedule a delete without a tracking number and a provider, the compiler refuses (Part VII shows where this type comes from).
  • the action is conditional, so the same transaction can write zero, one or many outbox rows depending on the business state at that moment.
  • notice what didn't happen: the cancel use case returns without ever talking to the shipping provider. if the provider is down for an hour, cancelling still works instantly, and the delete is retried until it goes through.

4.3 One transaction, three writes

Look at the cancel transaction as a whole: aggregates (order and products), domain events and an outbox action, all in one db.transaction. Either all three kinds of writes are committed, or none of them is.

That's the guarantee that closes the table from 1.1, there is no crash point left between "state changed" and "the outside world will be told": the row that will tell the outside world is part of the state change. After the commit, whatever happens to the process, the intent survives in Postgres. How reliably it gets delivered from there is the relay's job, the next part.


Part V, Relaying: the Polling Workers

5.1 Worker anatomy

Every background process in this system (the two processors and the two helpers) has the same shape: a class with start(), stop(), a private loop() and a public runIteration().

export class OutboxProcessorWorker {
  private running = false;
  private loopPromise: Promise<void> | null = null;
  private abortController = new AbortController();

  constructor(
    private container: Container,
    private options: OutboxProcessorWorkerOptions = defaultOutboxProcessorWorkerOptions,
  ) {}

  /** Runs exactly one iteration. Exposed directly so tests can call it
   *  without going through the infinite loop. */
  async runIteration(): Promise<void> {
    const iterationId = `iter_${crypto.randomUUID().replace(/-/g, "")}`;

    await runWithContext(
      { requestId: iterationId, startTime: performance.now() },
      async () => {
        const scope = this.container.createScope();

        try {
          const service = scope.resolve(OUTBOX_PROCESSOR_SERVICE);
          await service.execute(
            new OutboxProcessorCommand(
              this.options.maxPublicationAttempts,
              this.options.batchSize,
            ),
          );
        } finally {
          await scope.dispose();
        }
      },
    );
  }

  start(): void {
    if (this.running) return;
    this.running = true;
    this.abortController = new AbortController(); // fresh controller per run
    this.loopPromise = this.loop();
  }

  async stop(): Promise<void> {
    this.running = false;
    this.abortController.abort();
    await this.loopPromise;
  }

  private async loop(): Promise<void> {
    while (this.running) {
      try {
        await this.runIteration();
      } catch (error) {
        logger.error("Outbox processor iteration failed", error as Error);
        await sleep(this.options.sleepAfterFailMs, this.abortController.signal);
        continue;
      }

      await sleep(this.options.pollIntervalMs, this.abortController.signal);
    }
  }
}
  • it's a polling loop, not a cron job: run an iteration, sleep, repeat. the practical difference is that iterations can never overlap (the next one starts only after the previous one finished), a slow batch can't pile up concurrent runs of itself.
  • the sleep takes the abort signal, so stop() wakes the loop immediately instead of waiting out a 3 minute sleep, and then awaits loopPromise so the current iteration finishes cleanly. that's what makes SIGTERM handling in the entrypoint graceful.
  • a failed iteration sleeps a short sleepAfterFailMs instead of the full poll interval, so a transient failure (DB hiccup) recovers quickly.
  • each iteration gets its own scope and correlation id, same as §18 of the DI post, and runIteration() being public means tests run exactly one iteration without any timers.
  • the entrypoint is just composition + configuration + signal handling:
const container = buildOutboxProcessorContainer();

const outboxProcessorWorker = new OutboxProcessorWorker(container, {
  pollIntervalMs: 1000 * 60 * 3, // 3 minutes
  sleepAfterFailMs: 5000,
  maxPublicationAttempts: 5,
  batchSize: 40,
});

outboxProcessorWorker.start();

const shutdown = async (signal: string) => {
  logger.info(`Received ${signal}, starting graceful shutdown...`);
  try {
    await outboxProcessorWorker.stop();
    process.exit(0);
  } catch (error) {
    logger.error("Error during graceful shutdown", error as Error);
    process.exit(1);
  }
};

process.on("SIGTERM", () => shutdown("SIGTERM"));
process.on("SIGINT", () => shutdown("SIGINT"));

The poll interval is the main latency knob of the whole system: with 3 minutes, an email can be up to 3 minutes late. Lower it (even to a second) if you need it, a poll is one indexed query that returns nothing most of the time. (see the limitations part for the alternatives.)

5.2 Publishing outbox actions

export class OutboxProcessorService {
  constructor(
    private outboxRepository: OutboxRepository,
    private outboxQueue: Queue,
  ) {}

  async execute(command: OutboxProcessorCommand) {
    const pendingJobs = await this.outboxRepository.getPendingJobs(
      command.batchSize,
    );

    if (pendingJobs.length === 0) return;

    for (const job of pendingJobs) {
      const attempt = job.attempts + 1;

      // 1. claim
      try {
        const claimed = await this.outboxRepository.updateRowToProcessing({
          id: job.id,
          attempts: attempt,
        });

        if (!claimed) continue; // somebody else got it
      } catch (error) {
        this.logger.error("Failed to claim outbox job", error as Error, { id: job.id, attempt });
        continue;
      }

      // 2. publish
      try {
        await this.outboxQueue.add(job.eventType, job.payload, {
          jobId: job.id,
          attempts: 5, // BullMQ execution retry configuration (independent of DB)
          backoff: { type: "exponential", delay: 2000 },
        });
      } catch (error) {
        try {
          await this.handlePublishFailure(job.id, attempt, error, command);
        } catch (error) {
          // job is not published and the DB still says PROCESSING,
          // the stuck rows worker will reset it later
        }
        continue;
      }

      // 3. mark completed
      try {
        await this.outboxRepository.updateRowToCompleted({
          id: job.id,
          processedAt: new Date(),
        });
      } catch (error) {
        // job IS published but the DB still says PROCESSING, the stuck rows worker
        // will reset it and it will be published again: a duplicate delivery,
        // harmless because consumers are idempotent (at-least-once semantics)
      }
    }
  }
}

Three steps per row, and each one has its own failure story, which is why they have separate try blocks:

1. claim. updateRowToProcessing is a compare-and-set:

const [updatedRow] = await this.db
  .update(outbox)
  .set({
    status: OutboxStatus.PROCESSING,
    attempts: params.attempts,
    locked_at: new Date(),
  })
  .where(
    and(
      eq(outbox.id, params.id),
      eq(outbox.status, OutboxStatus.PENDING),
    ),
  )
  .returning();

return !!updatedRow;

the where status = 'PENDING' clause makes the update itself the lock: if two processors read the same row, Postgres lets exactly one update match it, the other gets false and skips the row. no explicit lock, no read-then-write race. (the same trick as claiming a conversation in the agent post.) this is why running more than one replica of a processor is safe, even though we run one.

2. publish. the row's id becomes the BullMQ jobId. BullMQ ignores an add whose jobId already exists in the queue, so a republished row can't create a second job while the first is still there, a cheap first line of defense, but not the real one (see Part VII, completed jobs are removed from Redis, only the idempotency table remembers).

If add fails (Redis down, timeout), we don't lose anything, we decide what to do with the row:

private async handlePublishFailure(jobId, attempt, error, command) {
  if (attempt >= command.maxPublicationAttempts) {
    await this.outboxRepository.updateRowToFailed({
      id: jobId,
      processedAt: new Date(),
      errorMessage: error instanceof Error ? error.message : "Unknown error",
    });
  } else {
    const retryDelayMs = Math.pow(2, attempt) * 1000;
    const nextRetryAt = new Date(Date.now() + retryDelayMs);

    await this.outboxRepository.updateRowToPending({
      id: jobId,
      scheduledAt: nextRetryAt,
      errorMessage: error instanceof Error ? error.message : "Unknown error",
    });
  }
}

the row goes back to PENDING with scheduled_at pushed into the future (2s, 4s, 8s,... exponential), and after maxPublicationAttempts it becomes FAILED and stays in the table for a human to look at.

3. mark completed. if this update fails, the job was published but the row says PROCESSING. we do nothing and let the stuck rows worker (Part VI) reset it, which means it will be published again. that's the at-least-once guarantee showing its price, and exactly why the consumers must be idempotent.

there are two independent retry layers in the system, and they're easy to mix up:

DB publication attemptsBullMQ execution attempts
Protects againstRedis unreachable when publishingthe handler failing when running
Counted inoutbox.attemptsBullMQ job attemptsMade
LimitmaxPublicationAttempts (5)attempts: 5 in the job options
Backoffscheduled_at in the futureBullMQ exponential backoff
After the limitrow becomes FAILEDthe job fails and is removed from the queue

5.3 Publishing domain events

DomainEventsProcessorService is the same three steps with two differences, it reads the other category, and instead of queue.add it calls a port:

const pendingJobs = await this.outboxRepository.getPendingEvents(command.batchSize);

// ... claim, exactly as before

await this.eventPublisher.publish(job.eventType, job.payload, job.id);

// ... on failure: handlePublishFailure, on success: updateRowToCompleted

and the polling query it relies on:

this.db.query.outbox.findMany({
  where: and(
    eq(outbox.category, OutboxCategory.DOMAIN_EVENT),
    eq(outbox.status, OutboxStatus.PENDING),
    lte(outbox.scheduledAt, new Date()),
  ),
  orderBy: (outbox, { asc }) => [asc(outbox.scheduledAt), asc(outbox.created_at)],
  limit,
})
  • scheduled_at <= now() is what makes backoff work: a row waiting out its retry delay simply isn't visible to the poll yet.
  • the service doesn't know BullMQ exists, it depends on EventPublisher, an application port. the action processor does hold a BullMQ Queue directly because an action always goes to exactly one known queue, there's nothing to route. an event needs a routing decision, and routing belongs behind a port.

5.4 The EventPublisher port + BullMqEventPublisher

export type EventPublisher = {
  publish(
    eventType: DomainEventCode,
    payload: unknown,
    jobId: string,
  ): Promise<void>;
};
export class BullMqEventPublisher implements EventPublisher {
  constructor(
    private flowProducer: FlowProducer,
    private emailQueue: Queue,
    private inventoryQueue: Queue,
    private analyticsQueue: Queue,
    private embeddingQueue: Queue,
  ) {}

  async publish(eventType: DomainEventCode, payload: unknown, jobId: string) {
    const queues = this.mapEventToQueues(eventType);

    const jobsArray: FlowJob[] = queues.map((queue) => ({
      name: eventType,
      data: payload,
      queueName: queue.name,
      opts: {
        jobId, // unique per queue
        attempts: 5,
        backoff: { type: "exponential", delay: 2000 },
      },
    }));

    if (jobsArray.length === 0) return;

    await this.flowProducer.addBulk(jobsArray);
  }

  private mapEventToQueues(eventType: DomainEventCode): Queue[] {
    const eventToQueuesMapper: Record<DomainEventCode, Queue[]> = {
      "order.created": [this.emailQueue, this.inventoryQueue, this.analyticsQueue],
      "order.cancelled": [this.emailQueue, this.inventoryQueue, this.analyticsQueue],
      "order.confirmed": [this.emailQueue, this.analyticsQueue],
      "product.created": [this.embeddingQueue],
      "product.updated": [this.embeddingQueue],
      "product.deleted": [this.embeddingQueue],
      // ... every other event maps to [] for now
    };

    return eventToQueuesMapper[eventType] ?? [];
  }
}
  • the event → queues mapper is the routing table of the whole event system, one place that answers "who cares about this event?". it's typed as Record<DomainEventCode, Queue[]>, so adding a new event code to the domain breaks compilation here until you decide where it goes ([] is a valid, explicit answer: nobody cares yet).
  • the same jobId is used for every queue, that's fine for BullMQ because job ids are unique per queue, and it means one outbox row can be traced through all the queues it fanned out to.
  • FlowProducer.addBulk is used for the fan-out, as far as i know it adds the whole batch in one Redis MULTI, so an event is either enqueued to all its queues or to none, instead of landing in the email queue and failing before the analytics one (a plain queue.add loop would leave exactly that half-published state, and the retry would double-publish the first queues).
  • an event mapped to [] still goes through the whole pipeline and ends COMPLETED, publishing to zero queues is a successful publication. the day you add a consumer, only the mapper changes.

and the composition root that wires it all up, buildDomainEventsProcessorContainer:

container.register(EMAIL_QUEUE, (scope) => createBullMqEmailQueue(scope.resolve(REDIS)), "singleton");
container.register(EMBEDDING_QUEUE, (scope) => createBullMqEmbeddingQueue(scope.resolve(REDIS)), "singleton");
// ... analytics + inventory queues, the flow producer

container.register(
  EVENT_PUBLISHER,
  (scope) =>
    new BullMqEventPublisher(
      scope.resolve(BULLMQ_FLOW_PRODUCER),
      scope.resolve(EMAIL_QUEUE),
      scope.resolve(INVENTORY_QUEUE),
      scope.resolve(ANALYTICS_QUEUE),
      scope.resolve(EMBEDDING_QUEUE),
    ),
  "singleton",
);

container.register(
  DOMAIN_EVENTS_PROCESSOR_SERVICE,
  (scope) =>
    new DomainEventsProcessorService(
      scope.resolve(OUTBOX_REPOSITORY),
      scope.resolve(EVENT_PUBLISHER),
    ),
  "scoped",
);

this is the smallest-graph idea from the DI post again: the process that publishes events knows about the queues and the outbox repository, it has never heard of CreateOrderService, nor of the shipping gateway.

5.5 The honest duplication

OutboxProcessorService and DomainEventsProcessorService are about 95% the same code: claim, publish, handle failure, mark completed. i kept them as two classes on purpose, the duplicated part is a state machine that's easy to read, and the differences (a queue vs a port, jobs vs events) would turn into configuration flags in a shared version. if a third category ever shows up, that's the moment to extract a generic relay that takes fetchPending and publish as functions.


Part VI, Failure Cases and Helper Workers

6.1 The failure matrix

The whole design can be judged by one question: what happens when the process dies right here?

Crash pointRow stateWho recovers itNet effect
during the business transactionnothing was written (rolled back)no one neededthe use case failed cleanly, nothing leaked
after the commit, before any worker pollsPENDINGthe next polldelivered, just late
after the claim, before the publishPROCESSINGthe stuck rows workerrow reset, then published
Redis down during the publishPENDING with a future scheduled_atthe next polls (backoff)delivered after Redis returns
Redis down for all 5 attemptsFAILEDa humannot delivered, visible in the table
after the publish, before marking COMPLETEDPROCESSINGthe stuck rows workerpublished again: a duplicate, absorbed by consumer idempotency
the handler throws(row is already COMPLETED)BullMQ retriesretried 5 times with backoff

Two things to notice. The only way to lose a side effect is the FAILED row, and it's loud (it sits in the table with its error message). And the only way to duplicate one is the "published but not marked" row, which is the price of at-least-once, paid for by idempotent consumers.

6.2 Reset stuck rows

a row in PROCESSING for too long means its processor died between the claim and the end of the iteration. the fix is a third polling worker, same anatomy as the others, with its own options:

export const defaultResetStuckOutboxRowsWorkerOptions = {
  pollIntervalMs: 1000 * 60 * 60 * 24, // 1 day
  sleepAfterFailMs: 5000,
  stuckFor: 1000 * 60 * 60 * 24, // 1 day
  batchSize: 100,
};

// inside runIteration
const stuckBefore = new Date(Date.now() - this.options.stuckFor);
await service.execute(new ResetStuckOutboxRowsCommand(this.options.batchSize, stuckBefore));

and the service is almost nothing:

const stuckRows = await this.outboxRepository.getStuckRows(
  command.batchSize,
  command.stuckBefore,
);

for (const row of stuckRows) {
  try {
    await this.outboxRepository.updateRowToPending({
      id: row.id,
      scheduledAt: row.scheduledAt,
      errorMessage: row.errorMessage ?? "row stuck",
    });
  } catch (error) {
    // next poll cycle will try again this row
    this.logger.error("Failed to reset stuck row", error as Error, { row });
    continue;
  }
}
  • getStuckRows is status = 'PROCESSING' and locked_at <= stuckBefore, for both categories, that's the payoff of the shared table: one recovery worker for events and actions.
  • resetting only flips the status back to PENDING, the normal processor picks the row up on its next poll, we don't publish from the helper, there's still exactly one place that publishes.
  • the threshold is a trade-off, too short and you reset rows that are merely slow (a duplicate), too long and a crash delays delivery. it's a duplicate-vs-latency dial, and duplicates are cheap here, so it's safe to err on the short side compared to the 1 day default above.

6.3 Clean outbox

without cleanup the table grows forever, and the polling queries slow down with it.

export const defaultCleanOutboxWorkerOptions = {
  pollIntervalMs: 1000 * 60 * 60 * 24, // 1 day
  sleepAfterFailMs: 5000,
  retentionMs: 5 * 24 * 60 * 60 * 1000, // 5 days
};

// inside runIteration
const cutoff = new Date(Date.now() - this.options.retentionMs);
await service.execute(new CleanOutboxCommand(cutoff));
await this.db.transaction(async (tx) => {
  await this.outboxRepository.deleteCompletedRows(command.olderThan, tx);
});
db.delete(outbox).where(
  and(
    eq(outbox.status, OutboxStatus.COMPLETED),
    lte(outbox.processed_at, olderThan),
  ),
)
  • only COMPLETED rows are deleted, and only after a retention period (5 days), so you keep a few days of history to debug with (aggregate_id + created_at makes "what happened to order X?" a single query).
  • FAILED rows are never deleted automatically, they're the table's dead letter queue, you want a human to see them.

6.4 FAILED rows

and speaking of which: right now nothing consumes the FAILED rows, they sit there until someone queries them. that's fine to start with, but the honest minimum for production is a metric or an alert on count(*) where status = 'FAILED', plus a way to push a row back to PENDING once the underlying problem is fixed (it's one update, since the processor already knows how to handle a PENDING row).


Part VII, Consuming Outbox Actions

7.1 outbox-queue and OutboxHandlerWorker

Everything up to now ended with a job sitting in a BullMQ queue. The consumer side is a BullMQ Worker per queue, with a scope per job (§19 of the DI post):

this.worker = new Worker(
  "outbox-queue",
  async (job) => {
    return runWithContext(
      { requestId: `job_${job.id}`, jobId: job.id ?? "unknown", queueName: "outbox-queue", startTime: performance.now() },
      async () => {
        const container = this.buildContainer();
        const scope = container.createScope();

        const jobId = job.id;
        const outboxAction = job.name as OutboxAction;

        if (!jobId) throw new BadRequestError("jobId is required");

        try {
          // 1. schema lookup + validation
          const payloadSchema = outboxJobPayloadsSchemas.shape[outboxAction];

          if (!payloadSchema) {
            throw new ValidationError("outboxAction", `Invalid outbox action: ${outboxAction}`);
          }

          const payload = payloadSchema.parse(job.data) as OutboxJobPayloadType<typeof outboxAction>;

          // 2. build typed command
          const command = buildOutboxCommand(outboxAction, payload);

          // 3. resolve service & execute (fully typed end-to-end)
          await executeOutboxHandler(outboxAction, scope, command, jobId);
        } catch (error) {
          this.logger.error("Outbox job failed", error as Error, { jobId, outboxAction });
          throw error; // BullMQ handles retries
        } finally {
          await scope.dispose();
        }
      },
    );
  },
  { connection: this.connection, concurrency: 3, lockDuration: 30000, stalledInterval: 30000 },
);

the job's name is the action (delete_order_in_shipping_api), its data is the payload we stored in the row, and the whole handling is three steps: validate, build a command, execute. the worker contains no business logic and no knowledge of any specific action, the next sections plug the knowledge in.

  • validate first. the payload crossed a serialization boundary (Postgres jsonb → Redis → JSON), and rows can outlive the code that wrote them (a deploy can change a schema while old rows are still pending), so we treat it as untrusted input and parse it with zod before it reaches any service.
  • rethrow so BullMQ retries. the worker never swallows an error, a throw is how BullMQ learns the job failed, and it re-runs it with the exponential backoff set when we published (attempts: 5).
  • concurrency: 3 means three jobs run interleaved in the same process, three scopes, zero shared state.

7.2 A service per action, plus idempotency

each action has exactly one handler service, here's the one that calls the external API:

export class DeleteOrderFromShippingProviderService {
  constructor(
    private db: DBClient,
    private shippingProviderGateway: ShippingProviderGateway,
    private idempotencyKeysRepository: IdempotencyKeysRepository,
  ) {}

  async execute(command: DeleteOrderFromShippingProviderCommand, jobId: string) {
    try {
      await this.db.transaction(async (tx) => {
        // create the idempotency key first with this jobId to make sure it wasn't
        // successfully processed before
        await this.idempotencyKeysRepository.create(
          jobId,
          "DeleteOrderFromShippingProviderService",
          tx,
        );

        const { success } = await this.shippingProviderGateway.deleteUnshippedShipment(
          command.trackingNumber,
        );

        if (!success)
          throw new GatewayError(
            "shippingProviderGateway",
            new Error("Failed to delete order from shipping provider"),
          );
      });
    } catch (error) {
      this.logger.error("Error deleting order from shipping provider", error as Error, { jobId });
      throw error; // so the worker knows the job failed
    }
  }
}

and the table behind that create:

export const idempotencyKeys = pgTable(
  "idempotency_keys",
  {
    id: varchar("id", { length: 40 }).notNull(), // BullMQ jobId / outbox row id
    handler_name: varchar("handler_name", { length: 100 }).notNull(),
    created_at: timestamp("created_at").notNull().defaultNow(),
    payload: jsonb("payload"),
  },
  (t) => [primaryKey({ columns: [t.id, t.handler_name] })],
);

how this works, step by step:

  1. the idempotency key is the job id, which is the outbox row id, the one identity that's stable across every republish and every BullMQ retry of the same row.
  2. inside one transaction we insert the key first. the primary key makes the insert fail if this handler already processed this job: a duplicate delivery throws on the insert, before the external call, and the work isn't repeated.
  3. then we call the gateway. if it fails, we throw, the transaction rolls back the key, and BullMQ's retry finds no key and tries again. if it succeeds, the transaction commits the key.
  4. two concurrent deliveries of the same job can't both pass step 2: the second insert waits for the first transaction, then fails on the key (or sees it committed).

the key is (id, handler_name) and not just id because one outbox row can fan out to several queues with the same job id (Part VIII), each handler needs its own key, otherwise the second handler would see the first one's key and silently skip its work.

now the part we should clarify: this handler makes an external call inside a database transaction, which is exactly what Part II said not to do. it's a deliberate trade-off, not an oversight:

  • the rule from 2.3 is about the business transaction, where the cost of a failed commit is a ghost shipment. here the transaction contains one row (the key), the failure we protect against is different: running the action twice.
  • the cost is real: a pooled connection is held while the HTTP call runs, and there's still one window left, the gateway call succeeds, then the commit fails, and the retry repeats the call. no database trick closes that window, because the other side of the call isn't in our database.
  • the way to close it for real is on the provider side: make the action itself idempotent (deleting an already deleted shipment must count as success, or send the provider an idempotency key if it supports one). idempotency keys in your DB reduce duplicates, they don't make a non-idempotent remote call safe.

the same table also guards the API itself: CreateOrderService stores the client's idempotencyKey with the created orderId as payload, and looks it up at the start of the request (parseOrderIdFromPayload), so a client retrying a timed-out POST /orders gets the original order id back instead of a second order. one table, two uses: dedupe requests at the edge, dedupe jobs at the consumers.

7.3 The registry pattern and job validation

the worker from 7.1 is generic, it handles every action through three lookups. let's build them, in order.

1. the payload schemas, one zod schema per action, in one object:

export const outboxJobPayloadsSchemas = z.object({
  [OutboxAction.CREATE_ORDER_IN_SHIPPING_API]: z.object({ orderId: z.string() }),
  [OutboxAction.CREATE_SHIPMENT_IN_SHIPPING_API]: z.object({ trackingNumber: z.string() }),
  [OutboxAction.UPDATE_ORDER_IN_SHIPPING_API]: z.object({ orderId: z.string() }),
  [OutboxAction.DELETE_ORDER_IN_SHIPPING_API]: z.object({
    trackingNumber: z.string(),
    shippingProvider: z.enum(ShippingProvider),
  }),
});

export type OutboxJobPayloadType<T extends OutboxAction> = z.infer<
  typeof outboxJobPayloadsSchemas
>[T];

this is the single source of truth: the worker does outboxJobPayloadsSchemas.shape[action].parse(...) at runtime, and the services that write actions use OutboxJobPayloadType<typeof OutboxAction.X> at compile time (that's the type from 4.2). the writer and the reader of a payload can't drift apart.

2. the command builder, which turns a validated payload into the typed command object the service expects:

export function buildOutboxCommand<T extends OutboxAction>(
  action: T,
  payload: OutboxJobPayloadType<T>,
): OutboxActionToCommand[T] {
  switch (action) {
    case OutboxAction.DELETE_ORDER_IN_SHIPPING_API: {
      const p = payload as OutboxJobPayloadType<typeof OutboxAction.DELETE_ORDER_IN_SHIPPING_API>;

      return new DeleteOrderFromShippingProviderCommand(
        p.trackingNumber,
        p.shippingProvider,
      ) as OutboxActionToCommand[T];
    }
    // ... one case per action

    default: {
      const _exhaustive: never = action;
      throw new Error(`Unhandled outbox action: ${_exhaustive}`);
    }
  }
}

the never in the default branch is the exhaustiveness check: add a new value to OutboxAction and this function stops compiling until you handle it.

3. the handler registry, mapping each action to the token of its service and how to call it:

type OutboxActionToCommand = {
  [OutboxAction.CREATE_ORDER_IN_SHIPPING_API]: CreateOrderInShippingProviderCommand;
  [OutboxAction.DELETE_ORDER_IN_SHIPPING_API]: DeleteOrderFromShippingProviderCommand;
  // ...
};

type OutboxActionToHandlerService = {
  [OutboxAction.DELETE_ORDER_IN_SHIPPING_API]: DeleteOrderFromShippingProviderService;
  // ...
};

type OutboxActionToToken = {
  [OutboxAction.DELETE_ORDER_IN_SHIPPING_API]: typeof DELETE_ORDER_FROM_SHIPPING_PROVIDER_SERVICE;
  // ...
};

type HandlerRegistryEntry<T extends OutboxAction> = {
  token: OutboxActionToToken[T];
  handlerMethod: (
    handler: OutboxActionToHandlerService[T],
    command: OutboxActionToCommand[T],
    jobId: string,
  ) => Promise<unknown>;
};

const handlerRegistry: { [K in OutboxAction]: HandlerRegistryEntry<K> } = {
  [OutboxAction.DELETE_ORDER_IN_SHIPPING_API]: {
    token: DELETE_ORDER_FROM_SHIPPING_PROVIDER_SERVICE,
    async handlerMethod(handler, command, jobId) {
      return handler.execute(command, jobId);
    },
  },
  // ... one entry per action
};

4. the executor, the one function the worker calls:

export async function executeOutboxHandler<T extends OutboxAction>(
  action: T,
  scope: Scope,
  command: OutboxActionToCommand[T],
  jobId: string,
): Promise<unknown> {
  const entry = handlerRegistry[action];
  const service = scope.resolve<OutboxActionToHandlerService[T]>(
    entry.token as InjectionToken<OutboxActionToHandlerService[T]>,
  );
  return entry.handlerMethod(service, command, jobId);
}

what the mapped type { [K in OutboxAction]: HandlerRegistryEntry<K> } buys us: every action must have an entry, and each entry is checked against its own command, service and token types. a handler method that takes the wrong command for its action doesn't compile.

so adding a new outbox action is a checklist the compiler walks you through:

  1. add the name to OutboxAction.
  2. add its payload schema to outboxJobPayloadsSchemas (the other places now demand it).
  3. write the command class and the handler service (with the idempotency key).
  4. add its cases to the three OutboxActionTo... maps and to buildOutboxCommand.
  5. add its registry entry, and register the service token in the outbox handler composition root.
  6. call saveJob from the use case that owes the outside world something.

one honest note on the type safety: the generic functions use as casts (as OutboxActionToCommand[T], as InjectionToken<...>), TypeScript can't prove that a switch on action narrows the generic T, so these are the unchecked seams of the registry. they're small and contained in one file, but they're real, a wrong command class in the right case still compiles.


Part VIII, Consuming Domain Events

8.1 Queue per consumer concern, plus its worker

the action side has one queue because every action has one handler. events are different, one event can matter to several independent concerns, so we create one queue per consumer concern: email-queue, embedding-queue, inventory-queue, analytics-queue.

export function createBullMqEmailQueue(connection: Redis) {
  return new Queue("email-queue", {
    connection,
    defaultJobOptions: {
      attempts: 3,
      backoff: { type: "exponential", delay: 1000 },
      priority: 0,
      removeOnComplete: true,
      removeOnFail: true,
    },
  });
}

(the queue factories are identical apart from the name, and the publisher overrides attempts/backoff per job.)

why a queue per concern instead of one events queue that everyone reads:

  • failure isolation. if the email provider is down, the email queue backs up and retries, while embeddings and analytics keep flowing.
  • independent scaling and concurrency. emails are cheap and numerous, embeddings call an LLM API and are slow, they deserve different worker counts and rate limits.
  • independent deployment. a consumer concern is its own process (Part X), you can ship, stop or restart the email handler without touching the rest.
  • fan-out for free. the same event lands in as many queues as the mapper says, each consumer gets its own copy and its own retries, which is the publish/subscribe shape without needing a broker that does pub/sub.

and each queue has its own worker, which is the 7.1 worker again with the event vocabulary:

this.worker = new Worker(
  "email-queue",
  async (job) => {
    return runWithContext({ requestId: `job_${job.id}`, jobId: job.id ?? "unknown", queueName: "email-queue", startTime: performance.now() }, async () => {
      const container = this.buildContainer();
      const scope = container.createScope();

      const jobId = job.id;
      const eventCode = job.name as EmailQueueDomainEvents;

      if (!jobId) throw new BadRequestError("jobId is required");

      try {
        // 1. schema lookup + validation
        const payloadSchema = domainEventsPayloadSchemas.shape[eventCode];

        if (!payloadSchema) {
          throw new ValidationError("Domain Event", `Invalid Domain Event: ${eventCode}`);
        }

        const payload = payloadSchema.parse(job.data) as EmailDomainEventsPayloadTypes<typeof eventCode>;

        // 2. build typed command
        const command = buildEmailQueueEventCommand(eventCode, payload);

        // 3. resolve service & execute (fully typed end-to-end)
        await executeEmailQueueEventHandler(eventCode, scope, command, jobId);
      } catch (error) {
        this.logger.error("Domain Event job failed", error as Error, { jobId, eventCode });
        throw error; // BullMQ handles retries
      } finally {
        await scope.dispose();
      }
    });
  },
  { connection: this.connection, concurrency: 3, lockDuration: 30000, stalledInterval: 30000 },
);

the embedding queue's worker is the same file with "embedding-queue", EmbeddingQueueDomainEvents, buildEmbeddingQueueEventCommand and executeEmbeddingQueueEventHandler swapped in. i'm not going to paste it again, the structure is identical, and that identical structure is itself a design point, a new consumer queue is mostly copy-and-rename (and a candidate for a generic worker factory if you add many more).

8.2 A service per event, per queue

inside one queue, every event it consumes has its own handler service, and the set of events a queue handles is a subset type of all domain events:

export type EmailQueueDomainEvents =
  | typeof DomainEventCode.ORDER_CREATED
  | typeof DomainEventCode.ORDER_CONFIRMED
  | typeof DomainEventCode.ORDER_CANCELLED
  | typeof DomainEventCode.ORDER_DELIVERED
  | typeof DomainEventCode.ORDER_RETURNED
  | typeof DomainEventCode.RATING_APPROVED
  | typeof DomainEventCode.RATING_REJECTED
  | typeof DomainEventCode.RATING_SUBMITTED
  | typeof DomainEventCode.USER_REGISTERED;

export type EmbeddingQueueDomainEvents =
  | typeof DomainEventCode.PRODUCT_CREATED
  | typeof DomainEventCode.PRODUCT_UPDATED
  | typeof DomainEventCode.PRODUCT_DELETED;
  • the subset is what keeps the registries small and exhaustive: the email registry must handle exactly these 9 events (not all 40+ domain events), the embedding registry exactly these 3.
  • the routing (mapper in the publisher) and the handling (subset + registry in the worker) are two halves of the same decision, and they must agree: an event you add to a queue's subset needs a mapper entry that routes it there, and an event the mapper routes there needs a handler. nothing connects those two files at compile time, so a test that checks "every event in each subset is routed to its queue" is worth writing.
  • a handler service is the same shape as 7.2: command + jobId in, idempotency key + work inside a transaction. the embedding one, for example, does the slow and fallible work (read the product, chunk, call the embedding API) before opening the transaction, then keeps the transaction to two writes: the idempotency key and the new chunks (the full story is in the agent post's RAG section).

8.3 Registry, validation and idempotency

the registries have the exact same four pieces as 7.3, just keyed by DomainEventCode:

type EmbeddingQueueEventToCommand = {
  [DomainEventCode.PRODUCT_CREATED]: EmbeddingQueueProductUpsertedEventsHandlerCommand;
  [DomainEventCode.PRODUCT_UPDATED]: EmbeddingQueueProductUpsertedEventsHandlerCommand;
  [DomainEventCode.PRODUCT_DELETED]: EmbeddingQueueProductDeletedEventHandlerCommand;
};

const embeddingQueueEventHandlerRegistry: {
  [K in EmbeddingQueueDomainEvents]: EmbeddingQueueEventHandlerRegistryEntry<K>;
} = {
  [DomainEventCode.PRODUCT_CREATED]: {
    token: EMBEDDING_QUEUE_PRODUCT_UPSERTED_EVENTS_HANDLER_SERVICE,
    async handlerMethod(handler, command, jobId) {
      return handler.execute(command, jobId);
    },
  },
  [DomainEventCode.PRODUCT_UPDATED]: {
    token: EMBEDDING_QUEUE_PRODUCT_UPSERTED_EVENTS_HANDLER_SERVICE, // same service as created
    async handlerMethod(handler, command, jobId) {
      return handler.execute(command, jobId);
    },
  },
  [DomainEventCode.PRODUCT_DELETED]: {
    token: EMBEDDING_QUEUE_PRODUCT_DELETED_EVENTS_HANDLER_SERVICE,
    async handlerMethod(handler, command, jobId) {
      return handler.execute(command, jobId);
    },
  },
};

things worth noting:

  • many events, one service. product.created and product.updated map to the same handler, because "upsert the embeddings of this product" is the same work for both. the registry is where that decision is visible, and it's why the command type for both is the same class.
  • one shared schema map. all queues validate against the same domainEventsPayloadSchemas, the event payload is a contract of the domain, not of the queue. each queue only derives its own typed view: EmailDomainEventsPayloadTypes<T extends EmailQueueDomainEvents> indexes the shared inferred type with the queue's subset, so a handler can't ask for a payload of an event its queue doesn't consume.
  • dates survive the round trip. an event holds a Date, but it crossed JSON twice, so its schema is z.iso.datetime().pipe(z.coerce.date()): validate the ISO string, then hand the handler a real Date again.
  • the casts again. buildEmbeddingQueueEventCommand has the same as ...[T] seams as the outbox one, and that's where a copy-paste mistake can hide: if the PRODUCT_DELETED case builds the upserted command (same constructor shape, so it compiles), a delete event would re-embed the product instead of removing its chunks. exhaustiveness checks prove every event is handled, they can't prove it's handled by the right class. per-event tests are the cheap insurance.
  • idempotency per handler. every handler inserts (jobId, handlerName) at the start of its transaction, exactly as in 7.2. fan-out is why the composite key matters: order.created reaches the email, inventory and analytics queues with the same job id, and each handler must be allowed to process it exactly once, independently of the others.

Part IX, Example Use Cases

with the pipeline in place, a new reaction to something that happened in the system stops being an architectural question and becomes a checklist (queue, mapper entry, subset, handler, token, process).

Email sending. the email queue consumes order lifecycle events, rating moderation and user registration. a use case like ConfirmOrderService has no idea an email exists, it records order.confirmed and commits, the email handler builds the message and sends it with its own retries. if the mail provider has an outage, orders keep flowing, emails go out when it's back, and the idempotency key keeps a retried job from emailing the customer twice.

Embeddings knowledge base updates. product.created, product.updated and product.deleted go to the embedding queue, whose handlers keep the pgvector table in sync with the catalog. the agent post's offline RAG pipeline is this pipeline, the only embedding-specific code is the handler services, the delivery guarantees are the ones from this post, which is why it can be "just another consumer".

Shipping provider actions. the outbox action side: create_order_in_shipping_api, update_order_in_shipping_api, create_shipment_in_shipping_api and delete_order_in_shipping_api. these are the external writes that can never happen inside the business transaction, each one is a row written next to the state change that caused it, delivered to exactly one handler, retried until the provider accepts it.

Inventory and analytics. the mapper already routes order.created and order.cancelled to the inventory and analytics queues, they're the next consumers in line (stock reports, sales dashboards, low-stock alerts). the point of listing them is how little is left to do: write the handler services and a registry, add the worker process, and nothing about publishing, retries or recovery changes.


Part X, One Process Per Concern

Same idea as the agent post: every concern in this post is its own process over one image, with its own composition root.

x-service-defaults: &service-defaults
  image: ghcr.io/djaouad10/ddd-e-commerce-backend:latest
  restart: unless-stopped
  env_file: .env

services:
  outbox-processor:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/outbox-processor.js"]
    scale: 1

  domain-events-processor:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/domain-events-processor.js"]
    scale: 1

  outbox-handler:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/outbox-handler.js"]

  email-queue-handler:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/email-queue-handler.js"]

  embedding-queue-handler:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/embedding-queue-handler.js"]

  clean-outbox:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/clean-outbox.js"]
    scale: 1

  reset-stuck-outbox-rows:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/reset-stuck-outbox-rows.js"]
    scale: 1
ProcessReadsWritesScaling
outbox-processoroutbox (actions)outbox-queue1 (more is safe, claim is atomic)
domain-events-processoroutbox (events)all event queues1 (same)
outbox-handleroutbox-queueexternal APIs, idempotency_keyshorizontal
*-queue-handlerits queueits own domain (email, embeddings,... etc.)horizontal, per concern
reset-stuck-outbox-rowsoutboxoutbox1
clean-outboxoutboxoutbox1
  • the relay processes stay at one replica by choice, not by necessity: the claim is an atomic compare-and-set, so two replicas would never process the same row, one is simply enough for this load and makes the behavior easier to reason about.
  • the handler processes scale horizontally, BullMQ distributes jobs between workers, and idempotency keys make an accidental double delivery harmless.
  • every process only builds the subgraph it touches: the processors never construct a handler service, the handlers never construct the publisher.

Part XI, Close

Recap in one breath: a state change and the intent to tell the world about it, written in one transaction → a single outbox table with two categories → polling workers that claim, publish, back off and mark → helpers that reset stuck rows and clean completed ones → BullMQ queues per consumer concern → handlers that validate through a registry and dedupe through an idempotency key → at-least-once delivery that behaves like exactly-once.

The repo: github.com/djaouad10/DDD-E-commerce-Backend. Further reading: Chris Richardson's Microservices Patterns (the Transactional Outbox chapter), and Designing Data-Intensive Applications by Martin Kleppmann for the delivery-guarantee fundamentals.

Thanks for reading 👋More posts →