use crate::aggregate::{Aggregate, AggregateRoot, EventOf};
use crate::event::DomainEvent;
use crate::events::Events;
use crate::message::Message;
use core::fmt::Debug;
use core::hash::Hash;
pub trait Saga: Aggregate {
type CorrelationKey: Clone + Eq + Hash + Send + Sync + Debug + 'static;
type Command: Message;
fn intent_for(event: &EventOf<Self>) -> Option<Self::Command>;
}
pub trait React<E: DomainEvent, const N: usize = 0>: Saga {
fn correlate(event: &E) -> Option<Self::CorrelationKey>;
fn react(
state: &Self::State,
event: &E,
) -> Result<Option<Events<EventOf<Self>, N>>, Self::Error>;
}
impl<A: Aggregate> AggregateRoot<A> {
pub fn react<E, const N: usize>(
&self,
event: &E,
) -> Result<Option<Events<EventOf<A>, N>>, A::Error>
where
E: DomainEvent,
A: React<E, N>,
{
A::react(self.state(), event)
}
}
#[cfg(test)]
#[allow(clippy::expect_used, reason = "test code")]
mod saga_dispatch_tests {
use super::{React, Saga};
use crate::aggregate::{Aggregate, AggregateRoot, AggregateState};
use crate::event::DomainEvent;
use crate::events;
use crate::events::Events;
use crate::message::Message;
use crate::version::Version;
#[derive(Debug, Clone, Hash, PartialEq, Eq, Default)]
struct OrderId([u8; 8]);
impl OrderId {
fn new(n: u64) -> Self {
Self(n.to_le_bytes())
}
}
impl core::fmt::Display for OrderId {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(f, "{}", u64::from_le_bytes(self.0))
}
}
impl AsRef<[u8]> for OrderId {
fn as_ref(&self) -> &[u8] {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum SagaEvent {
PaymentRequested,
OrderCompleted,
}
impl Message for SagaEvent {}
impl DomainEvent for SagaEvent {
fn name(&self) -> &'static str {
match self {
Self::PaymentRequested => "PaymentRequested",
Self::OrderCompleted => "OrderCompleted",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Intent {
TakePayment,
}
impl Message for Intent {}
#[derive(Debug)]
struct OrderSagaState {
payment_requested: bool,
completed: bool,
}
impl AggregateState for OrderSagaState {
type Event = SagaEvent;
fn initial() -> Self {
Self {
payment_requested: false,
completed: false,
}
}
fn apply(mut self, event: &SagaEvent) -> Self {
match event {
SagaEvent::PaymentRequested => self.payment_requested = true,
SagaEvent::OrderCompleted => self.completed = true,
}
self
}
}
#[derive(Debug, thiserror::Error, PartialEq)]
#[error("order saga error")]
enum OrderSagaError {
EmptyOrder,
}
#[derive(Debug)]
struct OrderPlaced {
id: u64,
total: u64,
}
impl Message for OrderPlaced {}
impl DomainEvent for OrderPlaced {
fn name(&self) -> &'static str {
"OrderPlaced"
}
}
#[derive(Debug)]
struct PaymentSettled {
id: u64,
}
impl Message for PaymentSettled {}
impl DomainEvent for PaymentSettled {
fn name(&self) -> &'static str {
"PaymentSettled"
}
}
struct OrderSaga;
impl Aggregate for OrderSaga {
type State = OrderSagaState;
type Error = OrderSagaError;
type Id = OrderId;
}
impl Saga for OrderSaga {
type CorrelationKey = u64;
type Command = Intent;
fn intent_for(event: &SagaEvent) -> Option<Intent> {
match event {
SagaEvent::PaymentRequested => Some(Intent::TakePayment),
SagaEvent::OrderCompleted => None, }
}
}
impl React<OrderPlaced> for OrderSaga {
fn correlate(event: &OrderPlaced) -> Option<u64> {
Some(event.id)
}
fn react(
state: &OrderSagaState,
event: &OrderPlaced,
) -> Result<Option<Events<SagaEvent>>, OrderSagaError> {
if event.total == 0 {
return Err(OrderSagaError::EmptyOrder);
}
if state.payment_requested {
return Ok(None); }
Ok(Some(events![SagaEvent::PaymentRequested]))
}
}
impl React<PaymentSettled> for OrderSaga {
fn correlate(event: &PaymentSettled) -> Option<u64> {
Some(event.id)
}
fn react(
_state: &OrderSagaState,
_event: &PaymentSettled,
) -> Result<Option<Events<SagaEvent>>, OrderSagaError> {
Ok(Some(events![SagaEvent::OrderCompleted]))
}
}
#[test]
fn react_produces_own_events_via_dispatch() {
let root = AggregateRoot::<OrderSaga>::new(OrderId::new(1));
let produced = root
.react(&OrderPlaced { id: 1, total: 100 })
.expect("ok")
.expect("some");
assert_eq!(
produced.into_iter().collect::<Vec<_>>(),
vec![SagaEvent::PaymentRequested]
);
}
#[test]
fn react_ignores_when_no_op() {
let mut root = AggregateRoot::<OrderSaga>::new(OrderId::new(1));
root.replay(Version::INITIAL, &SagaEvent::PaymentRequested)
.expect("replay");
let outcome = root.react(&OrderPlaced { id: 1, total: 100 }).expect("ok");
assert!(outcome.is_none());
}
#[test]
fn react_surfaces_saga_error() {
let root = AggregateRoot::<OrderSaga>::new(OrderId::new(1));
assert_eq!(
root.react::<OrderPlaced, 0>(&OrderPlaced { id: 1, total: 0 }),
Err(OrderSagaError::EmptyOrder)
);
}
#[test]
fn correlate_extracts_key_per_event_type() {
assert_eq!(
<OrderSaga as React<OrderPlaced>>::correlate(&OrderPlaced { id: 42, total: 9 }),
Some(42)
);
assert_eq!(
<OrderSaga as React<PaymentSettled>>::correlate(&PaymentSettled { id: 7 }),
Some(7)
);
}
#[test]
fn intent_for_projects_events_one_to_one() {
assert_eq!(
OrderSaga::intent_for(&SagaEvent::PaymentRequested),
Some(Intent::TakePayment)
);
assert_eq!(OrderSaga::intent_for(&SagaEvent::OrderCompleted), None);
}
}