use super::order::{
Inventory, Order, OrderFulfillmentSaga, OrderItem, OrderStatus, Payment, PaymentStatus,
SagaStatus,
};
use distributed::{AggregateBuilder, HashMapRepository};
#[tokio::test]
async fn saga_happy_path_completes_order() {
let order_repo = HashMapRepository::new().aggregate::<Order>();
let inventory_repo = HashMapRepository::new().aggregate::<Inventory>();
let payment_repo = HashMapRepository::new().aggregate::<Payment>();
let saga_repo = HashMapRepository::new().aggregate::<OrderFulfillmentSaga>();
let mut widget_inventory = Inventory::new();
widget_inventory
.initialize("WIDGET-001".to_string(), 100)
.unwrap();
inventory_repo.commit(&mut widget_inventory).await.unwrap();
let order_id = "order-123".to_string();
let items = vec![OrderItem {
sku: "WIDGET-001".to_string(),
quantity: 5,
price_cents: 1000,
}];
let mut order = Order::new();
order
.create(order_id.clone(), "customer-456".to_string(), items.clone())
.unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut order_fulfillment_saga = OrderFulfillmentSaga::new();
order_fulfillment_saga
.start(
"saga-123".to_string(),
order_id.clone(),
"customer-456".to_string(),
items,
5000, )
.unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Started);
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut inventory = inventory_repo.get("WIDGET-001").await.unwrap().unwrap();
assert!(inventory.can_reserve(5));
inventory.reserve(order_id.clone(), 5).unwrap();
inventory_repo.commit(&mut inventory).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-123").await.unwrap().unwrap();
order_fulfillment_saga.inventory_reserved().unwrap();
assert_eq!(
order_fulfillment_saga.status(),
SagaStatus::InventoryReserved
);
assert!(order_fulfillment_saga.compensation().inventory_reserved);
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut order = order_repo.get(&order_id).await.unwrap().unwrap();
order.mark_inventory_reserved().unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut payment = Payment::new();
payment
.initiate("payment-789".to_string(), order_id.clone(), 5000)
.unwrap();
payment.authorize("txn-abc123".to_string()).unwrap();
payment.capture().unwrap();
assert!(payment.is_successful());
payment_repo.commit(&mut payment).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-123").await.unwrap().unwrap();
order_fulfillment_saga.payment_succeeded().unwrap();
assert_eq!(
order_fulfillment_saga.status(),
SagaStatus::PaymentProcessed
);
assert!(order_fulfillment_saga.compensation().payment_processed);
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut order = order_repo.get(&order_id).await.unwrap().unwrap();
order.mark_payment_processed().unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-123").await.unwrap().unwrap();
order_fulfillment_saga.complete().unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Completed);
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut inventory = inventory_repo.get("WIDGET-001").await.unwrap().unwrap();
inventory.commit_reservation(order_id.clone()).unwrap();
inventory_repo.commit(&mut inventory).await.unwrap();
let mut order = order_repo.get(&order_id).await.unwrap().unwrap();
order.complete().unwrap();
order_repo.commit(&mut order).await.unwrap();
let final_order_fulfillment_saga = saga_repo.get("saga-123").await.unwrap().unwrap();
assert_eq!(final_order_fulfillment_saga.status(), SagaStatus::Completed);
assert!(final_order_fulfillment_saga.is_complete());
let final_order = order_repo.get(&order_id).await.unwrap().unwrap();
assert_eq!(final_order.status(), OrderStatus::Completed);
let final_inventory = inventory_repo.get("WIDGET-001").await.unwrap().unwrap();
assert_eq!(final_inventory.available(), 95); assert_eq!(final_inventory.reserved(), 0);
assert!(final_inventory.reservation_for_order(&order_id).is_none());
}
#[tokio::test]
async fn saga_compensates_on_payment_failure() {
let order_repo = HashMapRepository::new().aggregate::<Order>();
let inventory_repo = HashMapRepository::new().aggregate::<Inventory>();
let payment_repo = HashMapRepository::new().aggregate::<Payment>();
let saga_repo = HashMapRepository::new().aggregate::<OrderFulfillmentSaga>();
let mut widget_inventory = Inventory::new();
widget_inventory
.initialize("WIDGET-002".to_string(), 50)
.unwrap();
inventory_repo.commit(&mut widget_inventory).await.unwrap();
let order_id = "order-fail-456".to_string();
let items = vec![OrderItem {
sku: "WIDGET-002".to_string(),
quantity: 10,
price_cents: 500,
}];
let mut order = Order::new();
order
.create(order_id.clone(), "customer-789".to_string(), items.clone())
.unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut order_fulfillment_saga = OrderFulfillmentSaga::new();
order_fulfillment_saga
.start(
"saga-fail-456".to_string(),
order_id.clone(),
"customer-789".to_string(),
items,
5000,
)
.unwrap();
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut inventory = inventory_repo.get("WIDGET-002").await.unwrap().unwrap();
inventory.reserve(order_id.clone(), 10).unwrap();
inventory_repo.commit(&mut inventory).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-fail-456").await.unwrap().unwrap();
order_fulfillment_saga.inventory_reserved().unwrap();
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut order = order_repo.get(&order_id).await.unwrap().unwrap();
order.mark_inventory_reserved().unwrap();
order_repo.commit(&mut order).await.unwrap();
let inventory = inventory_repo.get("WIDGET-002").await.unwrap().unwrap();
assert_eq!(inventory.available(), 40); assert_eq!(inventory.reserved(), 10);
let mut payment = Payment::new();
payment
.initiate("payment-fail-xyz".to_string(), order_id.clone(), 5000)
.unwrap();
payment.fail("Insufficient funds".to_string()).unwrap();
assert!(!payment.is_successful());
assert_eq!(payment.status(), PaymentStatus::Failed);
payment_repo.commit(&mut payment).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-fail-456").await.unwrap().unwrap();
order_fulfillment_saga
.step_failed("Payment".to_string(), "Insufficient funds".to_string())
.unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Compensating);
assert!(order_fulfillment_saga.needs_inventory_compensation());
assert!(!order_fulfillment_saga.needs_payment_compensation()); saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut inventory = inventory_repo.get("WIDGET-002").await.unwrap().unwrap();
inventory.release_reservation(order_id.clone()).unwrap();
inventory_repo.commit(&mut inventory).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-fail-456").await.unwrap().unwrap();
order_fulfillment_saga.inventory_compensated().unwrap();
assert!(!order_fulfillment_saga.needs_inventory_compensation());
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut order = order_repo.get(&order_id).await.unwrap().unwrap();
order
.cancel("Payment failed: Insufficient funds".to_string())
.unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-fail-456").await.unwrap().unwrap();
order_fulfillment_saga.mark_failed().unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Failed);
assert!(order_fulfillment_saga.is_complete());
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let final_order_fulfillment_saga = saga_repo.get("saga-fail-456").await.unwrap().unwrap();
assert_eq!(final_order_fulfillment_saga.status(), SagaStatus::Failed);
assert_eq!(
final_order_fulfillment_saga.failure_reason(),
Some("Payment: Insufficient funds")
);
let final_order = order_repo.get(&order_id).await.unwrap().unwrap();
assert_eq!(final_order.status(), OrderStatus::Cancelled);
let final_inventory = inventory_repo.get("WIDGET-002").await.unwrap().unwrap();
assert_eq!(final_inventory.available(), 50); assert_eq!(final_inventory.reserved(), 0);
}
#[tokio::test]
async fn saga_compensates_on_inventory_failure() {
let order_repo = HashMapRepository::new().aggregate::<Order>();
let inventory_repo = HashMapRepository::new().aggregate::<Inventory>();
let saga_repo = HashMapRepository::new().aggregate::<OrderFulfillmentSaga>();
let mut widget_inventory = Inventory::new();
widget_inventory
.initialize("WIDGET-003".to_string(), 5)
.unwrap(); inventory_repo.commit(&mut widget_inventory).await.unwrap();
let order_id = "order-inv-fail-789".to_string();
let items = vec![OrderItem {
sku: "WIDGET-003".to_string(),
quantity: 10, price_cents: 500,
}];
let mut order = Order::new();
order
.create(order_id.clone(), "customer-xyz".to_string(), items.clone())
.unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut order_fulfillment_saga = OrderFulfillmentSaga::new();
order_fulfillment_saga
.start(
"saga-inv-fail".to_string(),
order_id.clone(),
"customer-xyz".to_string(),
items,
5000,
)
.unwrap();
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let inventory = inventory_repo.get("WIDGET-003").await.unwrap().unwrap();
assert!(!inventory.can_reserve(10));
let mut order_fulfillment_saga = saga_repo.get("saga-inv-fail").await.unwrap().unwrap();
order_fulfillment_saga
.step_failed(
"Inventory".to_string(),
"Insufficient stock for WIDGET-003".to_string(),
)
.unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Compensating);
assert!(!order_fulfillment_saga.needs_inventory_compensation());
assert!(!order_fulfillment_saga.needs_payment_compensation());
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let mut order = order_repo.get(&order_id).await.unwrap().unwrap();
order.cancel("Insufficient stock".to_string()).unwrap();
order_repo.commit(&mut order).await.unwrap();
let mut order_fulfillment_saga = saga_repo.get("saga-inv-fail").await.unwrap().unwrap();
order_fulfillment_saga.mark_failed().unwrap();
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let final_order_fulfillment_saga = saga_repo.get("saga-inv-fail").await.unwrap().unwrap();
assert_eq!(final_order_fulfillment_saga.status(), SagaStatus::Failed);
assert!(final_order_fulfillment_saga.is_complete());
let final_order = order_repo.get(&order_id).await.unwrap().unwrap();
assert_eq!(final_order.status(), OrderStatus::Cancelled);
let final_inventory = inventory_repo.get("WIDGET-003").await.unwrap().unwrap();
assert_eq!(final_inventory.available(), 5);
assert_eq!(final_inventory.reserved(), 0);
}
#[tokio::test]
async fn saga_is_replayable_from_events() {
let saga_repo = HashMapRepository::new().aggregate::<OrderFulfillmentSaga>();
let items = vec![OrderItem {
sku: "WIDGET-REPLAY".to_string(),
quantity: 3,
price_cents: 1500,
}];
let mut order_fulfillment_saga = OrderFulfillmentSaga::new();
order_fulfillment_saga
.start(
"saga-replay".to_string(),
"order-replay".to_string(),
"customer-replay".to_string(),
items,
4500,
)
.unwrap();
order_fulfillment_saga.inventory_reserved().unwrap();
order_fulfillment_saga.payment_succeeded().unwrap();
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let restored = saga_repo.get("saga-replay").await.unwrap().unwrap();
assert_eq!(restored.order_id(), "order-replay");
assert_eq!(restored.customer_id(), "customer-replay");
assert_eq!(restored.total_cents(), 4500);
assert_eq!(restored.status(), SagaStatus::PaymentProcessed);
assert!(restored.compensation().inventory_reserved);
assert!(restored.compensation().payment_processed);
let mut restored = restored;
restored.complete().unwrap();
saga_repo.commit(&mut restored).await.unwrap();
let final_order_fulfillment_saga = saga_repo.get("saga-replay").await.unwrap().unwrap();
assert_eq!(final_order_fulfillment_saga.status(), SagaStatus::Completed);
}
#[tokio::test]
async fn saga_tracks_compensation_state_correctly() {
let saga_repo = HashMapRepository::new().aggregate::<OrderFulfillmentSaga>();
let items = vec![OrderItem {
sku: "WIDGET-COMP".to_string(),
quantity: 1,
price_cents: 100,
}];
let mut order_fulfillment_saga = OrderFulfillmentSaga::new();
order_fulfillment_saga
.start(
"saga-comp".to_string(),
"order-comp".to_string(),
"customer-comp".to_string(),
items,
100,
)
.unwrap();
assert!(!order_fulfillment_saga.compensation().inventory_reserved);
assert!(!order_fulfillment_saga.compensation().payment_processed);
order_fulfillment_saga.inventory_reserved().unwrap();
assert!(order_fulfillment_saga.compensation().inventory_reserved);
assert!(!order_fulfillment_saga.compensation().payment_processed);
order_fulfillment_saga.payment_succeeded().unwrap();
assert!(order_fulfillment_saga.compensation().inventory_reserved);
assert!(order_fulfillment_saga.compensation().payment_processed);
order_fulfillment_saga
.step_failed("FinalStep".to_string(), "Something went wrong".to_string())
.unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Compensating);
assert!(order_fulfillment_saga.needs_inventory_compensation());
assert!(order_fulfillment_saga.needs_payment_compensation());
order_fulfillment_saga.payment_compensated().unwrap();
assert!(order_fulfillment_saga.needs_inventory_compensation());
assert!(!order_fulfillment_saga.needs_payment_compensation());
order_fulfillment_saga.inventory_compensated().unwrap();
assert!(!order_fulfillment_saga.needs_inventory_compensation());
assert!(!order_fulfillment_saga.needs_payment_compensation());
order_fulfillment_saga.mark_failed().unwrap();
assert_eq!(order_fulfillment_saga.status(), SagaStatus::Failed);
assert!(order_fulfillment_saga.is_complete());
saga_repo.commit(&mut order_fulfillment_saga).await.unwrap();
let restored = saga_repo.get("saga-comp").await.unwrap().unwrap();
assert_eq!(restored.status(), SagaStatus::Failed);
assert!(!restored.compensation().inventory_reserved);
assert!(!restored.compensation().payment_processed);
}