Sagas y Process Managers
Procesos de larga duración, compensaciones y coordinación entre contextos.
Una Saga coordina un proceso de negocio de larga duración que abarca múltiples aggregates o bounded contexts, asegurando consistencia a través de fronteras sin transacciones distribuidas. En lugar de un lock global, cada paso del proceso es una operación local y autónoma. Si el proceso falla a mitad, las Sagas revierten el trabajo completado mediante transacciones de compensación — acciones de dominio explícitas que deshacen el efecto de pasos anteriores.
Existen dos estilos de implementación. En la coreografía, los aggregates reaccionan directamente a los eventos de otros aggregates sin coordinador central: cada participante sabe qué hacer cuando recibe cierto evento. En la orquestación, un Process Manager centralizado recibe eventos y envía comandos para coordinar a los participantes: el estado del proceso vive explícitamente en el Process Manager, lo que facilita la trazabilidad y el diagnóstico.
Las transacciones de compensación son el mecanismo de 'rollback' de las Sagas. No son rollbacks de base de datos — son nuevas acciones de dominio que deshacen el efecto de pasos previos. Cada paso que muta estado debe tener diseñado explícitamente su inverso antes de implementar el flujo hacia adelante. La compensación es eventual, no atómica.
El Process Manager es un coordinador con estado. Su estado debe persistirse antes de enviar cualquier comando: si el proceso crashea entre pasos, debe poder reanudarse desde donde estaba. El Process Manager reacciona a eventos y decide el siguiente comando — nunca llama métodos de aggregates directamente.
Cuándo NO usar Sagas: añaden complejidad de estado distribuido, compensación explícita e idempotencia obligatoria. Para flujos simples dentro de un único bounded context, un Application Service con una transacción es suficiente y mucho más simple. Solo introduce Sagas para flujos genuinamente multi-contexto con más de dos pasos y necesidad real de compensación.
# application/sagas/<process>_saga.py — Saga / Process Manager
# application/sagas/<process>_state.py — Saga State (persisted)
@dataclass
class <Process>SagaState:
saga_id: str
step: str # 'pending' | 'step_a_done' | 'step_b_done' | 'completed' | 'compensating'
<step_a>_result: Optional[str] = None
<step_b>_result: Optional[str] = None
class <Process>Saga:
def __init__(self, cmd_bus: <CommandBus>, saga_store: <SagaStore>) -> None:
self._cmd_bus = cmd_bus
self._store = saga_store
def on_<process_started>(self, event: <ProcessStarted>) -> None:
state = <Process>SagaState(saga_id=event.saga_id, step='pending')
self._store.save(state)
self._cmd_bus.send(<DoStepA>(saga_id=event.saga_id, ...))
def on_<step_a_completed>(self, event: <StepACompleted>) -> None:
state = self._store.load(event.saga_id)
if state.step == 'step_a_done': # Idempotency check
return
state.step = 'step_a_done'
state.<step_a>_result = event.result
self._store.save(state)
self._cmd_bus.send(<DoStepB>(saga_id=event.saga_id, ...))
def on_<step_b_failed>(self, event: <StepBFailed>) -> None:
state = self._store.load(event.saga_id)
state.step = 'compensating'
self._store.save(state)
# Compensate step A because step B failed
self._cmd_bus.send(<UndoStepA>(saga_id=event.saga_id,
result=state.<step_a>_result))
def on_<step_a_undone>(self, event: <StepAUndone>) -> None:
state = self._store.load(event.saga_id)
state.step = 'compensated'
self._store.save(state)Debugging lab
Detecta y corrige el error o la violación de diseño.
- 7.3.5.1
# Bug: saga sends DoStepB command before persisting saga state after StepACompleted class ProcessSaga: def on_step_a_completed(self, event: StepACompleted) -> None: state = self._store.load(event.saga_id) # Bug: sending command BEFORE persisting updated state self._cmd_bus.send(DoStepB(saga_id=event.saga_id, data=event.result)) state.step = "step_a_done" state.step_a_result = event.result self._store.save(state)
- 7.3.5.2
# Bug: saga has no compensating command for step A when step B fails class ProcessSaga: def on_step_b_failed(self, event: StepBFailed) -> None: state = self._store.load(event.saga_id) # Bug: just marking as failed with no compensation state.step = "failed" self._store.save(state) # Step A was already completed — its effect is permanent, no undo
- 7.3.5.3
# Bug: saga directly calling aggregate method instead of sending a command class ProcessSaga: def on_step_b_failed(self, event: StepBFailed) -> None: state = self._store.load(event.saga_id) # Bug: directly calling aggregate method — tight coupling aggregate = self._aggregate_repo.find(AggregateId(state.aggregate_id)) aggregate.undo_step_a(state.step_a_result) self._aggregate_repo.save(aggregate) state.step = "compensating" self._store.save(state)
- 7.3.5.4
# Bug: saga handling on_step_a_completed without checking if already in that step (non-idempotent) class ProcessSaga: def on_step_a_completed(self, event: StepACompleted) -> None: state = self._store.load(event.saga_id) # Bug: no idempotency check — processing same event twice sends DoStepB twice state.step = "step_a_done" state.step_a_result = event.result self._store.save(state) self._cmd_bus.send(DoStepB(saga_id=event.saga_id, data=event.result))
- 7.3.5.5
# Bug: orchestration saga introduced for a 2-step flow within a single bounded context # Both steps touch aggregates in the same BC, same database, same transaction scope class TwoStepSaga: def on_action_started(self, event: ActionStarted) -> None: state = TwoStepSagaState(saga_id=event.saga_id, step="pending") self._store.save(state) self._cmd_bus.send(DoStepOne(saga_id=event.saga_id)) def on_step_one_done(self, event: StepOneDone) -> None: state = self._store.load(event.saga_id) state.step = "step_one_done" self._store.save(state) self._cmd_bus.send(DoStepTwo(saga_id=event.saga_id))