Siddhant Deval
Siddhant Deval
backend5 min read

Domain Events & The Mediator Pattern: Decoupled Cross-Aggregate Workflows

When an Order is placed, the Inventory, Billing, and Notification services must react — but they must not be directly coupled to the Order Aggregate. Domain Events and the Mediator pattern decouple these cross-aggregate side effects through an in-process event bus, enabling rich domain workflows without transaction scope violations.

Domain Events & The Mediator Pattern: Decoupled Cross-Aggregate Workflows

Domain objects are autonomous state machines with enforced invariants — they are not passive bags of data passed between controller functions. And when an Order is placed, three things must happen: inventory is reserved, payment is captured, and a confirmation email is sent. The temptation is to call all three services directly from inside place(). This is the exact coupling that Domain Events exist to prevent. An Order knows that something significant happened. It does not know — and must not know — who cares or what they will do about it.

This article implements the full Domain Event machinery for the billing engine: the event collection protocol inside AggregateRoot, the IEventBus port, concrete event handlers for inventory, payment, and email, and the precise timing rule for when events are dispatched relative to the database transaction.

Architectural Note

Prerequisites: Part 4 (Aggregates & Repositories) — specifically the AggregateRoot.pullDomainEvents() mechanism, which is the collection side of the pattern implemented here. The dispatch side (the IEventBus) is introduced in this article.


1. The Direct Coupling Anti-Pattern

This is the implementation most engineers write first, and it ships to production in most Node.js services:

TYPESCRIPT
// ❌ Anti-Pattern: Order.place() directly orchestrates cross-cutting side effects
export class Order extends AggregateRoot<OrderId> {
  constructor(
    // ❌ The Order domain object now depends on application-layer services
    private readonly inventoryService: InventoryService,
    private readonly paymentService: PaymentService,
    private readonly emailService: EmailService,
  ) { /* ... */ }

  async place(): Promise<void> {
    if (this._status !== OrderStatus.PENDING) throw new DomainError('...');
    if (this._items.length === 0) throw new DomainError('...');

    this._status = OrderStatus.PLACED;

    // ❌ The Order now directly triggers three external side effects
    await this.inventoryService.reserve(this._items);  // What if this fails?
    await this.paymentService.capture(this.total);     // What if payment fails after inventory?
    await this.emailService.sendConfirmation(this._customerId); // Tight coupling, no retry
  }
}

Count the problems:

  1. Domain object depends on application services — Order is no longer a pure domain object. It imports InventoryService and PaymentService, which import Stripe and Redis. The Order class can no longer be instantiated in a unit test without a live payment processor.
  2. No partial failure handling — if inventoryService.reserve() succeeds but paymentService.capture() throws, the inventory is reserved but the order is never paid. The domain is now in an inconsistent state with no compensation mechanism.
  3. Order of side effects is hardcoded — if the business decides to send the confirmation email before capturing payment (to inform the customer while payment processes asynchronously), the domain object must be modified.
  4. Adding a new side effect requires modifying the Order — when the fraud detection team wants to run a risk check when an order is placed, they open Order.ts and add another call. The Order class grows without bound.

Domain Events solve all four problems simultaneously.


2. Domain Events as First-Class Domain Concepts

2.1 What a Domain Event Is — and Is Not

A Domain Event is a fact that happened in the domain, named in past tense using the Ubiquitous Language of the business:

✅ Domain Event (fact, past tense) ❌ Not a Domain Event
OrderPlaced PlaceOrder (that is a Command)
PaymentCaptured PaymentCapture (ambiguous)
SubscriptionRenewed RenewSubscriptionEvent (redundant suffix)
CustomerAccountSuspended UserStatusChanged (too generic — no domain language)

A Domain Event is not a command (don't call OrderPlaced "a request to place an order"). It is not a notification (don't call it an "event notification"). It is a fact: something happened, it cannot be undone, and interested parties may react to it.

2.2 The DomainEvent Base Class

TYPESCRIPT
// src/domain/shared/DomainEvent.ts
import { randomUUID } from 'node:crypto';

export abstract class DomainEvent {
  /** Globally unique event ID for idempotency checks */
  public readonly eventId: string;
  /** When the event occurred — set at the moment the domain method executes */
  public readonly occurredAt: Date;
  /** The ID of the Aggregate Root that raised this event */
  public readonly aggregateId: string;
  /** Schema version — used for event versioning and migration (Part 11) */
  public readonly eventVersion: number;

