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.