CQRS + Event Sourcing — 05 — Projet réel : système de commandes complet

Assembler CQRS + Event Sourcing dans un projet TypeScript et Python complet : handlers, event store PostgreSQL, projections, queries. Architecture finale et structure de fichiers.

05 — Projet réel : système de commandes complet

Ce que tu vas apprendre

  • Assembler tous les concepts de la série dans un projet cohérent
  • La structure de fichiers recommandée
  • Les trois flux principaux : créer, confirmer, annuler une commande
  • Les queries : lire l'état courant et l'historique
  • Les points de vigilance en production

Prérequis


Ce qu'on va construire

Un système de gestion de commandes avec :

  • 4 commands : CreateOrder, ConfirmOrder, ShipOrder, CancelOrder
  • 5 events : OrderCreated, OrderConfirmed, OrderShipped, OrderCancelled, OrderDelivered
  • 2 projections : OrderSummary (liste), CustomerOrderHistory (profil)
  • 3 queries : lire une commande, lister par client, résumé des commandes

Structure de fichiers

src/
├── domain/
│   ├── order/
│   │   ├── order.aggregate.ts     — Agrégat + rehydration
│   │   ├── order.commands.ts      — Types des commands
│   │   ├── order.events.ts        — Types des domain events
│   │   └── order.handlers.ts      — Command handlers
│   └── shared/
│       └── result.ts              — Result type
├── infrastructure/
│   ├── event-store.ts             — PostgresEventStore
│   ├── snapshot-store.ts          — SnapshotStore
│   └── projection-engine.ts       — ProjectionEngine
├── projections/
│   ├── order-summary.projection.ts
│   └── customer-history.projection.ts
├── queries/
│   └── order.queries.ts           — QueryHandlers
└── migrations/
    └── 001_create_tables.sql

Les types : commands et events

typescript// order.commands.ts
export interface CreateOrderCommand {
  readonly customerId: string;
  readonly items: ReadonlyArray<{
    productId: string;
    quantity: number;
    unitPriceCents: number;
  }>;
  readonly shippingAddress: string;
}

export interface ConfirmOrderCommand {
  readonly orderId: string;
}

export interface ShipOrderCommand {
  readonly orderId: string;
  readonly trackingId: string;
}

export interface CancelOrderCommand {
  readonly orderId: string;
  readonly reason: string;
}
typescript// order.events.ts
export interface OrderCreatedEvent {
  readonly eventType: "OrderCreated";
  readonly orderId: string;
  readonly customerId: string;
  readonly items: ReadonlyArray<{ productId: string; quantity: number; unitPriceCents: number }>;
  readonly totalCents: number;
  readonly shippingAddress: string;
  readonly occurredAt: string;
}

export interface OrderConfirmedEvent {
  readonly eventType: "OrderConfirmed";
  readonly orderId: string;
  readonly occurredAt: string;
}

export interface OrderShippedEvent {
  readonly eventType: "OrderShipped";
  readonly orderId: string;
  readonly trackingId: string;
  readonly occurredAt: string;
}

export interface OrderCancelledEvent {
  readonly eventType: "OrderCancelled";
  readonly orderId: string;
  readonly reason: string;
  readonly occurredAt: string;
}

export type OrderEvent =
  | OrderCreatedEvent
  | OrderConfirmedEvent
  | OrderShippedEvent
  | OrderCancelledEvent;

L'agrégat Order complet

typescript// order.aggregate.ts
import type { StoredEvent } from "../infrastructure/event-store";
import type { OrderEvent } from "./order.events";

type OrderStatus = "pending" | "confirmed" | "shipped" | "delivered" | "cancelled";

interface OrderState {
  orderId: string;
  customerId: string;
  status: OrderStatus;
  totalCents: number;
  shippingAddress: string;
  trackingId: string | null;
  version: number;
}

const INITIAL_STATE: OrderState = {
  orderId: "",
  customerId: "",
  status: "pending",
  totalCents: 0,
  shippingAddress: "",
  trackingId: null,
  version: -1,
};

export class OrderAggregate {
  private _state: OrderState;
  private _pendingEvents: OrderEvent[] = [];

  private constructor(state: OrderState) {
    this._state = { ...state };
  }