  protected constructor(aggregateId: string, eventVersion: number = 1) {
    this.eventId   = randomUUID();
    this.occurredAt = new Date();
    this.aggregateId = aggregateId;
    this.eventVersion = eventVersion;
  }

  /** Type discriminator — used by event handlers for routing */
  abstract get eventType(): string;
}

2.3 Concrete Domain Events for the Billing Engine

TYPESCRIPT
// src/domain/order/events/OrderPlaced.ts
import { DomainEvent } from '../../shared/DomainEvent';
import { OrderId } from '../OrderId';
import { CustomerId } from '../../customer/CustomerId';
import { Money } from '../../payment/Money';

export class OrderPlaced extends DomainEvent {
  constructor(
    public readonly orderId: OrderId,
    public readonly customerId: CustomerId,
    public readonly orderTotal: Money,
    public readonly itemCount: number,
  ) {
    super(orderId);
  }

  get eventType() { return 'order.placed' as const; }
}
TYPESCRIPT
// src/domain/order/events/OrderCancelled.ts
import { DomainEvent } from '../../shared/DomainEvent';
import { OrderId } from '../OrderId';

export class OrderCancelled extends DomainEvent {
  constructor(
    public readonly orderId: OrderId,
    public readonly cancellationReason: string,
    public readonly cancelledAt: Date = new Date(),
  ) {
    super(orderId);
  }

  get eventType() { return 'order.cancelled' as const; }
}
TYPESCRIPT
// src/domain/payment/events/PaymentCaptured.ts
import { DomainEvent } from '../../shared/DomainEvent';
import { PaymentId } from '../PaymentId';
import { OrderId } from '../../order/OrderId';
import { Money } from '../Money';

export class PaymentCaptured extends DomainEvent {
  constructor(
    public readonly paymentId: PaymentId,
    public readonly orderId: OrderId,
    public readonly amount: Money,
    public readonly providerReference: string,
  ) {
    super(paymentId);
  }

  get eventType() { return 'payment.captured' as const; }
}

2.4 Raising Events Inside Aggregate Methods

Events are collected — not dispatched — during domain method execution:

TYPESCRIPT
// src/domain/order/Order.ts
export class Order extends AggregateRoot<OrderId> {

  place(): void {
    // Guard clauses first
    if (this._status !== OrderStatus.PENDING)
      throw new DomainError(`Order ${this.id} cannot be placed — status is ${this._status}`);
    if (this._items.length === 0)
      throw new DomainError(`Order ${this.id} cannot be placed with no line items`);
    if (!this._shippingAddress)
      throw new DomainError(`Order ${this.id} requires a shipping address`);

    // State transition
    this._status = OrderStatus.PLACED;

    // ✅ Raise the event — do NOT dispatch it
    // The event is queued internally; the Use Case dispatches it after the DB commit
    this.addDomainEvent(new OrderPlaced(
      this.id,
      this._customerId,
      this.total,
      this._items.length,
    ));
  }

  cancel(reason: string): void {
    const cancellableStatuses = [OrderStatus.PENDING, OrderStatus.PLACED];
    if (!cancellableStatuses.includes(this._status))
      throw new DomainError(`Order ${this.id} cannot be cancelled — status is ${this._status}`);

    this._status = OrderStatus.CANCELLED;
    this.addDomainEvent(new OrderCancelled(this.id, reason));
  }
}

The addDomainEvent() method (inherited from AggregateRoot) pushes the event onto an internal private array. The method returns synchronously. No side effects occur. The Order object has no knowledge of what will happen to the event.


3. The IEventBus Port and The Mediator Pattern

3.1 The IEventBus Domain Interface

The event bus interface is defined in the Application layer (it could also be in Domain — the key is that it is defined at or inside the layer that uses it, not in Infrastructure):

TYPESCRIPT
// src/application/shared/IEventBus.ts
import { DomainEvent } from '../../domain/shared/DomainEvent';

export interface IEventHandler<T extends DomainEvent> {
  handle(event: T): Promise<void>;
}

export interface IEventBus {
  /** Dispatch a single domain event to all registered handlers */
  publish(event: DomainEvent): Promise<void>;

