Skip to content

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

flux-order-engine

An event-driven order backbone that demonstrates the patterns a distributed order system actually needs to survive partial failure: a transactional outbox (never dual-write), idempotent consumers (at-least-once delivery, exactly-once effect), saga choreography with compensation, CQRS read models, and a dead-letter queue for poison messages.

Built with Python 3.13 · FastAPI · SQLAlchemy (async) · Kafka/Redpanda · RabbitMQ · Postgres · MongoDB · uv · Docker Compose · Ruff + mypy(strict) + pytest.


Architecture

flowchart LR
  Client -->|POST /orders| API[Order API]
  API -->|"order + outbox row<br/>(one transaction)"| PG[(Postgres)]
  Relay[Outbox Relay] -->|poll unpublished| PG
  Relay -->|publish| K{{Kafka: orders.placed}}
  K --> INV[inventory] -->|inventory.reserved| K2{{Kafka}}
  K2 --> PAY[payment] -->|payment.completed / order.confirmed| K3{{Kafka}}
  K3 --> PROJ[projector] --> MONGO[(Mongo read model)]
  Client -->|GET /orders/id| PROJ
  INV -. handler error .-> DLQ[(RabbitMQ retry → dlq)]
  PAY -. handler error .-> DLQ
Loading

The dependency arrow points inward (hexagonal / ports-and-adapters): domain/ knows nothing about FastAPI, Kafka, SQLAlchemy or Mongo. Adapters depend on the domain, never the reverse — which is what lets the saga rules and the use-case be unit-tested with zero I/O, and lets infrastructure be swapped without touching the core.

The saga (choreography — no central orchestrator)

stateDiagram-v2
  [*] --> placed: POST /orders
  placed --> inventory_reserved: inventory.reserved
  placed --> cancelled: inventory.failed (out of stock)
  inventory_reserved --> confirmed: payment.completed → order.confirmed
  inventory_reserved --> cancelled: payment.failed → inventory.released → order.cancelled
  confirmed --> [*]
  cancelled --> [*]
Loading

Each service owns exactly one transition and reacts to events; failures emit compensating events (inventory.released, order.cancelled). See ADR-0003.


The four senior moves

Pattern Problem it solves Where
Transactional outbox The dual-write bug — a crash between "save to DB" and "publish to Kafka" loses or duplicates events. place_order.py, relay.py · ADR-0001
Idempotent consumers At-least-once delivery means a consumer can see the same event twice. A (consumer, event_id) ledger makes the effect exactly-once. consumer.py
Saga choreography + compensation Distributed "transaction" across services with no 2PC; roll back via events. inventory.py, payment.py · ADR-0003
DLQ via RabbitMQ A poison message must not block a Kafka partition forever — retry with TTL, then quarantine. dlq.py · ADR-0002

CQRS read model: a projector turns the event stream into a denormalized Mongo document so queries never touch the write side — ADR-0004.


Run it

Requires Docker. Brings up Redpanda, RabbitMQ, Postgres, Mongo, the API, and the four workers (relay, inventory, payment, projector):

make up        # docker compose up -d --build
make logs      # tail the saga participants

Demo — a successful order flows through the saga

# place an order (returns 202 + order_id)
curl -s localhost:8000/orders \
  -H 'content-type: application/json' \
  -d '{"customer_id":"c1","amount_cents":4200}'

# a moment later, the read model shows status=confirmed
curl -s localhost:8000/orders/<order_id>

Or run the scripted walkthrough (places an order and polls the read model):

./scripts/demo.sh           # amount 4200 → confirmed
./scripts/demo.sh 600000    # above the payment limit → payment declined → cancelled
./scripts/demo.sh 2000000   # above the inventory limit → out of stock → cancelled

Watch make logs to see the event chain: orders.placed → inventory.reserved → payment.completed → order.confirmed, and the projector upserting Mongo.

Things to try (the failure-mode demos)

  • Idempotency: docker compose restart projector mid-flow — on restart it skips already-processed events (look for consumer.skip_duplicate in the logs) instead of double-applying them.
  • Compensation: place an order above the payment limit (600000) — payment declines and the saga emits inventory.released + order.cancelled.
  • DLQ: a handler exception routes the message to RabbitMQ's retry → dlq (queues visible at the management UI, http://localhost:15672, guest/guest) rather than blocking the partition.

Interactive API docs (OpenAPI/Swagger): http://localhost:8000/docs


Project layout

src/flux/
├── domain/            # pure: event vocabulary, saga rules, type aliases (no I/O)
│   ├── events.py      #   canonical type strings + type→status map
│   └── saga.py        #   inventory_can_reserve / payment_can_charge (pure functions)
├── application/
│   └── place_order.py #   use-case: order + outbox in ONE transaction
├── adapters/          # the outside world
│   ├── db/            #   SQLAlchemy models + async session/bootstrap
│   ├── messaging/     #   relay, idempotent consumer, publisher, RabbitMQ DLQ
│   └── read_model/    #   Mongo projector (CQRS query side)
├── services/          # saga participants: inventory, payment
├── config.py          # pydantic-settings (env prefix FLUX_)
├── logging.py         # structlog JSON logs
├── main.py            # composition root: the Order API
└── workers.py         # one entrypoint per worker role
k8s/                   # Deployments + HPA (scale inventory on Kafka consumer lag)
docs/adr/              # architecture decision records
tests/{unit,integration}/

Development

uv sync                 # install (Python 3.13)
make test               # pytest
make lint               # ruff check + format check
make typecheck          # mypy --strict on src
uv run pre-commit install

Tests

  • Unit (tests/unit/) — run anywhere, no infrastructure. Cover the pure saga rules, the event/status mapping, the projector's status derivation, and the place-order use-case (outbox written atomically) against in-memory SQLite.
  • Integration (tests/integration/) — exercise the write side against real Postgres via testcontainers; they validate server-side defaults, the JSON column under asyncpg, and the idempotency ledger's composite key. They skip automatically when Docker is unavailable, so the unit suite always stays green.

CI (.github/workflows/ci.yml) runs ruff, mypy, and pytest on every push.

Kubernetes

k8s/ holds Deployments for the API and the inventory worker, plus a HorizontalPodAutoscaler that scales inventory on Kafka consumer-group lag — the right signal for an event-driven consumer (queue depth, not CPU). Manifests are the artifact; no live cluster is required to read the intent.

Configuration

All config is 12-factor via environment (FLUX_ prefix) — see .env.example. Key vars: FLUX_PG_DSN, FLUX_KAFKA, FLUX_RABBIT, FLUX_MONGO.

Architecture decisions

Scope & honesty

This is a runnable, well-architected reference that demonstrates the signature patterns end-to-end on the happy path plus the key failure modes (compensation, idempotent replay, DLQ). It is not production-complete: the saga rules are toy thresholds, schema bootstrap uses create_all (use Alembic in prod), the relay polls (CDC/Debezium is the scale-up path), and the DLQ retry count is fixed. Each such shortcut is called out at its call site or in the relevant ADR.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages