diff --git a/README.md b/README.md
index 6479ec9..4d24147 100644
--- a/README.md
+++ b/README.md
@@ -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
diff --git a/saga/README.md b/saga/README.md
new file mode 100644
index 0000000..72d5026
--- /dev/null
+++ b/saga/README.md
@@ -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.
diff --git a/saga/output/00-choreography-happy-path.txt b/saga/output/00-choreography-happy-path.txt
new file mode 100644
index 0000000..03be1d5
--- /dev/null
+++ b/saga/output/00-choreography-happy-path.txt
@@ -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".
diff --git a/saga/output/01-choreography-compensation.txt b/saga/output/01-choreography-compensation.txt
new file mode 100644
index 0000000..2215c46
--- /dev/null
+++ b/saga/output/01-choreography-compensation.txt
@@ -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.
diff --git a/saga/output/02-orchestration-saga.txt b/saga/output/02-orchestration-saga.txt
new file mode 100644
index 0000000..6a2df6f
--- /dev/null
+++ b/saga/output/02-orchestration-saga.txt
@@ -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.
diff --git a/saga/output/03-idempotent-consumer-duplicate.txt b/saga/output/03-idempotent-consumer-duplicate.txt
new file mode 100644
index 0000000..8715bb4
--- /dev/null
+++ b/saga/output/03-idempotent-consumer-duplicate.txt
@@ -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.
diff --git a/saga/pom.xml b/saga/pom.xml
new file mode 100644
index 0000000..d8828b8
--- /dev/null
+++ b/saga/pom.xml
@@ -0,0 +1,74 @@
+
+ 4.0.0
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ 4.1.1
+
+
+
+ com.ankurm
+ saga
+ 1.0
+ jar
+
+
+ 25
+ UTF-8
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter
+
+
+
+ org.springframework.boot
+ spring-boot-starter-kafka
+
+
+ org.springframework.boot
+ spring-boot-starter-jackson
+
+
+ org.springframework.boot
+ spring-boot-starter-data-jpa
+
+
+ com.h2database
+ h2
+ runtime
+
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+ org.springframework.kafka
+ spring-kafka-test
+ test
+
+
+ org.awaitility
+ awaitility
+ test
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-maven-plugin
+
+
+
+
diff --git a/saga/src/main/java/com/ankurm/sagademo/SagaDemoApplication.java b/saga/src/main/java/com/ankurm/sagademo/SagaDemoApplication.java
new file mode 100644
index 0000000..5f2f998
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/SagaDemoApplication.java
@@ -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);
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/Topics.java b/saga/src/main/java/com/ankurm/sagademo/Topics.java
new file mode 100644
index 0000000..1a03466
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/Topics.java
@@ -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";
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/choreography/ChoreographySagaStarter.java b/saga/src/main/java/com/ankurm/sagademo/choreography/ChoreographySagaStarter.java
new file mode 100644
index 0000000..af23491
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/choreography/ChoreographySagaStarter.java
@@ -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 kafka;
+
+ public ChoreographySagaStarter(OrderService orders, KafkaTemplate 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;
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/choreography/InventoryChoreographyListener.java b/saga/src/main/java/com/ankurm/sagademo/choreography/InventoryChoreographyListener.java
new file mode 100644
index 0000000..62e34ad
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/choreography/InventoryChoreographyListener.java
@@ -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 kafka;
+
+ public InventoryChoreographyListener(InventoryService inventory, KafkaTemplate 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()));
+ }
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/choreography/OrderChoreographyListener.java b/saga/src/main/java/com/ankurm/sagademo/choreography/OrderChoreographyListener.java
new file mode 100644
index 0000000..307d45c
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/choreography/OrderChoreographyListener.java
@@ -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());
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/choreography/PaymentChoreographyListener.java b/saga/src/main/java/com/ankurm/sagademo/choreography/PaymentChoreographyListener.java
new file mode 100644
index 0000000..937a904
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/choreography/PaymentChoreographyListener.java
@@ -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 kafka;
+
+ public PaymentChoreographyListener(PaymentService payments, KafkaTemplate 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()));
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/InventoryRejected.java b/saga/src/main/java/com/ankurm/sagademo/events/InventoryRejected.java
new file mode 100644
index 0000000..0bd8c6c
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/InventoryRejected.java
@@ -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) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/InventoryReply.java b/saga/src/main/java/com/ankurm/sagademo/events/InventoryReply.java
new file mode 100644
index 0000000..f2f0f45
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/InventoryReply.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record InventoryReply(String commandId, Long orderId, boolean success, String reason) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/InventoryReserved.java b/saga/src/main/java/com/ankurm/sagademo/events/InventoryReserved.java
new file mode 100644
index 0000000..01fb56a
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/InventoryReserved.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record InventoryReserved(Long orderId) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/OrderCreated.java b/saga/src/main/java/com/ankurm/sagademo/events/OrderCreated.java
new file mode 100644
index 0000000..7129dbe
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/OrderCreated.java
@@ -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) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/PaymentRefunded.java b/saga/src/main/java/com/ankurm/sagademo/events/PaymentRefunded.java
new file mode 100644
index 0000000..52e474c
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/PaymentRefunded.java
@@ -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) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/PaymentReply.java b/saga/src/main/java/com/ankurm/sagademo/events/PaymentReply.java
new file mode 100644
index 0000000..faab8e9
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/PaymentReply.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record PaymentReply(String commandId, Long orderId, boolean success) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/PaymentReserved.java b/saga/src/main/java/com/ankurm/sagademo/events/PaymentReserved.java
new file mode 100644
index 0000000..d78db86
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/PaymentReserved.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record PaymentReserved(Long orderId, String sku, int qty) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/RefundPaymentCommand.java b/saga/src/main/java/com/ankurm/sagademo/events/RefundPaymentCommand.java
new file mode 100644
index 0000000..c4a1a5c
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/RefundPaymentCommand.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record RefundPaymentCommand(String commandId, Long orderId) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/RefundReply.java b/saga/src/main/java/com/ankurm/sagademo/events/RefundReply.java
new file mode 100644
index 0000000..d091f0d
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/RefundReply.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record RefundReply(String commandId, Long orderId) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/ReserveInventoryCommand.java b/saga/src/main/java/com/ankurm/sagademo/events/ReserveInventoryCommand.java
new file mode 100644
index 0000000..f527f2f
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/ReserveInventoryCommand.java
@@ -0,0 +1,4 @@
+package com.ankurm.sagademo.events;
+
+public record ReserveInventoryCommand(String commandId, Long orderId, String sku, int qty) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/events/ReservePaymentCommand.java b/saga/src/main/java/com/ankurm/sagademo/events/ReservePaymentCommand.java
new file mode 100644
index 0000000..9c2ee30
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/events/ReservePaymentCommand.java
@@ -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) {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/idempotency/IdempotencyGuard.java b/saga/src/main/java/com/ankurm/sagademo/idempotency/IdempotencyGuard.java
new file mode 100644
index 0000000..245fa39
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/idempotency/IdempotencyGuard.java
@@ -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;
+ }
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommand.java b/saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommand.java
new file mode 100644
index 0000000..846ffe3
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommand.java
@@ -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;
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommandRepository.java b/saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommandRepository.java
new file mode 100644
index 0000000..4a5d521
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommandRepository.java
@@ -0,0 +1,8 @@
+package com.ankurm.sagademo.idempotency;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+
+public interface ProcessedCommandRepository extends JpaRepository {
+
+ long countByCommandId(String commandId);
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/inventory/InventoryService.java b/saga/src/main/java/com/ankurm/sagademo/inventory/InventoryService.java
new file mode 100644
index 0000000..73db4aa
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/inventory/InventoryService.java
@@ -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 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;
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orchestration/InventoryOrchestrationHandler.java b/saga/src/main/java/com/ankurm/sagademo/orchestration/InventoryOrchestrationHandler.java
new file mode 100644
index 0000000..5bda80a
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orchestration/InventoryOrchestrationHandler.java
@@ -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 kafka;
+
+ public InventoryOrchestrationHandler(InventoryService inventory, IdempotencyGuard idempotencyGuard,
+ KafkaTemplate 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));
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orchestration/OrderSagaOrchestrator.java b/saga/src/main/java/com/ankurm/sagademo/orchestration/OrderSagaOrchestrator.java
new file mode 100644
index 0000000..c5c6f58
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orchestration/OrderSagaOrchestrator.java
@@ -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 kafka;
+
+ public OrderSagaOrchestrator(OrderService orders, OrderRepository orderRepository,
+ KafkaTemplate 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());
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orchestration/PaymentOrchestrationHandler.java b/saga/src/main/java/com/ankurm/sagademo/orchestration/PaymentOrchestrationHandler.java
new file mode 100644
index 0000000..8634f92
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orchestration/PaymentOrchestrationHandler.java
@@ -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 kafka;
+
+ public PaymentOrchestrationHandler(PaymentService payments, KafkaTemplate 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()));
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orders/Order.java b/saga/src/main/java/com/ankurm/sagademo/orders/Order.java
new file mode 100644
index 0000000..543018f
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orders/Order.java
@@ -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;
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orders/OrderRepository.java b/saga/src/main/java/com/ankurm/sagademo/orders/OrderRepository.java
new file mode 100644
index 0000000..ade917d
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orders/OrderRepository.java
@@ -0,0 +1,6 @@
+package com.ankurm.sagademo.orders;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+
+public interface OrderRepository extends JpaRepository {
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orders/OrderService.java b/saga/src/main/java/com/ankurm/sagademo/orders/OrderService.java
new file mode 100644
index 0000000..d1fd0b3
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orders/OrderService.java
@@ -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);
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/orders/OrderStatus.java b/saga/src/main/java/com/ankurm/sagademo/orders/OrderStatus.java
new file mode 100644
index 0000000..5e456d0
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/orders/OrderStatus.java
@@ -0,0 +1,7 @@
+package com.ankurm.sagademo.orders;
+
+public enum OrderStatus {
+ PENDING,
+ CONFIRMED,
+ CANCELLED
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/payments/PaymentRecord.java b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentRecord.java
new file mode 100644
index 0000000..f74d283
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentRecord.java
@@ -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;
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/payments/PaymentRepository.java b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentRepository.java
new file mode 100644
index 0000000..e54590c
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentRepository.java
@@ -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 {
+ Optional findByOrderId(Long orderId);
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/payments/PaymentService.java b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentService.java
new file mode 100644
index 0000000..c61f059
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentService.java
@@ -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);
+ }
+}
diff --git a/saga/src/main/java/com/ankurm/sagademo/payments/PaymentStatus.java b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentStatus.java
new file mode 100644
index 0000000..be1ca16
--- /dev/null
+++ b/saga/src/main/java/com/ankurm/sagademo/payments/PaymentStatus.java
@@ -0,0 +1,6 @@
+package com.ankurm.sagademo.payments;
+
+public enum PaymentStatus {
+ RESERVED,
+ REFUNDED
+}
diff --git a/saga/src/main/resources/application.properties b/saga/src/main/resources/application.properties
new file mode 100644
index 0000000..14e8ba0
--- /dev/null
+++ b/saga/src/main/resources/application.properties
@@ -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
diff --git a/saga/src/test/java/com/ankurm/sagademo/ChoreographyCompensationTest.java b/saga/src/test/java/com/ankurm/sagademo/ChoreographyCompensationTest.java
new file mode 100644
index 0000000..df2501a
--- /dev/null
+++ b/saga/src/test/java/com/ankurm/sagademo/ChoreographyCompensationTest.java
@@ -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);
+ }
+}
diff --git a/saga/src/test/java/com/ankurm/sagademo/ChoreographyHappyPathTest.java b/saga/src/test/java/com/ankurm/sagademo/ChoreographyHappyPathTest.java
new file mode 100644
index 0000000..99b4e48
--- /dev/null
+++ b/saga/src/test/java/com/ankurm/sagademo/ChoreographyHappyPathTest.java
@@ -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);
+ }
+}
diff --git a/saga/src/test/java/com/ankurm/sagademo/IdempotentConsumerTest.java b/saga/src/test/java/com/ankurm/sagademo/IdempotentConsumerTest.java
new file mode 100644
index 0000000..f423a6c
--- /dev/null
+++ b/saga/src/test/java/com/ankurm/sagademo/IdempotentConsumerTest.java
@@ -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 kafka;
+ @Autowired
+ InventoryService inventory;
+ @Autowired
+ ProcessedCommandRepository processedCommands;
+ @Autowired
+ org.springframework.kafka.test.EmbeddedKafkaBroker embeddedKafka;
+
+ Consumer 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(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 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 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);
+ }
+}
diff --git a/saga/src/test/java/com/ankurm/sagademo/OrchestrationSagaTest.java b/saga/src/test/java/com/ankurm/sagademo/OrchestrationSagaTest.java
new file mode 100644
index 0000000..6513320
--- /dev/null
+++ b/saga/src/test/java/com/ankurm/sagademo/OrchestrationSagaTest.java
@@ -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);
+ }
+}