From 4a5d3318bf6f427a065611bc6a898a7137e34589 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 3 Oct 2026 19:56:40 +0000 Subject: [PATCH] Add the saga module: Saga Pattern in Spring Boot, orchestration vs choreography with Kafka --- README.md | 8 +- saga/README.md | 61 ++++++++++ saga/output/00-choreography-happy-path.txt | 17 +++ saga/output/01-choreography-compensation.txt | 15 +++ saga/output/02-orchestration-saga.txt | 20 ++++ .../03-idempotent-consumer-duplicate.txt | 16 +++ saga/pom.xml | 74 ++++++++++++ .../ankurm/sagademo/SagaDemoApplication.java | 12 ++ .../main/java/com/ankurm/sagademo/Topics.java | 34 ++++++ .../choreography/ChoreographySagaStarter.java | 30 +++++ .../InventoryChoreographyListener.java | 40 +++++++ .../OrderChoreographyListener.java | 34 ++++++ .../PaymentChoreographyListener.java | 47 ++++++++ .../sagademo/events/InventoryRejected.java | 6 + .../sagademo/events/InventoryReply.java | 4 + .../sagademo/events/InventoryReserved.java | 4 + .../ankurm/sagademo/events/OrderCreated.java | 7 ++ .../sagademo/events/PaymentRefunded.java | 6 + .../ankurm/sagademo/events/PaymentReply.java | 4 + .../sagademo/events/PaymentReserved.java | 4 + .../sagademo/events/RefundPaymentCommand.java | 4 + .../ankurm/sagademo/events/RefundReply.java | 4 + .../events/ReserveInventoryCommand.java | 4 + .../events/ReservePaymentCommand.java | 8 ++ .../idempotency/IdempotencyGuard.java | 34 ++++++ .../idempotency/ProcessedCommand.java | 39 +++++++ .../ProcessedCommandRepository.java | 8 ++ .../sagademo/inventory/InventoryService.java | 37 ++++++ .../InventoryOrchestrationHandler.java | 45 ++++++++ .../orchestration/OrderSagaOrchestrator.java | 76 +++++++++++++ .../PaymentOrchestrationHandler.java | 38 +++++++ .../com/ankurm/sagademo/orders/Order.java | 66 +++++++++++ .../sagademo/orders/OrderRepository.java | 6 + .../ankurm/sagademo/orders/OrderService.java | 31 ++++++ .../ankurm/sagademo/orders/OrderStatus.java | 7 ++ .../sagademo/payments/PaymentRecord.java | 49 ++++++++ .../sagademo/payments/PaymentRepository.java | 9 ++ .../sagademo/payments/PaymentService.java | 31 ++++++ .../sagademo/payments/PaymentStatus.java | 6 + .../src/main/resources/application.properties | 12 ++ .../ChoreographyCompensationTest.java | 76 +++++++++++++ .../sagademo/ChoreographyHappyPathTest.java | 67 +++++++++++ .../sagademo/IdempotentConsumerTest.java | 105 ++++++++++++++++++ .../sagademo/OrchestrationSagaTest.java | 94 ++++++++++++++++ 44 files changed, 1296 insertions(+), 3 deletions(-) create mode 100644 saga/README.md create mode 100644 saga/output/00-choreography-happy-path.txt create mode 100644 saga/output/01-choreography-compensation.txt create mode 100644 saga/output/02-orchestration-saga.txt create mode 100644 saga/output/03-idempotent-consumer-duplicate.txt create mode 100644 saga/pom.xml create mode 100644 saga/src/main/java/com/ankurm/sagademo/SagaDemoApplication.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/Topics.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/choreography/ChoreographySagaStarter.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/choreography/InventoryChoreographyListener.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/choreography/OrderChoreographyListener.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/choreography/PaymentChoreographyListener.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/InventoryRejected.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/InventoryReply.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/InventoryReserved.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/OrderCreated.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/PaymentRefunded.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/PaymentReply.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/PaymentReserved.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/RefundPaymentCommand.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/RefundReply.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/ReserveInventoryCommand.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/events/ReservePaymentCommand.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/idempotency/IdempotencyGuard.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommand.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/idempotency/ProcessedCommandRepository.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/inventory/InventoryService.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orchestration/InventoryOrchestrationHandler.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orchestration/OrderSagaOrchestrator.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orchestration/PaymentOrchestrationHandler.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orders/Order.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orders/OrderRepository.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orders/OrderService.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/orders/OrderStatus.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/payments/PaymentRecord.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/payments/PaymentRepository.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/payments/PaymentService.java create mode 100644 saga/src/main/java/com/ankurm/sagademo/payments/PaymentStatus.java create mode 100644 saga/src/main/resources/application.properties create mode 100644 saga/src/test/java/com/ankurm/sagademo/ChoreographyCompensationTest.java create mode 100644 saga/src/test/java/com/ankurm/sagademo/ChoreographyHappyPathTest.java create mode 100644 saga/src/test/java/com/ankurm/sagademo/IdempotentConsumerTest.java create mode 100644 saga/src/test/java/com/ankurm/sagademo/OrchestrationSagaTest.java 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); + } +}