  /** Register a handler for a specific event type */
  subscribe<T extends DomainEvent>(
    eventType: string,
    handler: IEventHandler<T>,
  ): void;
}

3.2 The In-Memory Mediator (for Tests and Development)

TYPESCRIPT
// src/infrastructure/messaging/InMemoryEventBus.ts
import { IEventBus, IEventHandler } from '../../application/shared/IEventBus';
import { DomainEvent } from '../../domain/shared/DomainEvent';

export class InMemoryEventBus implements IEventBus {
  private readonly handlers = new Map<string, IEventHandler<DomainEvent>[]>();

  subscribe<T extends DomainEvent>(eventType: string, handler: IEventHandler<T>): void {
    const existing = this.handlers.get(eventType) ?? [];
    this.handlers.set(eventType, [...existing, handler as IEventHandler<DomainEvent>]);
  }

  async publish(event: DomainEvent): Promise<void> {
    const handlers = this.handlers.get(event.eventType) ?? [];
    // Sequential dispatch — for parallel, use Promise.allSettled
    for (const handler of handlers) {
      await handler.handle(event);
    }
  }
}

This implementation is synchronous and in-process. For development and testing, this is sufficient. For production (high-throughput billing), the Kafka-backed implementation dispatches to a durable message broker with at-least-once delivery (covered in Part 11 with the Outbox Pattern).

3.3 The Spy Event Bus (for Asserting in Tests)

TYPESCRIPT
// src/infrastructure/messaging/SpyEventBus.ts (test utility)
import { IEventBus, IEventHandler } from '../../application/shared/IEventBus';
import { DomainEvent } from '../../domain/shared/DomainEvent';

export class SpyEventBus implements IEventBus {
  private readonly published: DomainEvent[] = [];

  async publish(event: DomainEvent): Promise<void> {
    this.published.push(event);
  }

  subscribe<T extends DomainEvent>(_: string, __: IEventHandler<T>): void {
    // Spy does not dispatch to real handlers — it only records
  }

  // ── Test assertion helpers ──

  eventsOfType<T extends DomainEvent>(eventType: string): T[] {
    return this.published.filter(e => e.eventType === eventType) as T[];
  }

  get publishedCount(): number { return this.published.length; }

  wasPublished(eventType: string): boolean {
    return this.published.some(e => e.eventType === eventType);
  }

  clear(): void { this.published.length = 0; }
}

4. Concrete Event Handlers

4.1 InventoryReservationHandler

TYPESCRIPT
// src/application/order/handlers/InventoryReservationHandler.ts
import { IEventHandler } from '../../shared/IEventBus';
import { OrderPlaced } from '../../../domain/order/events/OrderPlaced';
import { IInventoryRepository } from '../../../domain/inventory/ports/IInventoryRepository';
import { DomainError } from '../../../domain/shared/DomainError';

export class InventoryReservationHandler implements IEventHandler<OrderPlaced> {
  constructor(private readonly inventory: IInventoryRepository) {}

  async handle(event: OrderPlaced): Promise<void> {
    // Load the Order's line items from the read model (more efficient than reloading Order)
    const items = await this.inventory.findAvailability(event.orderId);

    for (const item of items) {
      if (item.availableQuantity < item.requestedQuantity) {
        // Raise a compensating event — Part 11 handles saga compensation
        throw new DomainError(
          `Insufficient inventory for product ${item.productId}: ` +
          `requested ${item.requestedQuantity}, available ${item.availableQuantity}`
        );
      }
    }

    await this.inventory.reserve(event.orderId, items);
  }
}

4.2 OrderConfirmationEmailHandler

TYPESCRIPT
// src/application/order/handlers/OrderConfirmationEmailHandler.ts
import { IEventHandler } from '../../shared/IEventBus';
import { OrderPlaced } from '../../../domain/order/events/OrderPlaced';
import { ICustomerRepository } from '../../../domain/customer/ports/ICustomerRepository';
import { IEmailPort } from '../../../domain/shared/ports/IEmailPort';

export class OrderConfirmationEmailHandler implements IEventHandler<OrderPlaced> {
  constructor(
    private readonly customers: ICustomerRepository,
    private readonly email: IEmailPort,
  ) {}