  static rehydrate(events: StoredEvent[]): OrderAggregate {
    const agg = new OrderAggregate(INITIAL_STATE);
    for (const ev of events) {
      agg._applyStored(ev);
    }
    return agg;
  }

  static fromSnapshot(snapshot: OrderState): OrderAggregate {
    return new OrderAggregate(snapshot);
  }

  private _applyStored(ev: StoredEvent): void {
    this._mutate(ev.eventType, ev.data as Record<string, unknown>);
    this._state.version = ev.version;
  }

  private _mutate(type: string, data: Record<string, unknown>): void {
    switch (type) {
      case "OrderCreated":
        this._state.orderId       = data.orderId as string;
        this._state.customerId    = data.customerId as string;
        this._state.totalCents    = data.totalCents as number;
        this._state.shippingAddress = data.shippingAddress as string;
        this._state.status        = "pending";
        break;
      case "OrderConfirmed":
        this._state.status = "confirmed";
        break;
      case "OrderShipped":
        this._state.status    = "shipped";
        this._state.trackingId = data.trackingId as string;
        break;
      case "OrderDelivered":
        this._state.status = "delivered";
        break;
      case "OrderCancelled":
        this._state.status = "cancelled";
        break;
    }
  }

  private _emit(event: OrderEvent): void {
    this._mutate(event.eventType, event as unknown as Record<string, unknown>);
    this._pendingEvents.push(event);
  }

  // ---- Commands ----

  confirm(): void {
    if (this._state.status !== "pending") {
      throw new Error(`Impossible de confirmer une commande "${this._state.status}"`);
    }
    this._emit({
      eventType: "OrderConfirmed",
      orderId: this._state.orderId,
      occurredAt: new Date().toISOString(),
    });
  }

  ship(trackingId: string): void {
    if (this._state.status !== "confirmed") {
      throw new Error(`Impossible d'expédier une commande "${this._state.status}"`);
    }
    if (!trackingId?.trim()) {
      throw new Error("trackingId obligatoire");
    }
    this._emit({
      eventType: "OrderShipped",
      orderId: this._state.orderId,
      trackingId: trackingId.trim(),
      occurredAt: new Date().toISOString(),
    });
  }

  cancel(reason: string): void {
    if (["delivered", "cancelled"].includes(this._state.status)) {
      throw new Error(`Impossible d'annuler une commande "${this._state.status}"`);
    }
    if (!reason?.trim()) {
      throw new Error("La raison de l'annulation est obligatoire");
    }
    this._emit({
      eventType: "OrderCancelled",
      orderId: this._state.orderId,
      reason: reason.trim(),
      occurredAt: new Date().toISOString(),
    });
  }

  // ---- Accesseurs ----
  get orderId(): string { return this._state.orderId; }
  get status(): OrderStatus { return this._state.status; }
  get version(): number { return this._state.version; }
  get snapshot(): OrderState { return { ...this._state }; }
  get pendingEvents(): OrderEvent[] { return [...this._pendingEvents]; }
  clearPendingEvents(): void { this._pendingEvents = []; }
}

Les command handlers

typescript// order.handlers.ts
export class OrderCommandService {
  constructor(
    private readonly eventStore: PostgresEventStore,
    private readonly snapshotStore: SnapshotStore,
    private readonly projectionEngine: ProjectionEngine
  ) {}

  async createOrder(cmd: CreateOrderCommand): Promise<string> {
    // Validation
    if (!cmd.customerId?.trim()) throw new Error("customerId obligatoire");
    if (!cmd.items?.length) throw new Error("Au moins un article requis");
    for (const item of cmd.items) {
      if (item.quantity <= 0) throw new Error(`Quantité invalide : ${item.productId}`);
      if (item.unitPriceCents <= 0) throw new Error(`Prix invalide : ${item.productId}`);
    }

    const orderId = crypto.randomUUID();
    const totalCents = cmd.items.reduce(
      (sum, i) => sum + i.unitPriceCents * i.quantity,
      0
    );

    const event: OrderCreatedEvent = {
      eventType: "OrderCreated",
      orderId,
      customerId: cmd.customerId.trim(),
      items: [...cmd.items],
      totalCents,
      shippingAddress: cmd.shippingAddress,
      occurredAt: new Date().toISOString(),
    };

    const streamId = `order-${orderId}`;
    await this.eventStore.append(streamId, [{ eventType: event.eventType, data: event as any }], -1);

    // Mettre à jour les projections
    const { events } = await this.eventStore.loadFrom(streamId, 0);
    await this.projectionEngine.dispatch(events);

    return orderId;
  }

