Add the saga module: Saga Pattern in Spring Boot, orchestration vs choreography with Kafka

This commit is contained in:
Claude
2026-10-03 19:56:40 +00:00
parent c1b40b0db9
commit 4a5d3318bf
44 changed files with 1296 additions and 3 deletions
+5 -3
View File
@@ -1,9 +1,10 @@
# spring-messaging-demo
Companion code for the messaging series on [ankurm.com](https://ankurm.com). Each directory is a
self-contained Maven project for one article, with its own `pom.xml`, its own numbered
documentation chapters, and its own captured output under `docs/output/` — regenerated by that
module's `scripts/run-all.sh`, never typed by hand.
self-contained Maven project for one article, with its own `pom.xml` and its own captured
output, regenerated from the test suite rather than typed by hand. Most modules keep that
output under `docs/output/` alongside numbered documentation chapters; `saga/` keeps it directly
under `output/` instead, with the equivalent depth in the article's own accordion sections.
| Module | Article | What it demonstrates |
|---|---|---|
@@ -13,6 +14,7 @@ module's `scripts/run-all.sh`, never typed by hand.
| [`sse-websocket/`](sse-websocket/README.md) | [Server-Sent Events and WebSocket on Spring Boot 4: SseEmitter, STOMP, and Which to Pick](https://ankurm.com/spring-boot-4-sse-websocket-stomp/) | The two browser-facing options side by side: the SSE lifecycle and the three ways a stream ends, STOMP's routing defaults, and the size limit that is enforced by Tomcat rather than by Spring |
| [`protocol-comparison/`](protocol-comparison/README.md) | [RSocket vs gRPC vs WebSocket on Spring Boot 4.1: When Each One Wins](https://ankurm.com/rsocket-vs-grpc-vs-websocket-spring-boot-4-1/) | One application serving the same two operations over all three, benchmarked in one JVM, ending in the only measurement that transfers off the box: what each server produces for a consumer that asked for a hundred |
| [`broker-comparison/`](broker-comparison/README.md) | [Kafka vs RabbitMQ vs Pulsar for Java Teams: A Decision Framework with Benchmarks](https://ankurm.com/kafka-vs-rabbitmq-vs-pulsar-java-decision-framework/) | The three brokers measured side by side on ordering, replay, consumer scaling and operational footprint, ending in a decision table where every row has a transcript behind it |
| [`saga/`](saga/README.md) | Saga Pattern in Spring Boot: Orchestration vs Choreography with Kafka | The same order/payment/inventory saga built twice against a real broker — choreography and orchestration — driven through the same inventory failure so the compensating transaction can be compared side by side, plus a real duplicate-delivery test for the idempotent-consumer problem |
The brokers make an instructive set. Kafka's consumer holds an offset and the broker remembers
nothing about individual records; RabbitMQ's broker owns the message until it is acknowledged and
+61
View File
@@ -0,0 +1,61 @@
# saga
Companion code for *Saga Pattern in Spring Boot: Orchestration vs Choreography with Kafka* on
[ankurm.com](https://ankurm.com).
The same order/payment/inventory saga, implemented twice against a real (embedded) Kafka
broker: once as **choreography** (every participant reacts to the previous step's event and
decides for itself what happens next) and once as **orchestration** (one class issues commands
and makes every decision; participants only ever answer a command). Both versions are driven
through the same failure -- insufficient inventory -- so the compensating transaction that
undoes the reserved payment can be compared side by side. A fourth test demonstrates the
idempotent-consumer problem concretely: the same command, redelivered, reserving stock only
once.
## Versions
| Artifact | Version |
|---|---|
| Spring Boot | 4.1.1 (GA; verified against Maven Central's `maven-metadata.xml` -- `4.2.0-M2` is the newest entry there but is a milestone, not GA) |
| Spring Kafka | 4.1.1, managed by the Spring Boot BOM |
| JDK | 25 (LTS) |
| H2 | runtime, in-memory, for the demo only |
## Quickstart
```
mvn -pl saga test -Dtest=ChoreographyHappyPathTest
mvn -pl saga test -Dtest=ChoreographyCompensationTest
mvn -pl saga test -Dtest=OrchestrationSagaTest
mvn -pl saga test -Dtest=IdempotentConsumerTest
```
Each test class starts its own embedded Kafka broker and its own isolated in-memory H2
database (a random schema name per test, via `@DynamicPropertySource`), so running them
individually or together (`mvn -pl saga test`) gives the same results.
## What each class is
| Class | Role |
|---|---|
| `choreography.ChoreographySagaStarter` | Creates the order, publishes `OrderCreated`, and does nothing else -- the only choreography class that "starts" anything |
| `choreography.PaymentChoreographyListener` | Reserves payment on `OrderCreated`; refunds it on `InventoryRejected` (the compensating transaction) |
| `choreography.InventoryChoreographyListener` | The saga's one failure trigger: reacts to `PaymentReserved`, decides stock is or isn't available, publishes accordingly |
| `choreography.OrderChoreographyListener` | Confirms or cancels the order; never decides anything, only reacts |
| `orchestration.OrderSagaOrchestrator` | Every saga decision lives here: what to do after each reply, including issuing the compensating `RefundPaymentCommand` |
| `orchestration.PaymentOrchestrationHandler` / `InventoryOrchestrationHandler` | Carry out exactly one command each and reply; never decide what happens next |
| `idempotency.IdempotencyGuard` | `claim(commandId)` -- true the first time a command id is seen, false on every redelivery, backed by a real UNIQUE constraint |
| `inventory.InventoryService` | In-memory stock only -- this module is about saga coordination, not inventory persistence |
## Captured output (`output/`)
| File | What it captures |
|---|---|
| `00-choreography-happy-path.txt` | Enough stock: payment reserved, inventory reserved, order confirmed, with no single class aware of the whole sequence |
| `01-choreography-compensation.txt` | Not enough stock: inventory rejects, payment refunds itself in response, order ends cancelled |
| `02-orchestration-saga.txt` | Both scenarios again, driven by `OrderSagaOrchestrator` instead -- same outcomes, one class making every decision |
| `03-idempotent-consumer-duplicate.txt` | The same `ReserveInventoryCommand`, same commandId, delivered twice: one processed-command row, one stock reservation, one reply |
No `docs/` chapter directory in this repository -- the intermediate and reference-depth
material that would normally live there is in accordion sections inside the WordPress post
itself.
@@ -0,0 +1,17 @@
$ mvn -pl saga test -Dtest=ChoreographyHappyPathTest
(real embedded Kafka broker via @EmbeddedKafka; ChoreographySagaStarter.placeOrder() publishes
OrderCreated, and every step after that is a listener reacting to the previous one's event --
no single class decides the whole saga)
Order 1 status: CONFIRMED
Stock remaining for GADGET-1: 8
Payment record status: RESERVED
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
Stock started at 10, the order asked for 2, and 8 is what's left -- the reservation really ran.
Five separate listeners fired in sequence to get here: PaymentChoreographyListener reacted to
OrderCreated and reserved payment; InventoryChoreographyListener reacted to PaymentReserved,
found enough stock, and published InventoryReserved; OrderChoreographyListener reacted to that
and confirmed the order. None of those three classes knows about the other two's existence --
each one only knows "when I see event X, I do Y, then publish Z".
@@ -0,0 +1,15 @@
$ mvn -pl saga test -Dtest=ChoreographyCompensationTest
(same saga, stock deliberately too low: 1 unit on hand, order asks for 5)
Order 1 status: CANCELLED
Stock remaining for GADGET-1 (untouched by the rejected reservation): 1
Payment record status: REFUNDED
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
InventoryChoreographyListener's tryReserve() refused before touching the stock map at all --
the 1 unit is exactly where it started. That rejection published InventoryRejected, which
PaymentChoreographyListener is also subscribed to; it refunded the payment it had reserved
earlier and published PaymentRefunded, which OrderChoreographyListener used to cancel the
order. The compensating transaction (the refund) is ordinary code in a class that already
existed for an unrelated reason -- nothing new coordinates it.
+20
View File
@@ -0,0 +1,20 @@
$ mvn -pl saga test -Dtest=OrchestrationSagaTest
(the same two scenarios -- enough stock, and not enough stock -- run through
OrderSagaOrchestrator instead of through independent listeners)
[happy path] Order 1 status: CONFIRMED
[happy path] Stock remaining for GADGET-ORCH-OK: 8
[happy path] Payment record status: RESERVED
[compensation] Order 2 status: CANCELLED
[compensation] Stock remaining for GADGET-ORCH-FAIL (untouched): 1
[compensation] Payment record status: REFUNDED
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0
Same outcomes as the choreography transcripts, reached a structurally different way. Here,
PaymentOrchestrationHandler and InventoryOrchestrationHandler never decide what happens next --
they carry out exactly one command each and reply. OrderSagaOrchestrator is the only class that
reads a reply and decides what command to send next, including the decision, in the second
test, to send RefundPaymentCommand once InventoryReply says no. Compare this transcript's two
tests with the two choreography transcripts above: identical business outcomes, from listener
classes that know nothing about the saga as a whole versus one class that knows all of it.
@@ -0,0 +1,16 @@
$ mvn -pl saga test -Dtest=IdempotentConsumerTest
(the exact same ReserveInventoryCommand -- same commandId -- published twice by hand, to
simulate a redelivery; stock starts at 10, the command asks for 3)
ProcessedCommand rows for commandId f07ec35c-853c-4405-a024-b62d304dfabb: 1
Stock remaining for GADGET-DUP after two deliveries of the same command: 7
InventoryReply messages actually published: 1
reply: InventoryReply[commandId=f07ec35c-853c-4405-a024-b62d304dfabb, orderId=999, success=true, reason=null]
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0
Two deliveries, one row in processed_commands, one actual reservation: 10 - 3 = 7, not 4. The
second delivery hit IdempotencyGuard.claim(), got false back because the UNIQUE constraint on
commandId rejected the second insert, and returned from the listener before calling
InventoryService.tryReserve() at all -- it never touched stock, and it never published a second
InventoryReply. Exactly one reply reached the topic for two deliveries of the same command.
+74
View File
@@ -0,0 +1,74 @@
<project xmlns="http://maven.apache.org/POM/4.0.0">
<modelVersion>4.0.0</modelVersion>
<!-- Same convention as the rest of this repo: inherit spring-boot-starter-parent so Kafka,
JPA and the test stack are all Boot-managed. Boot 4.1.1 manages Spring Kafka 4.1.1. -->
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>4.1.1</version>
<relativePath/>
</parent>
<groupId>com.ankurm</groupId>
<artifactId>saga</artifactId>
<version>1.0</version>
<packaging>jar</packaging>
<properties>
<java.version>25</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<!-- spring-boot-starter-kafka, not a bare spring-kafka dependency: see
kafka-error-handling/pom.xml in this same repo. Boot 4's auto-configuration lives in
spring-boot-kafka, which the starter brings and spring-kafka alone does not. -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jackson</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-jpa</artifactId>
</dependency>
<dependency>
<groupId>com.h2database</groupId>
<artifactId>h2</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.awaitility</groupId>
<artifactId>awaitility</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
@@ -0,0 +1,12 @@
package com.ankurm.sagademo;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class SagaDemoApplication {
public static void main(String[] args) {
SpringApplication.run(SagaDemoApplication.class, args);
}
}
@@ -0,0 +1,34 @@
package com.ankurm.sagademo;
/**
* Every topic this module uses, in one place, so the two coordination styles in the post --
* choreography and orchestration -- never accidentally share a topic. The "choreo." group
* carries plain domain events; the "orch." group splits into commands (the orchestrator telling
* a participant what to do) and replies (the participant reporting back), which is the concrete
* difference between the two styles: choreography has only events, orchestration has commands
* and replies to them.
*/
public final class Topics {
private Topics() {
}
// Choreography: each participant listens for the previous step's event and decides for
// itself what to publish next. No single place knows the whole saga.
public static final String CHOREO_ORDER_CREATED = "choreo.order.created";
public static final String CHOREO_PAYMENT_RESERVED = "choreo.payment.reserved";
public static final String CHOREO_INVENTORY_RESERVED = "choreo.inventory.reserved";
public static final String CHOREO_INVENTORY_REJECTED = "choreo.inventory.rejected";
public static final String CHOREO_PAYMENT_REFUNDED = "choreo.payment.refunded";
public static final String CHOREO_ORDER_CONFIRMED = "choreo.order.confirmed";
public static final String CHOREO_ORDER_CANCELLED = "choreo.order.cancelled";
// Orchestration: the orchestrator issues commands and participants reply; only the
// orchestrator decides what happens next.
public static final String ORCH_CMD_RESERVE_PAYMENT = "orch.cmd.reserve-payment";
public static final String ORCH_REPLY_PAYMENT = "orch.reply.payment";
public static final String ORCH_CMD_RESERVE_INVENTORY = "orch.cmd.reserve-inventory";
public static final String ORCH_REPLY_INVENTORY = "orch.reply.inventory";
public static final String ORCH_CMD_REFUND_PAYMENT = "orch.cmd.refund-payment";
public static final String ORCH_REPLY_REFUND = "orch.reply.refund";
}
@@ -0,0 +1,30 @@
package com.ankurm.sagademo.choreography;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.OrderCreated;
import com.ankurm.sagademo.orders.Order;
import com.ankurm.sagademo.orders.OrderService;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
/** The only thing that starts a choreography saga: create the order, publish one event, and
* walk away. Everything after this call is other services reacting to events, not this class
* coordinating anything. */
@Component
public class ChoreographySagaStarter {
private final OrderService orders;
private final KafkaTemplate<String, Object> kafka;
public ChoreographySagaStarter(OrderService orders, KafkaTemplate<String, Object> kafka) {
this.orders = orders;
this.kafka = kafka;
}
public Order placeOrder(String sku, int qty, int amountCents) {
Order order = orders.createPendingOrder(sku, qty, amountCents, "CHOREOGRAPHY");
kafka.send(Topics.CHOREO_ORDER_CREATED, new OrderCreated(order.getId(), sku, qty, amountCents));
return order;
}
}
@@ -0,0 +1,40 @@
package com.ankurm.sagademo.choreography;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.InventoryRejected;
import com.ankurm.sagademo.events.InventoryReserved;
import com.ankurm.sagademo.events.PaymentReserved;
import com.ankurm.sagademo.inventory.InventoryService;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
/**
* This is the only place in the choreography saga that knows what "not enough stock" means --
* it reacts to a payment event, makes a decision with no other service's input, and publishes
* one of two outcomes. Whichever one it publishes, it has no idea what will happen next; it
* just knows it told the truth about its own stock.
*/
@Component
public class InventoryChoreographyListener {
private final InventoryService inventory;
private final KafkaTemplate<String, Object> kafka;
public InventoryChoreographyListener(InventoryService inventory, KafkaTemplate<String, Object> kafka) {
this.inventory = inventory;
this.kafka = kafka;
}
@KafkaListener(topics = Topics.CHOREO_PAYMENT_RESERVED, groupId = "choreo-inventory")
public void onPaymentReserved(PaymentReserved event) {
boolean reserved = inventory.tryReserve(event.sku(), event.qty());
if (reserved) {
kafka.send(Topics.CHOREO_INVENTORY_RESERVED, new InventoryReserved(event.orderId()));
} else {
kafka.send(Topics.CHOREO_INVENTORY_REJECTED,
new InventoryRejected(event.orderId(), "insufficient stock for " + event.sku()));
}
}
}
@@ -0,0 +1,34 @@
package com.ankurm.sagademo.choreography;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.InventoryReserved;
import com.ankurm.sagademo.events.PaymentRefunded;
import com.ankurm.sagademo.orders.OrderService;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
/**
* The order itself has no say in any of this -- it is confirmed or cancelled by whichever event
* happens to arrive. That passivity is the point: in choreography, the order record is just
* data other services' events update, not a thing that drives the saga.
*/
@Component
public class OrderChoreographyListener {
private final OrderService orders;
public OrderChoreographyListener(OrderService orders) {
this.orders = orders;
}
@KafkaListener(topics = Topics.CHOREO_INVENTORY_RESERVED, groupId = "choreo-order-confirm")
public void onInventoryReserved(InventoryReserved event) {
orders.confirm(event.orderId());
}
@KafkaListener(topics = Topics.CHOREO_PAYMENT_REFUNDED, groupId = "choreo-order-cancel")
public void onPaymentRefunded(PaymentRefunded event) {
orders.cancel(event.orderId());
}
}
@@ -0,0 +1,47 @@
package com.ankurm.sagademo.choreography;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.InventoryRejected;
import com.ankurm.sagademo.events.OrderCreated;
import com.ankurm.sagademo.events.PaymentRefunded;
import com.ankurm.sagademo.events.PaymentReserved;
import com.ankurm.sagademo.payments.PaymentService;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
/**
* Payment's two jobs in the choreography saga, and nothing else. This class never hears about
* the order being confirmed or cancelled -- it reacts to its own two events and moves on. Who
* decides what happens to the order as a whole is not answered anywhere in this class, which is
* exactly the property choreography trades for having no single orchestrator to maintain.
*/
@Component
public class PaymentChoreographyListener {
private final PaymentService payments;
private final KafkaTemplate<String, Object> kafka;
public PaymentChoreographyListener(PaymentService payments, KafkaTemplate<String, Object> kafka) {
this.payments = payments;
this.kafka = kafka;
}
@KafkaListener(topics = Topics.CHOREO_ORDER_CREATED, groupId = "choreo-payment")
public void onOrderCreated(OrderCreated event) {
payments.reserve(event.orderId(), event.amountCents());
kafka.send(Topics.CHOREO_PAYMENT_RESERVED,
new PaymentReserved(event.orderId(), event.sku(), event.qty()));
}
/** The compensating transaction. Nothing calls this directly -- it runs because this
* listener happens to also be subscribed to the rejection event, which is choreography's
* defining trait: the participant that caused the original effect is the one that notices
* the failure and undoes it, not a coordinator telling it to. */
@KafkaListener(topics = Topics.CHOREO_INVENTORY_REJECTED, groupId = "choreo-payment-compensation")
public void onInventoryRejected(InventoryRejected event) {
payments.refund(event.orderId());
kafka.send(Topics.CHOREO_PAYMENT_REFUNDED, new PaymentRefunded(event.orderId()));
}
}
@@ -0,0 +1,6 @@
package com.ankurm.sagademo.events;
/** The saga's one failure trigger in this demo, deliberately placed in inventory so both
* coordination styles can be compared against exactly the same failure. */
public record InventoryRejected(Long orderId, String reason) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record InventoryReply(String commandId, Long orderId, boolean success, String reason) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record InventoryReserved(Long orderId) {
}
@@ -0,0 +1,7 @@
package com.ankurm.sagademo.events;
/** Choreography's starting event. Nothing downstream of this was told what to do -- each
* listener decides for itself, from its own domain knowledge, what reacting to this event
* means. */
public record OrderCreated(Long orderId, String sku, int qty, int amountCents) {
}
@@ -0,0 +1,6 @@
package com.ankurm.sagademo.events;
/** The compensating transaction's own completion event -- published after undoing the
* payment reservation that an earlier, now-abandoned step in the saga made. */
public record PaymentRefunded(Long orderId) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record PaymentReply(String commandId, Long orderId, boolean success) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record PaymentReserved(Long orderId, String sku, int qty) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record RefundPaymentCommand(String commandId, Long orderId) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record RefundReply(String commandId, Long orderId) {
}
@@ -0,0 +1,4 @@
package com.ankurm.sagademo.events;
public record ReserveInventoryCommand(String commandId, Long orderId, String sku, int qty) {
}
@@ -0,0 +1,8 @@
package com.ankurm.sagademo.events;
/** commandId is the idempotency key -- see IdempotentCommand in the idempotency package. The
* orchestrator sets it once per logical step and reuses the exact same value if it ever
* resends the command, which is what lets a participant tell "the same command, redelivered"
* apart from "a brand new command that happens to target the same order". */
public record ReservePaymentCommand(String commandId, Long orderId, int amountCents) {
}
@@ -0,0 +1,34 @@
package com.ankurm.sagademo.idempotency;
import org.springframework.dao.DataIntegrityViolationException;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;
/**
* "Has this exact command already been carried out?" -- answered by trying to insert its id
* and seeing whether the database's own UNIQUE constraint lets it through. REQUIRES_NEW so a
* duplicate here rolls back only this one insert, never whatever business transaction the
* caller is also in the middle of.
*/
@Component
public class IdempotencyGuard {
private final ProcessedCommandRepository processed;
public IdempotencyGuard(ProcessedCommandRepository processed) {
this.processed = processed;
}
/** Returns true the first time a given commandId is claimed, false on every claim after
* that. Callers are expected to do their real work only when this returns true. */
@Transactional(propagation = Propagation.REQUIRES_NEW)
public boolean claim(String commandId) {
try {
processed.save(new ProcessedCommand(commandId));
return true;
} catch (DataIntegrityViolationException alreadyClaimed) {
return false;
}
}
}
@@ -0,0 +1,39 @@
package com.ankurm.sagademo.idempotency;
import jakarta.persistence.Column;
import jakarta.persistence.Entity;
import jakarta.persistence.GeneratedValue;
import jakarta.persistence.GenerationType;
import jakarta.persistence.Id;
import jakarta.persistence.Table;
import jakarta.persistence.UniqueConstraint;
/** One row per command id this consumer has actually carried out. The UNIQUE constraint on
* commandId, not application code, is what makes the dedupe check safe under concurrent
* delivery -- two threads racing to insert the same commandId can't both win. */
@Entity
@Table(name = "processed_commands", uniqueConstraints = @UniqueConstraint(columnNames = "commandId"))
public class ProcessedCommand {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(nullable = false)
private String commandId;
protected ProcessedCommand() {
}
public ProcessedCommand(String commandId) {
this.commandId = commandId;
}
public Long getId() {
return id;
}
public String getCommandId() {
return commandId;
}
}
@@ -0,0 +1,8 @@
package com.ankurm.sagademo.idempotency;
import org.springframework.data.jpa.repository.JpaRepository;
public interface ProcessedCommandRepository extends JpaRepository<ProcessedCommand, Long> {
long countByCommandId(String commandId);
}
@@ -0,0 +1,37 @@
package com.ankurm.sagademo.inventory;
import java.util.concurrent.ConcurrentHashMap;
import org.springframework.stereotype.Service;
/**
* Deliberately not a database table. Stock here is in-memory, reset per-test by calling
* {@link #setStock}, because the point of this module is saga coordination, not inventory
* persistence. The dual-write problem and the transactional outbox that actually fixes it for a
* real stock table are the subject of a different post in this series.
*/
@Service
public class InventoryService {
private final ConcurrentHashMap<String, Integer> stock = new ConcurrentHashMap<>();
public void setStock(String sku, int quantity) {
stock.put(sku, quantity);
}
public int getStock(String sku) {
return stock.getOrDefault(sku, 0);
}
/** Returns true and decrements if enough stock was available; returns false and leaves
* stock untouched otherwise. The check and the decrement happen under one lock per SKU so
* two concurrent reservations for the last unit can't both succeed. */
public synchronized boolean tryReserve(String sku, int quantity) {
int available = stock.getOrDefault(sku, 0);
if (available < quantity) {
return false;
}
stock.put(sku, available - quantity);
return true;
}
}
@@ -0,0 +1,45 @@
package com.ankurm.sagademo.orchestration;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.InventoryReply;
import com.ankurm.sagademo.events.ReserveInventoryCommand;
import com.ankurm.sagademo.idempotency.IdempotencyGuard;
import com.ankurm.sagademo.inventory.InventoryService;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
/**
* The one listener in this module where a duplicate delivery would be genuinely dangerous: a
* command, not a notification, and the action it triggers -- decrementing stock -- is not safe
* to run twice. {@link IdempotencyGuard#claim} is checked before touching inventory at all;
* a redelivered command with the same commandId is dropped, silently, with no second reply.
*/
@Component
public class InventoryOrchestrationHandler {
private final InventoryService inventory;
private final IdempotencyGuard idempotencyGuard;
private final KafkaTemplate<String, Object> kafka;
public InventoryOrchestrationHandler(InventoryService inventory, IdempotencyGuard idempotencyGuard,
KafkaTemplate<String, Object> kafka) {
this.inventory = inventory;
this.idempotencyGuard = idempotencyGuard;
this.kafka = kafka;
}
@KafkaListener(topics = Topics.ORCH_CMD_RESERVE_INVENTORY, groupId = "orch-inventory")
public void onReserveInventory(ReserveInventoryCommand command) {
if (!idempotencyGuard.claim(command.commandId())) {
// Already carried out once. The first delivery already sent the one reply this
// command will ever get; stock must not move a second time.
return;
}
boolean reserved = inventory.tryReserve(command.sku(), command.qty());
String reason = reserved ? null : "insufficient stock for " + command.sku();
kafka.send(Topics.ORCH_REPLY_INVENTORY,
new InventoryReply(command.commandId(), command.orderId(), reserved, reason));
}
}
@@ -0,0 +1,76 @@
package com.ankurm.sagademo.orchestration;
import java.util.UUID;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.InventoryReply;
import com.ankurm.sagademo.events.PaymentReply;
import com.ankurm.sagademo.events.RefundPaymentCommand;
import com.ankurm.sagademo.events.RefundReply;
import com.ankurm.sagademo.events.ReserveInventoryCommand;
import com.ankurm.sagademo.events.ReservePaymentCommand;
import com.ankurm.sagademo.orders.Order;
import com.ankurm.sagademo.orders.OrderRepository;
import com.ankurm.sagademo.orders.OrderService;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
/**
* Every decision in this saga is made here and nowhere else. Payment and inventory never see
* each other's events, never decide what happens next, and never know whether the order was
* ultimately confirmed or cancelled -- they only ever answer the one command this class sent
* them. Contrast with the choreography package, where no single class looks like this one.
*/
@Component
public class OrderSagaOrchestrator {
private final OrderService orders;
private final OrderRepository orderRepository;
private final KafkaTemplate<String, Object> kafka;
public OrderSagaOrchestrator(OrderService orders, OrderRepository orderRepository,
KafkaTemplate<String, Object> kafka) {
this.orders = orders;
this.orderRepository = orderRepository;
this.kafka = kafka;
}
public Order startSaga(String sku, int qty, int amountCents) {
Order order = orders.createPendingOrder(sku, qty, amountCents, "ORCHESTRATION");
kafka.send(Topics.ORCH_CMD_RESERVE_PAYMENT,
new ReservePaymentCommand(UUID.randomUUID().toString(), order.getId(), amountCents));
return order;
}
@KafkaListener(topics = Topics.ORCH_REPLY_PAYMENT, groupId = "orchestrator")
public void onPaymentReply(PaymentReply reply) {
if (!reply.success()) {
orders.cancel(reply.orderId());
return;
}
Order order = orderRepository.findById(reply.orderId()).orElseThrow();
kafka.send(Topics.ORCH_CMD_RESERVE_INVENTORY,
new ReserveInventoryCommand(UUID.randomUUID().toString(), order.getId(), order.getSku(),
order.getQty()));
}
@KafkaListener(topics = Topics.ORCH_REPLY_INVENTORY, groupId = "orchestrator")
public void onInventoryReply(InventoryReply reply) {
if (reply.success()) {
orders.confirm(reply.orderId());
return;
}
// Compensation: the orchestrator is the one that notices the failure and decides to
// undo payment -- payment itself never finds out inventory said no except by being
// told to refund.
kafka.send(Topics.ORCH_CMD_REFUND_PAYMENT,
new RefundPaymentCommand(UUID.randomUUID().toString(), reply.orderId()));
}
@KafkaListener(topics = Topics.ORCH_REPLY_REFUND, groupId = "orchestrator")
public void onRefundReply(RefundReply reply) {
orders.cancel(reply.orderId());
}
}
@@ -0,0 +1,38 @@
package com.ankurm.sagademo.orchestration;
import com.ankurm.sagademo.Topics;
import com.ankurm.sagademo.events.PaymentReply;
import com.ankurm.sagademo.events.RefundPaymentCommand;
import com.ankurm.sagademo.events.RefundReply;
import com.ankurm.sagademo.events.ReservePaymentCommand;
import com.ankurm.sagademo.payments.PaymentService;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
/** Payment under orchestration does only what it's told: carry out the command, reply with the
* outcome. It never decides what the reply should lead to. */
@Component
public class PaymentOrchestrationHandler {
private final PaymentService payments;
private final KafkaTemplate<String, Object> kafka;
public PaymentOrchestrationHandler(PaymentService payments, KafkaTemplate<String, Object> kafka) {
this.payments = payments;
this.kafka = kafka;
}
@KafkaListener(topics = Topics.ORCH_CMD_RESERVE_PAYMENT, groupId = "orch-payment")
public void onReservePayment(ReservePaymentCommand command) {
payments.reserve(command.orderId(), command.amountCents());
kafka.send(Topics.ORCH_REPLY_PAYMENT, new PaymentReply(command.commandId(), command.orderId(), true));
}
@KafkaListener(topics = Topics.ORCH_CMD_REFUND_PAYMENT, groupId = "orch-payment-refund")
public void onRefundPayment(RefundPaymentCommand command) {
payments.refund(command.orderId());
kafka.send(Topics.ORCH_REPLY_REFUND, new RefundReply(command.commandId(), command.orderId()));
}
}
@@ -0,0 +1,66 @@
package com.ankurm.sagademo.orders;
import jakarta.persistence.Entity;
import jakarta.persistence.GeneratedValue;
import jakarta.persistence.GenerationType;
import jakarta.persistence.Id;
import jakarta.persistence.Table;
@Entity
@Table(name = "orders")
public class Order {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
private String sku;
private int qty;
private int amountCents;
private OrderStatus status;
/** Which coordination style this particular order's saga is running under -- "CHOREOGRAPHY"
* or "ORCHESTRATION". Both styles share this same Order table and the same InventoryService
* stock, which is exactly why each test resets its own sku's stock rather than relying on
* the other style to have left it alone. */
private String sagaMode;
protected Order() {
}
public Order(String sku, int qty, int amountCents, String sagaMode) {
this.sku = sku;
this.qty = qty;
this.amountCents = amountCents;
this.sagaMode = sagaMode;
this.status = OrderStatus.PENDING;
}
public Long getId() {
return id;
}
public String getSku() {
return sku;
}
public int getQty() {
return qty;
}
public int getAmountCents() {
return amountCents;
}
public OrderStatus getStatus() {
return status;
}
public void setStatus(OrderStatus status) {
this.status = status;
}
public String getSagaMode() {
return sagaMode;
}
}
@@ -0,0 +1,6 @@
package com.ankurm.sagademo.orders;
import org.springframework.data.jpa.repository.JpaRepository;
public interface OrderRepository extends JpaRepository<Order, Long> {
}
@@ -0,0 +1,31 @@
package com.ankurm.sagademo.orders;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@Service
public class OrderService {
private final OrderRepository orders;
public OrderService(OrderRepository orders) {
this.orders = orders;
}
@Transactional
public Order createPendingOrder(String sku, int qty, int amountCents, String sagaMode) {
return orders.save(new Order(sku, qty, amountCents, sagaMode));
}
@Transactional
public void confirm(Long orderId) {
Order order = orders.findById(orderId).orElseThrow();
order.setStatus(OrderStatus.CONFIRMED);
}
@Transactional
public void cancel(Long orderId) {
Order order = orders.findById(orderId).orElseThrow();
order.setStatus(OrderStatus.CANCELLED);
}
}
@@ -0,0 +1,7 @@
package com.ankurm.sagademo.orders;
public enum OrderStatus {
PENDING,
CONFIRMED,
CANCELLED
}
@@ -0,0 +1,49 @@
package com.ankurm.sagademo.payments;
import jakarta.persistence.Entity;
import jakarta.persistence.GeneratedValue;
import jakarta.persistence.GenerationType;
import jakarta.persistence.Id;
import jakarta.persistence.Table;
@Entity
@Table(name = "payment_records")
public class PaymentRecord {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
private Long orderId;
private int amountCents;
private PaymentStatus status;
protected PaymentRecord() {
}
public PaymentRecord(Long orderId, int amountCents) {
this.orderId = orderId;
this.amountCents = amountCents;
this.status = PaymentStatus.RESERVED;
}
public Long getId() {
return id;
}
public Long getOrderId() {
return orderId;
}
public int getAmountCents() {
return amountCents;
}
public PaymentStatus getStatus() {
return status;
}
public void setStatus(PaymentStatus status) {
this.status = status;
}
}
@@ -0,0 +1,9 @@
package com.ankurm.sagademo.payments;
import java.util.Optional;
import org.springframework.data.jpa.repository.JpaRepository;
public interface PaymentRepository extends JpaRepository<PaymentRecord, Long> {
Optional<PaymentRecord> findByOrderId(Long orderId);
}
@@ -0,0 +1,31 @@
package com.ankurm.sagademo.payments;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
/**
* Payment always succeeds in this demo -- the saga's single point of failure is deliberately
* placed in inventory instead, so both coordination styles can be compared against exactly the
* same failure, with payment's compensation (a refund) as the thing that actually gets exercised
* when inventory says no.
*/
@Service
public class PaymentService {
private final PaymentRepository payments;
public PaymentService(PaymentRepository payments) {
this.payments = payments;
}
@Transactional
public PaymentRecord reserve(Long orderId, int amountCents) {
return payments.save(new PaymentRecord(orderId, amountCents));
}
@Transactional
public void refund(Long orderId) {
PaymentRecord payment = payments.findByOrderId(orderId).orElseThrow();
payment.setStatus(PaymentStatus.REFUNDED);
}
}
@@ -0,0 +1,6 @@
package com.ankurm.sagademo.payments;
public enum PaymentStatus {
RESERVED,
REFUNDED
}
@@ -0,0 +1,12 @@
spring.application.name=saga
spring.datasource.url=jdbc:h2:mem:sagadb;DB_CLOSE_DELAY=-1
spring.jpa.hibernate.ddl-auto=update
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties.spring.json.trusted.packages=com.ankurm.sagademo.*
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
@@ -0,0 +1,76 @@
package com.ankurm.sagademo;
import java.time.Duration;
import java.util.UUID;
import com.ankurm.sagademo.choreography.ChoreographySagaStarter;
import com.ankurm.sagademo.inventory.InventoryService;
import com.ankurm.sagademo.orders.Order;
import com.ankurm.sagademo.orders.OrderRepository;
import com.ankurm.sagademo.orders.OrderStatus;
import com.ankurm.sagademo.payments.PaymentRepository;
import com.ankurm.sagademo.payments.PaymentStatus;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
/**
* Same saga, one fact changed: not enough stock. Nobody told payment to expect a refund --
* PaymentChoreographyListener is simply also subscribed to the rejection event, and reacts to
* it the same way it reacts to anything else. That is choreography's compensation story in
* full: a participant's own subscription list is the only place the "undo" logic lives.
*/
@SpringBootTest
@EmbeddedKafka(partitions = 1, bootstrapServersProperty = "spring.kafka.bootstrap-servers", topics = {
Topics.CHOREO_ORDER_CREATED, Topics.CHOREO_PAYMENT_RESERVED, Topics.CHOREO_INVENTORY_RESERVED,
Topics.CHOREO_INVENTORY_REJECTED, Topics.CHOREO_PAYMENT_REFUNDED
})
class ChoreographyCompensationTest {
@Autowired
ChoreographySagaStarter sagaStarter;
@Autowired
InventoryService inventory;
@Autowired
OrderRepository orders;
@Autowired
PaymentRepository payments;
@DynamicPropertySource
static void kafkaProperties(DynamicPropertyRegistry registry) {
registry.add("spring.datasource.url",
() -> "jdbc:h2:mem:" + UUID.randomUUID() + ";DB_CLOSE_DELAY=-1");
}
@Test
void insufficientStockTriggersPaymentRefundAndOrderCancellation() {
inventory.setStock("GADGET-1", 1);
Order order = sagaStarter.placeOrder("GADGET-1", 5, 2999);
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
Order reloaded = orders.findById(order.getId()).orElseThrow();
assertThat(reloaded.getStatus()).isEqualTo(OrderStatus.CANCELLED);
});
System.out.println("Order " + order.getId() + " status: "
+ orders.findById(order.getId()).orElseThrow().getStatus());
System.out.println("Stock remaining for GADGET-1 (untouched by the rejected reservation): "
+ inventory.getStock("GADGET-1"));
System.out.println("Payment record status: "
+ payments.findByOrderId(order.getId()).orElseThrow().getStatus());
// The stock was never decremented -- tryReserve's own check-then-decrement refused
// before touching the map.
assertThat(inventory.getStock("GADGET-1")).isEqualTo(1);
assertThat(payments.findByOrderId(order.getId()).orElseThrow().getStatus())
.isEqualTo(PaymentStatus.REFUNDED);
}
}
@@ -0,0 +1,67 @@
package com.ankurm.sagademo;
import java.time.Duration;
import java.util.UUID;
import com.ankurm.sagademo.choreography.ChoreographySagaStarter;
import com.ankurm.sagademo.inventory.InventoryService;
import com.ankurm.sagademo.orders.Order;
import com.ankurm.sagademo.orders.OrderRepository;
import com.ankurm.sagademo.orders.OrderStatus;
import com.ankurm.sagademo.payments.PaymentRepository;
import com.ankurm.sagademo.payments.PaymentStatus;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
@SpringBootTest
@EmbeddedKafka(partitions = 1, bootstrapServersProperty = "spring.kafka.bootstrap-servers", topics = {
Topics.CHOREO_ORDER_CREATED, Topics.CHOREO_PAYMENT_RESERVED, Topics.CHOREO_INVENTORY_RESERVED,
Topics.CHOREO_INVENTORY_REJECTED, Topics.CHOREO_PAYMENT_REFUNDED
})
class ChoreographyHappyPathTest {
@Autowired
ChoreographySagaStarter sagaStarter;
@Autowired
InventoryService inventory;
@Autowired
OrderRepository orders;
@Autowired
PaymentRepository payments;
@DynamicPropertySource
static void kafkaProperties(DynamicPropertyRegistry registry) {
registry.add("spring.datasource.url",
() -> "jdbc:h2:mem:" + UUID.randomUUID() + ";DB_CLOSE_DELAY=-1");
}
@Test
void orderGoesThroughPaymentAndInventoryAndEndsConfirmed() {
inventory.setStock("GADGET-1", 10);
Order order = sagaStarter.placeOrder("GADGET-1", 2, 2999);
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
Order reloaded = orders.findById(order.getId()).orElseThrow();
assertThat(reloaded.getStatus()).isEqualTo(OrderStatus.CONFIRMED);
});
System.out.println("Order " + order.getId() + " status: "
+ orders.findById(order.getId()).orElseThrow().getStatus());
System.out.println("Stock remaining for GADGET-1: " + inventory.getStock("GADGET-1"));
System.out.println("Payment record status: "
+ payments.findByOrderId(order.getId()).orElseThrow().getStatus());
assertThat(inventory.getStock("GADGET-1")).isEqualTo(8);
assertThat(payments.findByOrderId(order.getId()).orElseThrow().getStatus())
.isEqualTo(PaymentStatus.RESERVED);
}
}
@@ -0,0 +1,105 @@
package com.ankurm.sagademo;
import java.time.Duration;
import java.util.UUID;
import com.ankurm.sagademo.events.InventoryReply;
import com.ankurm.sagademo.events.ReserveInventoryCommand;
import com.ankurm.sagademo.idempotency.ProcessedCommandRepository;
import com.ankurm.sagademo.inventory.InventoryService;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.kafka.test.utils.KafkaTestUtils;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
/**
* Publishes the exact same {@link ReserveInventoryCommand} -- same commandId -- twice, by
* hand, to simulate exactly the kind of redelivery a consumer rebalance or a producer retry
* can cause. This is the one listener in the saga where a duplicate is dangerous: it decrements
* real stock, not just a log line.
*/
@SpringBootTest
@EmbeddedKafka(partitions = 1, bootstrapServersProperty = "spring.kafka.bootstrap-servers", topics = {
Topics.ORCH_CMD_RESERVE_INVENTORY, Topics.ORCH_REPLY_INVENTORY
})
class IdempotentConsumerTest {
@Autowired
KafkaTemplate<String, Object> kafka;
@Autowired
InventoryService inventory;
@Autowired
ProcessedCommandRepository processedCommands;
@Autowired
org.springframework.kafka.test.EmbeddedKafkaBroker embeddedKafka;
Consumer<String, Object> replyConsumer;
@DynamicPropertySource
static void kafkaProperties(DynamicPropertyRegistry registry) {
registry.add("spring.datasource.url",
() -> "jdbc:h2:mem:" + UUID.randomUUID() + ";DB_CLOSE_DELAY=-1");
}
@AfterEach
void closeConsumer() {
if (replyConsumer != null) {
replyConsumer.close();
}
}
@Test
void sameCommandIdDeliveredTwiceOnlyReservesStockOnce() {
inventory.setStock("GADGET-DUP", 10);
var consumerProps = KafkaTestUtils.consumerProps("idempotency-test-group", "true",
embeddedKafka);
consumerProps.put("value.deserializer", "org.springframework.kafka.support.serializer.JsonDeserializer");
consumerProps.put("spring.json.trusted.packages", "com.ankurm.sagademo.*");
consumerProps.put("auto.offset.reset", "earliest");
var cf = new org.springframework.kafka.core.DefaultKafkaConsumerFactory<String, Object>(consumerProps);
replyConsumer = cf.createConsumer();
embeddedKafka.consumeFromAnEmbeddedTopic(replyConsumer, Topics.ORCH_REPLY_INVENTORY);
String duplicatedCommandId = UUID.randomUUID().toString();
ReserveInventoryCommand command = new ReserveInventoryCommand(duplicatedCommandId, 999L, "GADGET-DUP", 3);
// The exact same command, same commandId, sent twice -- this is the redelivery.
kafka.send(Topics.ORCH_CMD_RESERVE_INVENTORY, command);
kafka.send(Topics.ORCH_CMD_RESERVE_INVENTORY, command);
await().atMost(Duration.ofSeconds(10)).untilAsserted(() ->
assertThat(processedCommands.countByCommandId(duplicatedCommandId)).isEqualTo(1));
// Give any second (wrongly-processed) reply a moment to arrive, if it was going to.
ConsumerRecords<String, Object> records = KafkaTestUtils.getRecords(replyConsumer, Duration.ofSeconds(5));
int replyCount = records.count();
System.out.println("ProcessedCommand rows for commandId " + duplicatedCommandId + ": "
+ processedCommands.countByCommandId(duplicatedCommandId));
System.out.println("Stock remaining for GADGET-DUP after two deliveries of the same command: "
+ inventory.getStock("GADGET-DUP"));
System.out.println("InventoryReply messages actually published: " + replyCount);
for (ConsumerRecord<String, Object> record : records) {
System.out.println(" reply: " + record.value());
}
// 10 - 3 = 7 if the reservation ran exactly once. If the duplicate had also been
// processed, this would be 4.
assertThat(inventory.getStock("GADGET-DUP")).isEqualTo(7);
assertThat(replyCount).isEqualTo(1);
}
}
@@ -0,0 +1,94 @@
package com.ankurm.sagademo;
import java.time.Duration;
import java.util.UUID;
import com.ankurm.sagademo.inventory.InventoryService;
import com.ankurm.sagademo.orchestration.OrderSagaOrchestrator;
import com.ankurm.sagademo.orders.Order;
import com.ankurm.sagademo.orders.OrderRepository;
import com.ankurm.sagademo.orders.OrderStatus;
import com.ankurm.sagademo.payments.PaymentRepository;
import com.ankurm.sagademo.payments.PaymentStatus;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
/**
* The same two scenarios as the choreography tests, run through {@link OrderSagaOrchestrator}
* instead. Every decision -- what to do after payment succeeds, what to do after inventory
* fails -- is made inside that one class; payment and inventory only ever answer a command.
*/
@SpringBootTest
@EmbeddedKafka(partitions = 1, bootstrapServersProperty = "spring.kafka.bootstrap-servers", topics = {
Topics.ORCH_CMD_RESERVE_PAYMENT, Topics.ORCH_REPLY_PAYMENT, Topics.ORCH_CMD_RESERVE_INVENTORY,
Topics.ORCH_REPLY_INVENTORY, Topics.ORCH_CMD_REFUND_PAYMENT, Topics.ORCH_REPLY_REFUND
})
class OrchestrationSagaTest {
@Autowired
OrderSagaOrchestrator orchestrator;
@Autowired
InventoryService inventory;
@Autowired
OrderRepository orders;
@Autowired
PaymentRepository payments;
@DynamicPropertySource
static void kafkaProperties(DynamicPropertyRegistry registry) {
registry.add("spring.datasource.url",
() -> "jdbc:h2:mem:" + UUID.randomUUID() + ";DB_CLOSE_DELAY=-1");
}
@Test
void ordersConfirmedByTheOrchestratorAfterBothStepsSucceed() {
inventory.setStock("GADGET-ORCH-OK", 10);
Order order = orchestrator.startSaga("GADGET-ORCH-OK", 2, 2999);
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
Order reloaded = orders.findById(order.getId()).orElseThrow();
assertThat(reloaded.getStatus()).isEqualTo(OrderStatus.CONFIRMED);
});
System.out.println("[happy path] Order " + order.getId() + " status: "
+ orders.findById(order.getId()).orElseThrow().getStatus());
System.out.println("[happy path] Stock remaining for GADGET-ORCH-OK: "
+ inventory.getStock("GADGET-ORCH-OK"));
System.out.println("[happy path] Payment record status: "
+ payments.findByOrderId(order.getId()).orElseThrow().getStatus());
assertThat(inventory.getStock("GADGET-ORCH-OK")).isEqualTo(8);
}
@Test
void insufficientStockMakesTheOrchestratorIssueARefundCommandAndCancelTheOrder() {
inventory.setStock("GADGET-ORCH-FAIL", 1);
Order order = orchestrator.startSaga("GADGET-ORCH-FAIL", 5, 2999);
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
Order reloaded = orders.findById(order.getId()).orElseThrow();
assertThat(reloaded.getStatus()).isEqualTo(OrderStatus.CANCELLED);
});
System.out.println("[compensation] Order " + order.getId() + " status: "
+ orders.findById(order.getId()).orElseThrow().getStatus());
System.out.println("[compensation] Stock remaining for GADGET-ORCH-FAIL (untouched): "
+ inventory.getStock("GADGET-ORCH-FAIL"));
System.out.println("[compensation] Payment record status: "
+ payments.findByOrderId(order.getId()).orElseThrow().getStatus());
assertThat(inventory.getStock("GADGET-ORCH-FAIL")).isEqualTo(1);
assertThat(payments.findByOrderId(order.getId()).orElseThrow().getStatus())
.isEqualTo(PaymentStatus.REFUNDED);
}
}