  async handle(event: OrderPlaced): Promise<void> {
    const customer = await this.customers.findById(event.customerId);
    if (!customer) return; // Defensive — customer may have been deleted

    await this.email.send({
      to: customer.email.value,
      subject: `Order Confirmation — ${event.orderId}`,
      template: 'order-confirmation',
      variables: {
        customerName: customer.name,
        orderId: event.orderId,
        total: event.orderTotal.toString(),
        itemCount: event.itemCount,
      },
    });
  }
}

Notice: OrderConfirmationEmailHandler is in the Application layer (src/application/). It depends on ICustomerRepository and IEmailPort — both Domain interfaces. It knows nothing about nodemailer or SendGrid. Those implementations are in Infrastructure.


5. The Dispatch Timing Rule: After Commit, Not During

This is the most common implementation mistake with Domain Events, and it causes real production incidents.

5.1 The Wrong Way: Dispatch During the Transaction

TYPESCRIPT
// ❌ Anti-Pattern: Dispatch events BEFORE the database commits
async execute(command: PlaceOrderCommand): Promise<OrderId> {
  const order = Order.create(command.customerId);
  order.place();

  // ❌ Events dispatched before save() — if save() fails, events already fired
  const events = order.pullDomainEvents();
  for (const event of events) await this.eventBus.publish(event);

  await this.orderRepo.save(order); // If this throws, events were already dispatched
  return order.id;
}

If orderRepo.save() fails (database connection drops, optimistic concurrency conflict), the events have already been dispatched. The InventoryReservationHandler has already reserved inventory for an order that was never saved. The confirmation email has already been sent. The billing system is now in an inconsistent state.

5.2 The Correct Way: Dispatch After Commit

TYPESCRIPT
// ✅ Correct: Dispatch events AFTER the database commits
async execute(command: PlaceOrderCommand): Promise<OrderId> {
  // 1. Load or create the Aggregate
  const order = Order.create(command.customerId);
  for (const item of command.items) {
    order.addItem(item.productId, item.quantity, item.unitPrice);
  }
  order.setShippingAddress(command.shippingAddress);

  // 2. Execute domain logic — events are COLLECTED here, not dispatched
  order.place();

  // 3. Persist the Aggregate — this is the COMMIT point
  await this.orderRepo.save(order);

  // 4. Pull events AFTER successful commit
  const domainEvents = order.pullDomainEvents();

  // 5. Dispatch events — handlers see a state that is now permanently in the database
  for (const event of domainEvents) {
    await this.eventBus.publish(event);
  }

  return order.id;
}

Step 3 is the transaction boundary. Steps 4 and 5 happen only if Step 3 succeeds. If the database commit fails (Step 3 throws), Steps 4 and 5 are never reached — no events are dispatched for an order that was never saved.

Crucial Requirement

The post-commit dispatch guarantee is only eventually consistent, not atomic. If the server crashes between Step 3 (successful DB commit) and Step 5 (event dispatch), the events are lost. For billing systems that require at-least-once event delivery, the Outbox Pattern (Part 11) writes events to the database in the same transaction as the Aggregate save, and a separate background process publishes them durably. For the purposes of Parts 5 through 10, in-process post-commit dispatch is sufficient.


6. Wiring Handlers in the DI Container

In the composition root (and more precisely, in the DI container module from Part 9), handlers are registered against their event types:

TYPESCRIPT
// src/infrastructure/container/eventHandlerModule.ts
import { InMemoryEventBus } from '../messaging/InMemoryEventBus';
import { InventoryReservationHandler } from '../../application/order/handlers/InventoryReservationHandler';
import { OrderConfirmationEmailHandler } from '../../application/order/handlers/OrderConfirmationEmailHandler';
// ... additional imports

export function wireEventHandlers(
  eventBus: InMemoryEventBus,
  deps: { inventory: IInventoryRepository; customers: ICustomerRepository; email: IEmailPort }
): void {
  eventBus.subscribe('order.placed', new InventoryReservationHandler(deps.inventory));
  eventBus.subscribe('order.placed', new OrderConfirmationEmailHandler(deps.customers, deps.email));
  // Multiple handlers for the same event type — all are called in sequence
}

Adding a new handler for order.placed — say, a fraud detection check — requires:

  1. Create FraudDetectionHandler implementing IEventHandler<OrderPlaced>
  2. Register it: eventBus.subscribe('order.placed', new FraudDetectionHandler(...))

