Implement event sourcing pattern where state is derived from an immutable sequence of events. Outputs event store design, aggregate patterns, projection builders, and command/event handlers.
Implement event sourcing pattern where state is derived from an immutable sequence of events. Outputs event store design, aggregate patterns, projection builders, and command/event handlers.
argument-hint
["domain","aggregate types","projection requirements","event store technology"]
allowed-tools
Read, Write, Bash
Event Sourcing
Event sourcing stores every state change as an immutable event rather than overwriting current state. The current state is always derived by replaying events. This enables a complete audit log, temporal queries, and the ability to rebuild any projection from scratch.
When to Use Event Sourcing
Use when:
Complete audit log is required (finance, healthcare, compliance)
Temporal queries: "what was the state at 3pm last Tuesday?"
Multiple read models needed from the same write data
Complex domain with many state transitions
Don't use when:
Simple CRUD with no history requirements
Team unfamiliar with the pattern — learning curve is steep
High write throughput + simple queries — overhead isn't worth it
Process
Define the domain events — past-tense, immutable facts (OrderPlaced, PaymentProcessed).
Design aggregates — domain objects that enforce invariants and emit events.
Implement the event store — append-only log with optimistic concurrency.
Build projections — read models derived from event streams.
# domain/order_aggregate.pyfrom typing importOptionalfrom dataclasses import dataclass, field
from .events import *
classInvalidStateTransition(Exception):
passclassOrderAggregate:
"""
Order aggregate — enforces business invariants.
State is derived entirely from replayed events.
Never persisted directly — only events are stored.
"""def__init__(self, order_id: str):
self.order_id = order_id
self.version = 0# Optimistic concurrency versionself._events: list[DomainEvent] = [] # Uncommitted events# State derived from eventsself.status: Optional[str] = Noneself.user_id: Optional[str] = Noneself.items: list = []
self.total_cents: int = 0self.payment_id: Optional[str] = Noneself.is_paid: bool = False# ── Command handlers ──────────────────────────────────defplace(self, user_id: str, items: list, total_cents: int, shipping_address: dict):
"""Command: place a new order."""ifself.status isnotNone:
raise InvalidStateTransition(f"Cannot place order in state: {self.status}")
# Validate invariantsifnot items:
raise ValueError("Order must have at least one item")
if total_cents <= 0:
raise ValueError("Order total must be positive")
self._apply(OrderPlaced(
aggregate_id=self.order_id,
user_id=user_id,
items=tuple(items),
total_cents=total_cents,
shipping_address=shipping_address,
))
defprocess_payment(self, payment_provider: str, payment_id: str, amount_cents: int):
ifself.status != "pending":
raise InvalidStateTransition(f"Cannot process payment in state: {self.status}")
self._apply(PaymentProcessed(
aggregate_id=self.order_id,
payment_provider=payment_provider,
payment_id=payment_id,
amount_cents=amount_cents,
))
deffail_payment(self, reason: str, error_code: str):
ifself.status != "pending":
raise InvalidStateTransition(f"Cannot fail payment in state: {self.status}")
self._apply(PaymentFailed(
aggregate_id=self.order_id,
reason=reason,
error_code=error_code,
))
defship(self, carrier: str, tracking_number: str, estimated_delivery: str):
ifself.status != "paid":
raise InvalidStateTransition(f"Cannot ship order in state: {self.status}")
self._apply(OrderShipped(
aggregate_id=self.order_id,
carrier=carrier,
tracking_number=tracking_number,
estimated_delivery=estimated_delivery,
))
defcancel(self, reason: str, cancelled_by: str):
ifself.status in ("shipped", "delivered", "cancelled"):
raise InvalidStateTransition(f"Cannot cancel order in state: {self.status}")
self._apply(OrderCancelled(
aggregate_id=self.order_id,
reason=reason,
cancelled_by=cancelled_by,
))
# ── Event application (state mutations) ───────────────def_apply(self, event: DomainEvent, is_replay: bool = False):
"""Apply an event — mutates state. Called for both new events and replays."""
handler = {
OrderPlaced: self._on_order_placed,
PaymentProcessed: self._on_payment_processed,
PaymentFailed: self._on_payment_failed,
OrderShipped: self._on_order_shipped,
OrderCancelled: self._on_order_cancelled,
}.get(type(event))
if handler:
handler(event)
self.version += 1ifnot is_replay:
self._events.append(event) # Buffer uncommitted eventsdef_on_order_placed(self, event: OrderPlaced):
self.status = "pending"self.user_id = event.user_id
self.items = list(event.items)
self.total_cents = event.total_cents
def_on_payment_processed(self, event: PaymentProcessed):
self.status = "paid"self.payment_id = event.payment_id
self.is_paid = Truedef_on_payment_failed(self, event: PaymentFailed):
self.status = "payment_failed"def_on_order_shipped(self, event: OrderShipped):
self.status = "shipped"def_on_order_cancelled(self, event: OrderCancelled):
self.status = "cancelled"# ── Reconstruction from events ──────────────────────── @classmethoddefload(cls, order_id: str, events: list[DomainEvent]) -> 'OrderAggregate':
"""Reconstruct aggregate state by replaying all historical events."""
aggregate = cls(order_id)
for event in events:
aggregate._apply(event, is_replay=True)
return aggregate
defuncommitted_events(self) -> list[DomainEvent]:
returnlist(self._events)
defmark_committed(self):
self._events.clear()
Event Store
# infrastructure/event_store.pyimport json
import asyncpg
from datetime import datetime, timezone
from typing importTypeclassOptimisticConcurrencyError(Exception):
passclassEventStore:
"""
Append-only event store backed by PostgreSQL.
Optimistic concurrency via expected_version.
"""def__init__(self, pool: asyncpg.Pool):
self.pool = pool
asyncdefappend_events(
self,
aggregate_id: str,
events: list,
expected_version: int,
):
"""
Append events to the store.
expected_version: the version the caller believes the aggregate is at.
Raises OptimisticConcurrencyError if another writer has appended events since.
"""asyncwithself.pool.acquire() as conn:
asyncwith conn.transaction():
# Check current version (optimistic lock)
current_version = await conn.fetchval(
"SELECT COALESCE(MAX(version), 0) FROM events WHERE aggregate_id = $1",
aggregate_id
)
if current_version != expected_version:
raise OptimisticConcurrencyError(
f"Concurrency conflict: expected version {expected_version}, "f"current version {current_version}"
)
# Append eventsfor i, event inenumerate(events):
version = expected_version + i + 1await conn.execute(
"""
INSERT INTO events (
event_id, aggregate_id, version,
event_type, event_data, occurred_at
) VALUES ($1, $2, $3, $4, $5, $6)
""",
event.event_id,
aggregate_id,
version,
type(event).__name__,
json.dumps(self._serialize(event)),
event.occurred_at,
)
asyncdefload_events(self, aggregate_id: str, from_version: int = 0) -> list:
"""Load all events for an aggregate, optionally from a specific version."""asyncwithself.pool.acquire() as conn:
rows = await conn.fetch(
"""
SELECT event_type, event_data, version
FROM events
WHERE aggregate_id = $1 AND version > $2
ORDER BY version ASC
""",
aggregate_id, from_version
)
return [self._deserialize(row["event_type"], row["event_data"]) for row in rows]
asyncdefload_events_by_type(
self, event_type: str, after: datetime = None, limit: int = 1000) -> list:
"""Load events of a specific type — for projections."""asyncwithself.pool.acquire() as conn:
query = "SELECT event_data, occurred_at FROM events WHERE event_type = $1"
params = [event_type]
if after:
query += " AND occurred_at > $2"
params.append(after)
query += f" ORDER BY occurred_at ASC LIMIT {limit}"
rows = await conn.fetch(query, *params)
return [self._deserialize(event_type, row["event_data"]) for row in rows]
def_serialize(self, event) -> dict:
"""Convert event to storable dict."""import dataclasses
return {k: v for k, v in dataclasses.asdict(event).items()
if k notin ("event_id", "occurred_at", "aggregate_id", "aggregate_version")}
def_deserialize(self, event_type: str, data: str) -> object:
"""Reconstruct event from stored data."""
event_classes = {
"OrderPlaced": OrderPlaced,
"PaymentProcessed": PaymentProcessed,
"PaymentFailed": PaymentFailed,
"OrderShipped": OrderShipped,
"OrderCancelled": OrderCancelled,
}
cls = event_classes.get(event_type)
ifnot cls:
raise ValueError(f"Unknown event type: {event_type}")
return cls(**json.loads(data))
# Schema
CREATE_EVENTS_TABLE = """
CREATE TABLE IF NOT EXISTS events (
event_id UUID PRIMARY KEY,
aggregate_id VARCHAR(255) NOT NULL,
version INTEGER NOT NULL,
event_type VARCHAR(255) NOT NULL,
event_data JSONB NOT NULL,
occurred_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE(aggregate_id, version) -- Optimistic concurrency constraint
);
CREATE INDEX idx_events_aggregate ON events(aggregate_id, version);
CREATE INDEX idx_events_type ON events(event_type, occurred_at);
CREATE INDEX idx_events_occurred ON events(occurred_at);
"""
Projections
# projections/order_read_model.pyclassOrderProjection:
"""
Read model for order queries.
Built by processing event stream — completely denormalized for fast reads.
Can be rebuilt from scratch by replaying all events.
"""def__init__(self, db):
self.db = db
asyncdefhandle(self, event):
"""Route events to appropriate handlers."""
handlers = {
"OrderPlaced": self.on_order_placed,
"PaymentProcessed": self.on_payment_processed,
"PaymentFailed": self.on_payment_failed,
"OrderShipped": self.on_order_shipped,
"OrderCancelled": self.on_order_cancelled,
}
handler = handlers.get(type(event).__name__)
if handler:
await handler(event)
asyncdefon_order_placed(self, event: OrderPlaced):
awaitself.db.execute(
"""
INSERT INTO order_read_model (
order_id, user_id, status, total_cents,
items, created_at, updated_at
) VALUES ($1, $2, 'pending', $3, $4, $5, $5)
ON CONFLICT (order_id) DO NOTHING
""",
event.aggregate_id, event.user_id, event.total_cents,
json.dumps(list(event.items)), event.occurred_at,
)
asyncdefon_payment_processed(self, event: PaymentProcessed):
awaitself.db.execute(
"UPDATE order_read_model SET status='paid', payment_id=$2, updated_at=$3 WHERE order_id=$1",
event.aggregate_id, event.payment_id, event.occurred_at,
)
asyncdefon_order_cancelled(self, event: OrderCancelled):
awaitself.db.execute(
"UPDATE order_read_model SET status='cancelled', cancel_reason=$2, updated_at=$3 WHERE order_id=$1",
event.aggregate_id, event.reason, event.occurred_at,
)
asyncdefrebuild(self, event_store: EventStore):
"""Rebuild entire read model from event history."""awaitself.db.execute("TRUNCATE order_read_model")
all_events = await event_store.load_events_by_type("OrderPlaced", limit=100000)
# Then load other event types and merge by time...# (In practice: use a single sorted stream of all events)for event in all_events:
awaitself.handle(event)
print(f"Rebuilt projection from {len(all_events)} events")
Rules
Events are past-tense, immutable facts — OrderPlaced, not PlaceOrder; never modify stored events.
Aggregates emit events, never mutate directly — all state changes flow through event application.
Optimistic concurrency at the event store — the expected_version check prevents lost updates.
Projections are disposable — they're derived; rebuild them anytime from the event log.
Separate command model from query model — aggregates handle writes; projections handle reads (CQRS).
Event schema migration is hard — version your events from day one (OrderPlaced.v2).
Snapshots for long-lived aggregates — store periodic snapshots to avoid replaying 10,000 events on every load.
Idempotent event handlers — projections may receive the same event twice; handle duplicates gracefully.
Don't put behavior in projections — they transform events to read models; no business logic.
Test by replaying — event sourcing tests replay events and assert on derived state, not mocked calls.