Module 6 · Async processing and messaging

Lesson 32 — Eventual consistency and sagas

The payment, the reservation and the email: three systems, one coherent story.

Published
In this lesson
  1. Exercise 1 — The map
  2. Exercise 2 — The saga
  3. Exercise 3 — The UNKNOWN
  4. Exercise 4 — Compensations
  5. Exercise 5 — Choreography and E2E
  6. Professor's summary

Exercise 1 — The map

OperationLevelConvergence
Seat availabilitySTRONG (read in TX with lock, 10)—
Reservation + seats + outboxSTRONG (one TX, 10/25)—
ChargeSTRONG-in-its-system; the local saga convergesreconciliation by intent (14)
Confirmation emaileventualat-least-once + dedup, minutes (29)
Outgoing webhookseventualbackoff + DLQ, hours (17/30)
Sales dashboardeventualstream → consumer, minutes
GDPR reporteventualdaily job (31), 24 h

The classic finding: the confirmation email in the request (the 300 ms-to-3 s SMTP inside the purchase's p95) and, worse, the synchronous outgoing webhook (the third party's 5 s timeout in your p99). Both to the queue with on_commit. The inverse sin (eventual where you wait): the cached availability (12) used to DECIDE the purchase — the cache informs the listing, the TX with lock decides the seat.

Exercise 2 — The saga

python
class PurchaseSaga(models.Model):
    class Estado(models.TextChoices):
        PENDING = "PENDING"; PAID = "PAID"; CONFIRMED = "CONFIRMED"
        PAYMENT_FAILED = "PAYMENT_FAILED"; COMPENSATED = "COMPENSATED"

    saga_id = models.UUIDField(default=uuid.uuid4, editable=False, unique=True)
    reserva = models.ForeignKey("Reservation", on_delete=models.PROTECT)
    intent_id = models.CharField(max_length=64)
    estado = models.CharField(max_length=20, choices=Estado.choices, default=Estado.PENDING)
    intentos = models.IntegerField(default=0)

    def transicionar(self, nuevo):
        if nuevo not in self.TRANSICIONES[self.estado]:   # dict: valid state → target states
            raise DomainError(f"transición inválida {self.estado} → {nuevo}")
        self.estado = nuevo
        self.save(update_fields=["estado", "updated_at"])

The orchestrator:

python
def iniciar_compra(user, event, seat_refs, *, clock, gateway):
    with transaction.atomic():
        r = reservar(user, event, seat_refs, clock=clock)      # strong TX (10)
        saga = PurchaseSaga.objects.create(reserva=r, intent_id=f"int_{r.public_ref}")
        transaction.on_commit(lambda: cobrar.delay(saga.saga_id))  # 29: never the ghost email/charge
    return saga
  1. The rollback before enqueueing: 29's test (ghost email) applies to the charge — the duplicate-seat DomainError leaves zero saga rows and zero tasks. 3. The resume: reanudar_sagas() filters estado=PENDING, updated_at < now - 5min, and calls the charge step with the EXISTING intent_id — the fake gateway with the Idempotency-Key returns the original charge: zero double charges. The table is the truth; the process, disposable.

Exercise 3 — The UNKNOWN

python
def paso_cobrar(saga_id, *, gateway, clock, max_reconcilia=6):
    saga = PurchaseSaga.objects.get(saga_id=saga_id)
    try:
        out = gateway.cobrar(saga.intent_id, saga.reserva.total())
    except GatewayTimeout:
        out = reconciliar(saga, gateway, clock, max_reconcilia)
    if out.status == "SUCCEEDED":
        saga.transicionar(PurchaseSaga.Estado.CONFIRMED); publicar(PaymentSucceeded(...))
    elif out.status == "DECLINED":
        saga.transicionar(PurchaseSaga.Estado.PAYMENT_FAILED); compensar(saga)
    # UNKNOWN unresolved: the task relaunches with backoff (29); when exhausted, runbook
  1. The test's three cases: the configurable fake gateway (fake.charge_succeeds_once_then_timeout) produces charged-declined-still-unknown. 2. The deadline: 6 reconciliations × backoff 1-2-4-8-16-32 min ≈ 1 h; on expiry: the saga stays UNKNOWN_EXPIRED and the runbook — (1) alert the on-call with saga_id; (2) check the gateway via panel/support; (3) if it charged: confirm manually and note it in the ledger; (4) if not: compensate and postmortem (47). Money in limbo > 1 h gets paid for in support: the deadline is the internal promise.
  1. The sin's red test:
python
def test_compensar_en_ciego_pierde_dinero(self):
    gateway.que_timeout_pero_cobra()          # the charge REALLY went through
    with self.assertRaises(AssertionError):
        compensar_sin_reconciliar(saga)        # refunds and frees: the seat sold 2×

The documented warning: compensating on UNKNOWN assumes "it didn't charge" without evidence — the worst case (it charged AND you freed the seat) is the only one that turns a technical incident into a net loss.

Exercise 4 — Compensations

python
def compensar(saga):
    if saga.estado == PurchaseSaga.Estado.COMPENSATED:
        return                                  # idempotent by state (state machine)
    saga.transicionar(PurchaseSaga.Estado.COMPENSATED)
    with transaction.atomic():
        LedgerEntry.reembolsar(saga.reserva)    # inverse entry (08); unique per (intent, kind)
        saga.reserva.cancelar(motivo="saga_compensated")   # state machine (00b)
        outbox.create(event_type="ReservationCancelled", payload={...})
    transaction.on_commit(lambda: email_fallo.delay(saga.reserva.public_ref))

Two calls → the second exits through the state guard: one refund. The reverse order: refund → free seats → email; the stack undoes the latest first because the late steps depend on the early ones (you can't free the seat of a reservation the refund will annul... or you can, but the FAILURE email must go out when the state is already stable). Non-compensable effects (the "purchase started" email already went out) are not un-sent: they are over-corrected with the failure email — the user sees the correction, not the deletion. 3. The chaser: filter(estado=COMPENSATION_PENDING, updated_at < now-24h) → metric + alert; with FakeClock the test advances 25 h without sleeping.

Exercise 5 — Choreography and E2E

  1. The analytics consumer accumulates PaymentSucceeded in its own table (or better still: materializes from the ledger): if analytics is down 2 h, on return the XAUTOCLAIM (30) delivers the backlog and it converges — the test: pause the consumer (it doesn't drain), publish 5 events, reactivate → the day's sales end up correct. Nothing in the purchase flow waits for it: that is why it is choreography.
  1. The saga's E2E (three scenarios, one story test):
python
def test_historia_completa(self):
    saga = iniciar_compra(self.user, self.event, ["A1"], ...)     # happy
    procesar_colas()      # eager + drain the tasks
    self.assertEqual(saga.refresh.estado, CONFIRMED); self.assertEqual(len(mailbox), 1)

def test_historia_declinada(self):
    self.gateway.declina_todo = True
    ... → PAYMENT_FAILED + compensated + failure email + seat A1 free

def test_historia_timeout_cobro(self):
    self.gateway.timeout_pero_cobra = True
    ... → reconciled → CONFIRMED + ONE charge (not two) + seat sold

The third scenario validates the saga's entire value: the system recovers from the "don't know" without a human and without loss.


Professor's summary

  • Strong for money and inventory (local TX), eventual for the notifiable, with defined and watched convergence.
  • The orchestrated saga keeps its state in a table (resumable) and its remote boundary has three outcomes; on UNKNOWN: reconcile, never compensate blind.
  • The compensation is an inverse fact, idempotent, in reverse order and with a deadline; non-compensable effects are over-corrected, not deleted.