  async confirmOrder(cmd: ConfirmOrderCommand): Promise<void> {
    await this._executeOnAggregate(`order-${cmd.orderId}`, agg => agg.confirm());
  }

  async shipOrder(cmd: ShipOrderCommand): Promise<void> {
    await this._executeOnAggregate(`order-${cmd.orderId}`, agg => agg.ship(cmd.trackingId));
  }

  async cancelOrder(cmd: CancelOrderCommand): Promise<void> {
    await this._executeOnAggregate(`order-${cmd.orderId}`, agg => agg.cancel(cmd.reason));
  }

  private async _executeOnAggregate(
    streamId: string,
    action: (agg: OrderAggregate) => void
  ): Promise<void> {
    // Charger avec snapshot si disponible
    const snapshot = await this.snapshotStore.load(streamId);
    let agg: OrderAggregate;
    let fromVersion: number;

    if (snapshot) {
      agg = OrderAggregate.fromSnapshot(snapshot.state as any);
      fromVersion = snapshot.version + 1;
    } else {
      const { events } = await this.eventStore.load(streamId);
      if (!events.length) throw new Error(`Commande introuvable : ${streamId}`);
      agg = OrderAggregate.rehydrate(events);
      fromVersion = 0;
    }

    // Charger les events depuis le snapshot
    if (snapshot) {
      const { events } = await this.eventStore.loadFrom(streamId, fromVersion);
      agg = OrderAggregate.rehydrate(events);
    }

    const expectedVersion = agg.version;

    // Exécuter la command
    action(agg);

    const pending = agg.pendingEvents;
    if (!pending.length) return;

    // Persister les nouveaux events
    await this.eventStore.append(
      streamId,
      pending.map(e => ({ eventType: e.eventType, data: e as any })),
      expectedVersion
    );
    agg.clearPendingEvents();

    // Snapshot tous les 50 events
    if ((expectedVersion + 1) % 50 === 0) {
      await this.snapshotStore.save(streamId, agg.snapshot, agg.version);
    }

    // Mettre à jour les projections
    const newVersion = expectedVersion + pending.length;
    const { events: newEvents } = await this.eventStore.loadFrom(streamId, expectedVersion + 1);
    await this.projectionEngine.dispatch(newEvents);
  }
}

Les queries

typescript// order.queries.ts
export class OrderQueryService {
  constructor(private readonly pool: Pool) {}

  async getOrder(orderId: string): Promise<OrderSummaryDTO | null> {
    const { rows } = await this.pool.query(
      `SELECT order_id, customer_id, status, total_cents, item_count, created_at, updated_at
       FROM order_summaries
       WHERE order_id = $1`,
      [orderId]
    );

    if (!rows.length) return null;

    const r = rows[0];
    return {
      orderId: r.order_id,
      customerId: r.customer_id,
      status: r.status,
      totalEuros: r.total_cents / 100,
      itemCount: r.item_count,
      createdAt: r.created_at,
      updatedAt: r.updated_at,
    };
  }

  async listCustomerOrders(
    customerId: string,
    status?: string
  ): Promise<OrderSummaryDTO[]> {
    const { rows } = await this.pool.query(
      `SELECT order_id, customer_id, status, total_cents, item_count, created_at, updated_at
       FROM order_summaries
       WHERE customer_id = $1
         AND ($2::varchar IS NULL OR status = $2)
       ORDER BY created_at DESC
       LIMIT 50`,
      [customerId, status ?? null]
    );

    return rows.map(r => ({
      orderId: r.order_id,
      customerId: r.customer_id,
      status: r.status,
      totalEuros: r.total_cents / 100,
      itemCount: r.item_count,
      createdAt: r.created_at,
      updatedAt: r.updated_at,
    }));
  }

  // Lire l'historique complet depuis l'event store (pour l'audit)
  async getOrderHistory(
    orderId: string,
    eventStore: PostgresEventStore
  ): Promise<StoredEvent[]> {
    const { events } = await eventStore.load(`order-${orderId}`);
    return events;
  }
}

Migration complète

sql-- migrations/001_create_tables.sql

-- Event store
CREATE TABLE IF NOT EXISTS events (
  id          UUID         PRIMARY KEY DEFAULT gen_random_uuid(),
  stream_id   VARCHAR(255) NOT NULL,
  event_type  VARCHAR(255) NOT NULL,
  data        JSONB        NOT NULL,
  metadata    JSONB        NOT NULL DEFAULT '{}',
  version     INTEGER      NOT NULL,
  created_at  TIMESTAMPTZ  NOT NULL DEFAULT NOW()
);

CREATE INDEX IF NOT EXISTS events_stream_id_idx ON events (stream_id);
CREATE INDEX IF NOT EXISTS events_event_type_idx ON events (event_type);
CREATE UNIQUE INDEX IF NOT EXISTS events_stream_version_unique ON events (stream_id, version);

-- Snapshots
CREATE TABLE IF NOT EXISTS snapshots (
  stream_id  VARCHAR(255) PRIMARY KEY,
  state      JSONB        NOT NULL,
  version    INTEGER      NOT NULL,
  created_at TIMESTAMPTZ  NOT NULL DEFAULT NOW()
);

-- Read models
CREATE TABLE IF NOT EXISTS order_summaries (
  order_id      VARCHAR(255) PRIMARY KEY,
  customer_id   VARCHAR(255) NOT NULL,
  status        VARCHAR(50)  NOT NULL,
  total_cents   INTEGER      NOT NULL,
  item_count    INTEGER      NOT NULL DEFAULT 0,
  created_at    TIMESTAMPTZ  NOT NULL,
  updated_at    TIMESTAMPTZ  NOT NULL
);

CREATE INDEX IF NOT EXISTS os_customer_id ON order_summaries (customer_id, created_at DESC);
CREATE INDEX IF NOT EXISTS os_status ON order_summaries (status);

Points de vigilance en production

Concurrence Projections Taille des streams
Retry sur ConcurrencyError Idempotentes avec ON CONFLICT Snapshot après 50 events

Concurrence : quand deux requêtes modifient le même agrégat simultanément, l'une obtiendra un ConcurrencyError. Le caller doit retenter — généralement 2-3 fois suffisent. Si les invariants ne sont plus satisfaits après rechargement, l'opération échoue proprement.

Idempotence des projections : si une projection est rejouée deux fois pour le même event (bug, restart), elle ne doit pas créer de doublon. ON CONFLICT DO UPDATE ou ON CONFLICT DO NOTHING dans les INSERT le garantit.

Taille de l'event store : les events ne sont jamais supprimés. En production, partitionner la table par created_at (PostgreSQL table partitioning) si le volume devient important.


Résumé de la série

Article Contenu
00 — Introduction Définitions, philosophie, quand utiliser
01 — Commands et events Handlers, domain events, validation
02 — Event store PostgreSQL Schéma, append-only, concurrence optimiste
03 — Projections Read models, ProjectionEngine, rebuild
04 — Replay et snapshots Rehydration, snapshots, versioning
05 — Projet réel Système complet TypeScript + Python

CQRS + Event Sourcing ajoute de la complexité — une complexité qui vaut le coût dans les domaines où l'historique est précieux et les besoins de lecture et d'écriture divergent. Pour un CRUD simple, cette architecture est excessive. Pour un système financier, un e-commerce ou un workflow métier complexe, elle apporte une traçabilité et une flexibilité difficiles à obtenir autrement.


Sources

  • Young, G. (2010). CQRS Documents. cqrs.files.wordpress.com.
  • Fowler, M. (2011). CQRS. martinfowler.com.
  • Fowler, M. (2005). Event Sourcing. martinfowler.com.
  • Vernon, V. (2013). Implementing Domain-Driven Design. Addison-Wesley.
  • Richardson, C. (2019). Microservices Patterns. Manning.
  • Betts, D., Dominguez, J., Melnik, G., Simonazzi, F., & Subramanian, M. (2012). Exploring CQRS and Event Sourcing. Microsoft patterns & practices.
  • Evans, E. (2003). Domain-Driven Design. Addison-Wesley.

Réservez un audit gratuit de 30 minutes. Je vous montre concrètement ce qu'on peut automatiser.