Módulo 6 · Procesamiento asíncrono y mensajería

Lección 32 — Consistencia eventual y saga

El pago, la reserva y el correo: tres sistemas, una historia coherente.

Publicada
En esta lección
  1. Ejercicio 1 — El mapa
  2. Ejercicio 2 — La saga
  3. Ejercicio 3 — El UNKNOWN
  4. Ejercicio 4 — Compensaciones
  5. Ejercicio 5 — Coreografía y E2E
  6. Resumen del profesor

Ejercicio 1 — El mapa

OperaciónNivelConvergencia
Disponibilidad de asientosFUERTE (lectura en TX con lock, 10)—
Reserva + asientos + outboxFUERTE (una TX, 10/25)—
CobroFUERTE-en-su-sistema; saga local convergereconciliación por intent (14)
Email confirmacióneventualat-least-once + dedup, minutos (29)
Webhooks salienteseventualbackoff + DLQ, horas (17/30)
Dashboard ventaseventualstream → consumidor, minutos
Informe RGPDeventualjob diario (31), 24 h

El finding clásico: el email de confirmación en el request (el SMTP de 300 ms a 3 s en el p95 de la compra) y, peor, el webhook saliente síncrono (el timeout de 5 s del tercero en tu p99). Ambos a la cola con on_commit. El pecado inverso (eventual donde esperas): la disponibilidad cacheada (12) usada para DECIDIR la compra — el caché informa el listado, la TX con lock decide el asiento.

Ejercicio 2 — La 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: estado válido → estados destino
            raise DomainError(f"transición inválida {self.estado} → {nuevo}")
        self.estado = nuevo
        self.save(update_fields=["estado", "updated_at"])

El orquestador:

python
def iniciar_compra(user, event, seat_refs, *, clock, gateway):
    with transaction.atomic():
        r = reservar(user, event, seat_refs, clock=clock)      # TX fuerte (10)
        saga = PurchaseSaga.objects.create(reserva=r, intent_id=f"int_{r.public_ref}")
        transaction.on_commit(lambda: cobrar.delay(saga.saga_id))  # 29: nunca el email/cobro fantasma
    return saga
  1. El rollback antes de encolar: el test de la 29 (email fantasma) aplica al cobro — el DomainError de asiento duplicado deja cero filas saga y cero tareas. 3. La reanudación: reanudar_sagas() filtra estado=PENDING, updated_at < now - 5min, y llama al paso de cobro con el intent_id EXISTENTE — la pasarela fake con Idempotency-Key devuelve el cargo original: cero doble cobro. La tabla es la verdad; el proceso, desechable.

Ejercicio 3 — El 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 sin resolver: la tarea relanza con backoff (29); al agotar, runbook
  1. Los tres casos del test: la pasarela fake configurable (fake.charge_succeeds_once_then_timeout) produce cobró-declined-sigue-unknown. 2. El plazo: 6 reconciliaciones × backoff 1-2-4-8-16-32 min ≈ 1 h; al expirar: la saga queda UNKNOWN_EXPIRED y el runbook — (1) alerta al on-call con saga_id; (2) consultar la pasarela por panel/soporte; (3) si cobró: confirmar manualmente y anotar en el ledger; (4) si no: compensar y postmortem (47). El dinero en el limbo > 1 h se paga con soporte: el plazo es la promesa interna.
  1. El test rojo del pecado:
python
def test_compensar_en_ciego_pierde_dinero(self):
    gateway.que_timeout_pero_cobra()          # el cargo REALMENTE entró
    with self.assertRaises(AssertionError):
        compensar_sin_reconciliar(saga)        # reembolsa y libera: el asiento se vendió 2×

La advertencia documentada: compensar ante UNKNOWN asume "no cobró" sin evidencia — el caso peor (cobró Y liberaste) es el único que convierte un incidente técnico en pérdida neta.

Ejercicio 4 — Compensaciones

python
def compensar(saga):
    if saga.estado == PurchaseSaga.Estado.COMPENSATED:
        return                                  # idempotente por estado (máquina de estados)
    saga.transicionar(PurchaseSaga.Estado.COMPENSATED)
    with transaction.atomic():
        LedgerEntry.reembolsar(saga.reserva)    # entrada inversa (08); única por (intent, kind)
        saga.reserva.cancelar(motivo="saga_compensated")   # máquina de estados (00b)
        outbox.create(event_type="ReservationCancelled", payload={...})
    transaction.on_commit(lambda: email_fallo.delay(saga.reserva.public_ref))

Dos llamadas → la segunda sale por el guard del estado: un reembolso. El orden inverso: reembolso → liberar asientos → email; la pila deshace lo último primero porque los pasos tardíos dependen de los tempranos (no puedes liberar el asiento de una reserva que el refund va a anular... o sí, pero el EMAIL de fallo debe salir cuando el estado ya es estable). Los efectos no compensables (el email de "compra iniciada" ya salió) no se desenvían: se sobre-corre con el email de fallo — el usuario ve la corrección, no el borrado. 3. El perseguidor: filter(estado=COMPENSATION_PENDING, updated_at < now-24h) → métrica + alerta; con FakeClock el test avanza 25 h sin dormir.

Ejercicio 5 — Coreografía y E2E

  1. El consumidor de analytics acumula PaymentSucceeded en su propia tabla (o even better: materializa desde el ledger): si analytics cae 2 h, al volver el XAUTOCLAIM (30) entrega el backlog y converge — el test: pausa el consumidor (no drena), publica 5 eventos, reactiva → las ventas del día terminan correctas. Nada del flujo de compra lo espera: por eso es coreografía.
  1. La E2E de la saga (tres escenarios, un test de historia):
python
def test_historia_completa(self):
    saga = iniciar_compra(self.user, self.event, ["A1"], ...)     # feliz
    procesar_colas()      # eager + drena las tareas
    self.assertEqual(saga.refresh.estado, CONFIRMED); self.assertEqual(len(mailbox), 1)

def test_historia_declinada(self):
    self.gateway.declina_todo = True
    ... → PAYMENT_FAILED + compensada + email de fallo + asiento A1 libre

def test_historia_timeout_cobro(self):
    self.gateway.timeout_pero_cobra = True
    ... → reconciliada → CONFIRMED + UN cargo (no dos) + asiento vendido

El tercer escenario es el que valida el valor entero de la saga: el sistema se recupera del "no sé" sin humano y sin pérdida.


Resumen del profesor

  • Fuerte para dinero e inventario (TX local), eventual para lo notificable, con convergencia definida y vigilada.
  • La saga orquestada tiene su estado en una tabla (reanudable) y su borde remoto con tres salidas; ante UNKNOWN: reconciliar, jamás compensar en ciego.
  • La compensación es un hecho inverso, idempotente, en orden inverso y con plazo; los efectos no compensables se sobre-corren, no se borran.