Zero changes to Order.ts, zero changes to PlaceOrderUseCase.ts, zero changes to any existing handler.


7. Testing the Full Event Lifecycle

TYPESCRIPT
// tests/unit/application/PlaceOrderUseCase.events.test.ts
import { PlaceOrderUseCase } from '../../../src/application/order/PlaceOrderUseCase';
import { InMemoryOrderRepository } from '../../../src/infrastructure/persistence/InMemoryOrderRepository';
import { SpyEventBus } from '../../../src/infrastructure/messaging/SpyEventBus';
import { OrderPlaced } from '../../../src/domain/order/events/OrderPlaced';
import { Money } from '../../../src/domain/payment/Money';

describe('PlaceOrderUseCase — domain event dispatch', () => {
  let orderRepo: InMemoryOrderRepository;
  let eventBus: SpyEventBus;
  let useCase: PlaceOrderUseCase;

  beforeEach(() => {
    orderRepo = new InMemoryOrderRepository();
    eventBus = new SpyEventBus();
    useCase = new PlaceOrderUseCase(orderRepo, eventBus);
  });

  it('should raise exactly one OrderPlaced event after placing an order', async () => {
    const orderId = await useCase.execute({
      customerId: 'cust_001' as any,
      items: [{ productId: 'prod_a' as any, quantity: 2, unitPrice: Money.of('USD', 49.99) }],
      shippingAddress: testAddress(),
    });

    // Assert the event was published
    expect(eventBus.wasPublished('order.placed')).toBe(true);
    expect(eventBus.publishedCount).toBe(1);

    // Assert event payload is correct
    const [event] = eventBus.eventsOfType<OrderPlaced>('order.placed');
    expect(event.orderId).toBe(orderId);
    expect(event.itemCount).toBe(1);
    expect(event.orderTotal.amount).toBeCloseTo(99.98);
  });

  it('should NOT publish events if order placement fails (domain guard)', async () => {
    // Empty items — place() will throw DomainError
    await expect(useCase.execute({
      customerId: 'cust_001' as any,
      items: [], // ← triggers DomainError in order.place()
      shippingAddress: testAddress(),
    })).rejects.toThrow(DomainError);

    expect(eventBus.publishedCount).toBe(0); // No events for a failed placement
  });

  it('should NOT publish events if repository.save() throws', async () => {
    // Simulate a database failure
    orderRepo.save = async () => { throw new Error('DB connection lost'); };

    await expect(useCase.execute({
      customerId: 'cust_001' as any,
      items: [{ productId: 'prod_a' as any, quantity: 1, unitPrice: Money.of('USD', 10.00) }],
      shippingAddress: testAddress(),
    })).rejects.toThrow('DB connection lost');

    expect(eventBus.publishedCount).toBe(0); // Events only dispatch after successful commit
  });
});

These three tests verify all three critical behaviors: events are raised on success, events are not raised when the domain rejects the command, and events are not raised when persistence fails. Each runs in under 2ms with no external dependencies.


Summary

Concept Domain Rule
Domain Event A past-tense fact in the Ubiquitous Language; never a command or a request
addDomainEvent() Called inside Aggregate methods; events are collected, not dispatched
pullDomainEvents() Called by the Use Case after repository.save() commits
Post-Commit Dispatch Events are published only after the database transaction succeeds — not during
IEventBus An Application-layer interface; in-process for dev/tests, Kafka-backed in production
Mediator Decoupling Order raises OrderPlaced; it does not know that InventoryReservationHandler exists
SpyEventBus Records published events without dispatching to real handlers; enables precise event assertions

What's Next

In Part 6, we build the Application layer — implementing PlaceOrderUseCase, CancelOrderUseCase, and ProcessRefundUseCase as thin orchestrators that coordinate Domain and Infrastructure without containing any business rules themselves. Part 6: Application Layer →

Research & Synthesis Note

This article was developed with AI-assisted deep search, specification cross-referencing, and technical research synthesis.

#Domain Events#Mediator Pattern#Domain-Driven Design#TypeScript#Node.js#Event-Driven Architecture#OOP
Siddhant Deval

Written by Siddhant Deval

Senior Full-Stack Engineer building high-scale architectures, browser performance engineering systems, and SaaS platforms.