From b4a623b8893a12fc33077b692face658c4369a22 Mon Sep 17 00:00:00 2001 From: asmhatre Date: Sat, 3 Oct 2026 20:59:55 +0000 Subject: [PATCH] Add cqrs module: CQRS in Spring Boot Without a Framework Separate write/read DataSources, a JPA command side, and two interchangeable projection listeners (sync and @Async) demonstrating the real latency-versus- freshness trade-off CQRS forces. Includes a reflection-based proof that the query side has no dependency on the write side, and a real failure/fix transcript for the -parameters compiler flag this standalone reactor doesn't inherit from spring-boot-starter-parent. --- cqrs/README.md | 102 +++++++++++++++++ cqrs/output/00-sync-profile-latency.txt | 20 ++++ .../01-async-profile-staleness-window.txt | 27 +++++ .../02-query-side-architecture-proof.txt | 17 +++ .../03-missing-parameters-flag-failure.txt | 21 ++++ cqrs/pom.xml | 107 ++++++++++++++++++ .../java/com/ankurm/cqrsdemo/Application.java | 26 +++++ .../cqrsdemo/command/OrderCancelled.java | 5 + .../command/OrderCommandController.java | 53 +++++++++ .../cqrsdemo/command/OrderCommandService.java | 65 +++++++++++ .../ankurm/cqrsdemo/command/OrderPlaced.java | 14 +++ .../ankurm/cqrsdemo/command/OrderShipped.java | 5 + .../cqrsdemo/command/internal/Order.java | 81 +++++++++++++ .../cqrsdemo/command/internal/OrderLine.java | 43 +++++++ .../command/internal/OrderRepository.java | 13 +++ .../cqrsdemo/config/DataSourceConfig.java | 49 ++++++++ .../cqrsdemo/query/OrderQueryController.java | 51 +++++++++ .../ankurm/cqrsdemo/query/OrderSummary.java | 8 ++ .../read/AsyncOrderSummaryProjection.java | 48 ++++++++ .../cqrsdemo/read/OrderSummaryWriter.java | 56 +++++++++ .../ankurm/cqrsdemo/read/ReadSideConfig.java | 37 ++++++ .../read/SyncOrderSummaryProjection.java | 43 +++++++ .../src/main/resources/application.properties | 23 ++++ .../AsyncProfileStalenessWindowTest.java | 76 +++++++++++++ .../cqrsdemo/QuerySideArchitectureTest.java | 43 +++++++ .../cqrsdemo/SyncProfileLatencyTest.java | 51 +++++++++ pom.xml | 2 + 27 files changed, 1086 insertions(+) create mode 100644 cqrs/README.md create mode 100644 cqrs/output/00-sync-profile-latency.txt create mode 100644 cqrs/output/01-async-profile-staleness-window.txt create mode 100644 cqrs/output/02-query-side-architecture-proof.txt create mode 100644 cqrs/output/03-missing-parameters-flag-failure.txt create mode 100644 cqrs/pom.xml create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/Application.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCancelled.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandController.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandService.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderPlaced.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderShipped.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/Order.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderLine.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderRepository.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/config/DataSourceConfig.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderQueryController.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderSummary.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/read/AsyncOrderSummaryProjection.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/read/OrderSummaryWriter.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/read/ReadSideConfig.java create mode 100644 cqrs/src/main/java/com/ankurm/cqrsdemo/read/SyncOrderSummaryProjection.java create mode 100644 cqrs/src/main/resources/application.properties create mode 100644 cqrs/src/test/java/com/ankurm/cqrsdemo/AsyncProfileStalenessWindowTest.java create mode 100644 cqrs/src/test/java/com/ankurm/cqrsdemo/QuerySideArchitectureTest.java create mode 100644 cqrs/src/test/java/com/ankurm/cqrsdemo/SyncProfileLatencyTest.java diff --git a/cqrs/README.md b/cqrs/README.md new file mode 100644 index 0000000..6477852 --- /dev/null +++ b/cqrs/README.md @@ -0,0 +1,102 @@ +# cqrs + +Companion project for the article **[CQRS in Spring Boot Without a Framework: Separate Read Models and Projections](https://ankurm.com/cqrs-spring-boot-without-a-framework-separate-read-models-projections/)** on **[ankurm.com](https://ankurm.com)**. + +There is deliberately **no `docs/` folder**: the deeper material lives in collapsible "going +deeper" sections inside the article itself, next to the paragraph each one extends. + +"Without a framework" means exactly that: no Axon, no Spring Modulith event publication +registry (that's the companion [`outbox`](../outbox) module, a different post), no event +store. Just plain Spring -- `ApplicationEventPublisher`, `@TransactionalEventListener`, +`@Async`, two `DataSource` beans, and `JdbcTemplate`. + +## Versions + +| | | +|---|---| +| Spring Boot | 4.1.1 | +| Spring Framework | 7.0.9 | +| JDK | 25 (Temurin 25.0.4.1+1) | +| Maven | 3.9 | +| H2 | 2.4.240 | +| Hibernate ORM | 7.4.5.Final | + +## The shape + +Two databases, one Spring Boot application: + +| | Database | Reached through | Who writes to it | +|---|---|---|---| +| Write side | `jdbc:h2:mem:cqrs-write` | Spring Data JPA (`OrderRepository`) | `OrderCommandService` only | +| Read side | `jdbc:h2:mem:cqrs-read` | a single `JdbcTemplate` bean (`readJdbcTemplate`) | the projection listener only | + +`OrderCommandService` changes the write side and publishes an event in the same +transaction. A projection listener reacts to that event after the transaction commits and +updates the read side. `OrderQueryController` reads only from the read side -- see +[`QuerySideArchitectureTest`](src/test/java/com/ankurm/cqrsdemo/QuerySideArchitectureTest.java), +which checks that by reflection rather than by convention. + +## The two profiles + +The same projection write (`OrderSummaryWriter`, with a configurable simulated delay to +stand in for whatever a real read-model write costs) is called two different ways: + +| Profile | Listener | Command path | Read model | +|---|---|---|---| +| *(default)* | [`SyncOrderSummaryProjection`](src/main/java/com/ankurm/cqrsdemo/read/SyncOrderSummaryProjection.java) | pays the projection's write cost | consistent the instant the command returns | +| `async-projection` | [`AsyncOrderSummaryProjection`](src/main/java/com/ankurm/cqrsdemo/read/AsyncOrderSummaryProjection.java) | does not pay it | consistent a short, measurable while later | + +```bash +mvn -pl cqrs test -Dtest=SyncProfileLatencyTest # default profile +mvn -pl cqrs test -Dtest=AsyncProfileStalenessWindowTest # -Dspring.profiles.active=async-projection, set in the test +``` + +## Quickstart + +```bash +export JAVA_HOME=/path/to/jdk-25 +mvn -pl cqrs -am spring-boot:run # default (sync) profile +mvn -pl cqrs -am spring-boot:run -Dspring-boot.run.profiles=async-projection + +curl -X POST localhost:8080/orders \ + -H 'Content-Type: application/json' \ + -d '{"customerName":"Priya","lines":[{"sku":"USB-C-HUB","quantity":1,"unitPriceCents":4999}]}' + +curl localhost:8080/order-summaries +``` + +## Captured output + +| File | What it shows | +|---|---| +| [`output/00-sync-profile-latency.txt`](output/00-sync-profile-latency.txt) | the default profile: command latency includes the projection write, and the read model is already correct the instant the command returns | +| [`output/01-async-profile-staleness-window.txt`](output/01-async-profile-staleness-window.txt) | the `async-projection` profile: the same command returns in milliseconds, and a query that lands inside the staleness window gets a real 404 for an order that was just placed | +| [`output/02-query-side-architecture-proof.txt`](output/02-query-side-architecture-proof.txt) | reflection over the compiled `OrderQueryController` class proving its only dependency is the read-side `JdbcTemplate` | +| [`output/03-missing-parameters-flag-failure.txt`](output/03-missing-parameters-flag-failure.txt) | the real 500 and `IllegalArgumentException` that come back from `GET /order-summaries/{orderId}` if this module is built without `true` -- captured by actually removing it and running the test, not reconstructed from memory | + +## A note for anyone copying this shape + +- **`DataSourceProperties` moved** in Spring Boot 4: it's now + `org.springframework.boot.jdbc.autoconfigure.DataSourceProperties` in artifact + `spring-boot-jdbc`, not `org.springframework.boot.autoconfigure.jdbc.DataSourceProperties`. + Same family of split as `@DataJpaTest` and `@WebMvcTest` in the + [`hexagonal`](../../spring-boot-demo/hexagonal) module's notes. +- **This module, like `order-fulfillment` and `outbox`, does not inherit + `spring-boot-starter-parent`**, so it does not get `-parameters` passed to `javac` for + free. Any `@PathVariable`/`@RequestParam` relying on its own parameter name (no explicit + `name = "..."`) fails at request time, not at compile time -- see + `output/03-missing-parameters-flag-failure.txt`. Worth flagging on the way past: + `order-fulfillment`'s `OrderController` has the exact same unqualified `@PathVariable`, + in the exact same kind of module, and has never hit this -- its own test suite exercises + that controller only through Spring Modulith's `PublishedEventsExtension`, never a real + HTTP call to `GET /orders/{id}`. The bug is latent there; this module doesn't touch that + one, this just flags it. +- **The first HTTP request against a freshly started test context is not a fair baseline.** + `AsyncProfileStalenessWindowTest` originally measured that first request directly and + the ~370ms of Hikari/Hibernate first-use cost looked like it came from the projection. + A throwaway warm-up request before the measured one fixed it -- see the comment at the + top of that test and the note in `output/01-async-profile-staleness-window.txt`. + +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/cqrs/output/00-sync-profile-latency.txt b/cqrs/output/00-sync-profile-latency.txt new file mode 100644 index 0000000..589267b --- /dev/null +++ b/cqrs/output/00-sync-profile-latency.txt @@ -0,0 +1,20 @@ +$ mvn -pl cqrs test -Dtest=SyncProfileLatencyTest +(default profile -- SyncOrderSummaryProjection active; its write to order_summary runs + on the same thread that is about to answer the command's HTTP request) + +2026-10-04T02:27:21.596+05:30 INFO 17202 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-1 - Added connection conn0: url=jdbc:h2:mem:cqrs-write user=SA +2026-10-04T02:27:22.729+05:30 INFO 17202 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-2 - Added connection conn10: url=jdbc:h2:mem:cqrs-read user=SA + +POST /orders (sync profile) took 670 ms, orderId=9ee5b016-7a81-40f9-8d7a-d7d1959c09be +Immediately after POST returned, GET /order-summaries/9ee5b016-7a81-40f9-8d7a-d7d1959c09be -> status=200 OK, body=OrderSummary[orderId=9ee5b016-7a81-40f9-8d7a-d7d1959c09be, customerName=Priya, itemCount=2, totalCents=7597, status=PLACED, updatedAt=2026-10-03T20:57:24.283700Z] + +[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0 + +Two real HikariCP pools, one per DataSource bean -- that line alone is a cheap way to +confirm the write and read sides really are two separate connection pools to two separate +H2 databases, not one database with two tables. The 670ms includes the 300ms the projection's +simulated write delay adds on top of real JPA/Hibernate first-request costs (see +output/01-async-profile-staleness-window.txt's warm-up note for what that baseline actually +is) -- and because that write finished before the HTTP response did, the very next GET +against the read model, with no wait at all, already sees itemCount, totalCents and status +correctly populated. diff --git a/cqrs/output/01-async-profile-staleness-window.txt b/cqrs/output/01-async-profile-staleness-window.txt new file mode 100644 index 0000000..ee2f4f5 --- /dev/null +++ b/cqrs/output/01-async-profile-staleness-window.txt @@ -0,0 +1,27 @@ +$ mvn -pl cqrs test -Dtest=AsyncProfileStalenessWindowTest +(async-projection profile -- AsyncOrderSummaryProjection active; @Async hands the same + write off to a separate thread instead of running it inline. A throwaway warm-up POST + runs first and is not measured: the first HTTP request against a freshly started test + context pays Hikari/Hibernate first-use costs that have nothing to do with this + profile's actual behaviour -- measuring that cold first request directly made POST + /orders look like it took ~370ms even in this profile, which would have buried the + exact thing the test exists to show.) + +2026-10-04T02:27:31.206+05:30 INFO 17346 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-1 - Added connection conn0: url=jdbc:h2:mem:cqrs-write user=SA +2026-10-04T02:27:32.304+05:30 INFO 17346 --- [cqrs] [ main] com.zaxxer.hikari.pool.HikariPool : HikariPool-2 - Added connection conn10: url=jdbc:h2:mem:cqrs-read user=SA + +POST /orders (async profile) took 14 ms, orderId=8451a829-7ef6-4d19-b839-2d9ce5885437 +Immediately after POST returned (inside the staleness window), GET /order-summaries/8451a829-7ef6-4d19-b839-2d9ce5885437 -> status=404 NOT_FOUND, body=null +After waiting for the projection (outside the staleness window), GET /order-summaries/8451a829-7ef6-4d19-b839-2d9ce5885437 -> status=200 OK, body=OrderSummary[orderId=8451a829-7ef6-4d19-b839-2d9ce5885437, customerName=Dev, itemCount=1, totalCents=8999, status=PLACED, updatedAt=2026-10-03T20:57:33.971972Z] + +[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0 + +This is the trade-off in one transcript, warmed-up and measured honestly. The identical +command, in the identical module, against the identical 300ms simulated write cost, took +14ms here instead of 670ms -- because the projection's write moved to another thread and +the HTTP response stopped waiting for it. The price is the 404 on the very next line: a +client that queries in that window, a few dozen milliseconds wide in this run, gets a +"not found" for an order that was just placed successfully. Only after the test waits +(Awaitility, polling every 50ms, up to 2 seconds) does the same GET return 200 with the +order correctly summarized. Neither profile is "more correct" than the other; they are +the same correctness guarantee, paid for at a different point in the request. diff --git a/cqrs/output/02-query-side-architecture-proof.txt b/cqrs/output/02-query-side-architecture-proof.txt new file mode 100644 index 0000000..cdf7ee8 --- /dev/null +++ b/cqrs/output/02-query-side-architecture-proof.txt @@ -0,0 +1,17 @@ +$ mvn -pl cqrs test -Dtest=QuerySideArchitectureTest +(plain JUnit reflection against the compiled OrderQueryController class -- no Spring + context started, because this is a claim about the class file, not about runtime wiring) + +field readJdbcTemplate : org.springframework.jdbc.core.JdbcTemplate +OrderQueryController constructor parameter types: [class org.springframework.jdbc.core.JdbcTemplate] + +[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0 + +The whole proof is in those two lines: OrderQueryController has exactly one declared +field and exactly one constructor parameter, and both are JdbcTemplate -- not +OrderRepository, not an EntityManager, not anything from the com.ankurm.cqrsdemo.command +package. A reviewer does not have to trust the Javadoc comment on the class or remember +to check it by hand on the next change; this test fails the build the moment anyone adds +a write-side dependency to the query side, the same way the hexagonal-architecture post's +`mvn dependency:tree` turns "the core has no Spring dependency" from a promise into +something Maven enforces. diff --git a/cqrs/output/03-missing-parameters-flag-failure.txt b/cqrs/output/03-missing-parameters-flag-failure.txt new file mode 100644 index 0000000..efa1f03 --- /dev/null +++ b/cqrs/output/03-missing-parameters-flag-failure.txt @@ -0,0 +1,21 @@ +Captured by temporarily building this module without the true +maven-compiler-plugin configuration (this module does not inherit spring-boot-starter-parent, +which is what normally adds that flag for free), then running the SyncProfileLatencyTest, +which calls GET /order-summaries/{orderId} over real HTTP. The real, unedited response and +server-side exception: + +DEBUG raw body: {"timestamp":"2026-10-03T20:54:06.783Z","status":500,"error":"Internal Server Error","path":"/order-summaries/c6fc4ad0-f98e-4e05-92a8-8717e9bffced"} +2026-10-04T02:24:06.801+05:30 ERROR 16353 --- [cqrs] [o-auto-1-exec-3] o.a.c.c.C.[.[.[/].[dispatcherServlet] : Servlet.service() for servlet [dispatcherServlet] in context with path [] threw exception [Request processing failed: java.lang.IllegalArgumentException: Name for argument of type [java.lang.String] not specified, and parameter name information not available via reflection. Ensure that the compiler uses the '-parameters' flag.] with root cause + +java.lang.IllegalArgumentException: Name for argument of type [java.lang.String] not specified, and parameter name information not available via reflection. Ensure that the compiler uses the '-parameters' flag. + +Re-adding true to cqrs/pom.xml's maven-compiler-plugin configuration +and re-running the exact same test produces the green run captured in +00-sync-profile-latency.txt, with a 200 and a real body instead of this 500. + +Worth noting on the way past: order-fulfillment's OrderController (post #31) has the exact +same unqualified @PathVariable Long id, in the exact same kind of standalone-reactor module +with no inherited parent -- it has just never been caught, because that module's own test +suite exercises its event-publication behaviour through Spring Modulith's +PublishedEventsExtension rather than through a real HTTP call to GET /orders/{id}. The bug is +latent there, not fixed; flagged separately rather than changed here. diff --git a/cqrs/pom.xml b/cqrs/pom.xml new file mode 100644 index 0000000..92d3a3c --- /dev/null +++ b/cqrs/pom.xml @@ -0,0 +1,107 @@ + + + 4.0.0 + + + com.ankurm + spring-modulith-demo + 1.0.0 + + + cqrs + jar + cqrs + Command/query separation with plain Spring: a JPA write model, a second H2 database + holding a denormalized read model that the query side reaches only through its own JdbcTemplate, + and two interchangeable projection listeners (synchronous and @Async) that show the actual + latency-versus-freshness trade-off CQRS makes you choose, not just assert. + + + 25 + + + + + + org.springframework.boot + spring-boot-dependencies + ${spring-boot.version} + pom + import + + + + + + + org.springframework.boot + spring-boot-starter-web + + + org.springframework.boot + spring-boot-starter-data-jpa + + + org.springframework.boot + spring-boot-starter-jdbc + + + com.h2database + h2 + runtime + + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.awaitility + awaitility + test + + + org.springframework.boot + spring-boot-resttestclient + test + + + org.springframework.boot + spring-boot-restclient + test + + + org.springframework.boot + spring-boot-configuration-processor + provided + + + + + cqrs + + + org.apache.maven.plugins + maven-compiler-plugin + + + true + + + + org.springframework.boot + spring-boot-maven-plugin + + + + diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/Application.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/Application.java new file mode 100644 index 0000000..d7a155c --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/Application.java @@ -0,0 +1,26 @@ +package com.ankurm.cqrsdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.scheduling.annotation.EnableAsync; + +/** + * Two stores, one Spring Boot app. The write side is a normal + * {@code spring-boot-starter-data-jpa} setup against {@code jdbc:h2:mem:cqrs-write}. + * The read side is a second, independent {@link javax.sql.DataSource} against + * {@code jdbc:h2:mem:cqrs-read}, reached only through {@link org.springframework.jdbc.core.JdbcTemplate} + * — see {@link com.ankurm.cqrsdemo.read.ReadSideConfig}. + * + *

{@code @EnableAsync} backs the {@code async-projection} profile's + * {@link com.ankurm.cqrsdemo.read.AsyncOrderSummaryProjection}. Without it the + * {@code @Async} annotation on that listener is silently ignored and it runs + * synchronously anyway — a real way to lose the "fast write" half of this + * trade-off without any error telling you so. + */ +@SpringBootApplication +@EnableAsync +public class Application { + public static void main(String[] args) { + SpringApplication.run(Application.class, args); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCancelled.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCancelled.java new file mode 100644 index 0000000..a26dc49 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCancelled.java @@ -0,0 +1,5 @@ +package com.ankurm.cqrsdemo.command; + +/** Published after a {@code cancel()} command commits. */ +public record OrderCancelled(String orderId) { +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandController.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandController.java new file mode 100644 index 0000000..3dbd8a1 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandController.java @@ -0,0 +1,53 @@ +package com.ankurm.cqrsdemo.command; + +import com.ankurm.cqrsdemo.command.internal.OrderLine; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RestController; + +import java.util.List; + +/** + * The only HTTP entry point into the write side. Note what it returns: just the new + * order's id, nothing shaped for display. If a caller wants to show the order back to + * someone, that is a query — {@link com.ankurm.cqrsdemo.query.OrderQueryController}, + * a different controller talking to a different store. + */ +@RestController +public class OrderCommandController { + + private final OrderCommandService commandService; + + public OrderCommandController(OrderCommandService commandService) { + this.commandService = commandService; + } + + public record LineRequest(String sku, int quantity, long unitPriceCents) { + } + + public record PlaceOrderRequest(String customerName, List lines) { + } + + public record PlaceOrderResponse(String orderId) { + } + + @PostMapping("/orders") + public PlaceOrderResponse placeOrder(@RequestBody PlaceOrderRequest request) { + List lines = request.lines().stream() + .map(l -> new OrderLine(l.sku(), l.quantity(), l.unitPriceCents())) + .toList(); + String orderId = commandService.placeOrder(request.customerName(), lines); + return new PlaceOrderResponse(orderId); + } + + @PostMapping("/orders/{orderId}/ship") + public void ship(@PathVariable String orderId) { + commandService.shipOrder(orderId); + } + + @PostMapping("/orders/{orderId}/cancel") + public void cancel(@PathVariable String orderId) { + commandService.cancelOrder(orderId); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandService.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandService.java new file mode 100644 index 0000000..39b528b --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderCommandService.java @@ -0,0 +1,65 @@ +package com.ankurm.cqrsdemo.command; + +import com.ankurm.cqrsdemo.command.internal.Order; +import com.ankurm.cqrsdemo.command.internal.OrderLine; +import com.ankurm.cqrsdemo.command.internal.OrderRepository; +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import java.util.List; +import java.util.NoSuchElementException; +import java.util.UUID; + +/** + * Every command method does exactly two things in one transaction: change the + * {@code orders}/{@code order_lines} tables, and publish the event that says so. + * It never touches a read-model table and has no idea one exists. + * + *

The events are published through plain {@link ApplicationEventPublisher}, not + * Spring Modulith's event publication registry that the companion {@code outbox} + * module (post #33) + * builds around. That is a deliberate scope line: this module is about the shape of + * CQRS itself — two stores, one path in, one path out, a projection connecting them — + * not about guaranteeing that projection against a crash between commit and delivery. + * Put them together and the registry is what plugs the hole this module's events + * leave open; see the "going deeper" link at the end of the first section. + */ +@Service +public class OrderCommandService { + + private final OrderRepository orderRepository; + private final ApplicationEventPublisher events; + + public OrderCommandService(OrderRepository orderRepository, ApplicationEventPublisher events) { + this.orderRepository = orderRepository; + this.events = events; + } + + @Transactional + public String placeOrder(String customerName, List lines) { + String orderId = UUID.randomUUID().toString(); + Order order = new Order(orderId, customerName, lines); + orderRepository.save(order); + events.publishEvent(new OrderPlaced(orderId, customerName, lines, order.totalCents())); + return orderId; + } + + @Transactional + public void shipOrder(String orderId) { + Order order = orderRepository.findById(orderId) + .orElseThrow(() -> new NoSuchElementException("no such order: " + orderId)); + order.ship(); + orderRepository.save(order); + events.publishEvent(new OrderShipped(orderId)); + } + + @Transactional + public void cancelOrder(String orderId) { + Order order = orderRepository.findById(orderId) + .orElseThrow(() -> new NoSuchElementException("no such order: " + orderId)); + order.cancel(); + orderRepository.save(order); + events.publishEvent(new OrderCancelled(orderId)); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderPlaced.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderPlaced.java new file mode 100644 index 0000000..f7f50fa --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderPlaced.java @@ -0,0 +1,14 @@ +package com.ankurm.cqrsdemo.command; + +import com.ankurm.cqrsdemo.command.internal.OrderLine; + +import java.util.List; + +/** + * Published after the write-side transaction that created the order commits. + * This is the only thing the read side ever learns about an order being placed — + * it carries everything {@code OrderSummaryProjection} needs so the projection + * never has to call back into the write side to fill in a blank. + */ +public record OrderPlaced(String orderId, String customerName, List lines, long totalCents) { +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderShipped.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderShipped.java new file mode 100644 index 0000000..901ec5d --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/OrderShipped.java @@ -0,0 +1,5 @@ +package com.ankurm.cqrsdemo.command; + +/** Published after a {@code ship()} command commits. */ +public record OrderShipped(String orderId) { +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/Order.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/Order.java new file mode 100644 index 0000000..0d69923 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/Order.java @@ -0,0 +1,81 @@ +package com.ankurm.cqrsdemo.command.internal; + +import jakarta.persistence.CollectionTable; +import jakarta.persistence.ElementCollection; +import jakarta.persistence.Entity; +import jakarta.persistence.EnumType; +import jakarta.persistence.Enumerated; +import jakarta.persistence.GeneratedValue; +import jakarta.persistence.GenerationType; +import jakarta.persistence.Id; +import jakarta.persistence.JoinColumn; +import jakarta.persistence.Table; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * The write-side aggregate. Normalized across two tables ({@code orders} and + * {@code orders_lines}), because the write side's only customer is the command + * service — nobody queries this shape directly. {@link #status} is an enum stored + * as a string so a glance at the {@code orders} table tells you what it means + * without a lookup table. + */ +@Entity +@Table(name = "orders") +public class Order { + + public enum Status { PLACED, SHIPPED, CANCELLED } + + @Id + private String id; + + private String customerName; + + @Enumerated(EnumType.STRING) + private Status status; + + @ElementCollection + @CollectionTable(name = "order_lines", joinColumns = @JoinColumn(name = "order_id")) + private List lines = new ArrayList<>(); + + protected Order() { + // JPA + } + + public Order(String id, String customerName, List lines) { + this.id = id; + this.customerName = customerName; + this.lines = new ArrayList<>(lines); + this.status = Status.PLACED; + } + + public void ship() { + this.status = Status.SHIPPED; + } + + public void cancel() { + this.status = Status.CANCELLED; + } + + public String id() { + return id; + } + + public String customerName() { + return customerName; + } + + public Status status() { + return status; + } + + public List lines() { + return Collections.unmodifiableList(lines); + } + + public long totalCents() { + return lines.stream().mapToLong(OrderLine::lineTotalCents).sum(); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderLine.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderLine.java new file mode 100644 index 0000000..7d51e35 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderLine.java @@ -0,0 +1,43 @@ +package com.ankurm.cqrsdemo.command.internal; + +import jakarta.persistence.Embeddable; + +/** + * A single line of a write-side {@link Order}. {@code @Embeddable} rather than its own + * entity+table on purpose — the write side's job is to be correct and normalized for + * the one thing it does (accept commands), not to be convenient to query. That's the + * read side's job. + */ +@Embeddable +public class OrderLine { + + private String sku; + private int quantity; + private long unitPriceCents; + + protected OrderLine() { + // JPA + } + + public OrderLine(String sku, int quantity, long unitPriceCents) { + this.sku = sku; + this.quantity = quantity; + this.unitPriceCents = unitPriceCents; + } + + public String sku() { + return sku; + } + + public int quantity() { + return quantity; + } + + public long unitPriceCents() { + return unitPriceCents; + } + + public long lineTotalCents() { + return (long) quantity * unitPriceCents; + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderRepository.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderRepository.java new file mode 100644 index 0000000..6c0e6e1 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/command/internal/OrderRepository.java @@ -0,0 +1,13 @@ +package com.ankurm.cqrsdemo.command.internal; + +import org.springframework.data.jpa.repository.JpaRepository; + +/** + * The write side's only repository. It is package-private in spirit even though Java + * cannot enforce that across a single module the way Spring Modulith's {@code internal} + * package convention does in the sibling {@code order-fulfillment} module — see + * {@link com.ankurm.cqrsdemo.query.OrderQueryController} for the test that checks the + * query side never gets a reference to this interface. + */ +public interface OrderRepository extends JpaRepository { +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/config/DataSourceConfig.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/config/DataSourceConfig.java new file mode 100644 index 0000000..698abaf --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/config/DataSourceConfig.java @@ -0,0 +1,49 @@ +package com.ankurm.cqrsdemo.config; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.jdbc.autoconfigure.DataSourceProperties; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; + +import javax.sql.DataSource; + +/** + * Two independent {@link DataSource} beans from two independent {@code spring.datasource.*} + * trees — {@code write} and {@code read} — rather than Boot's single autoconfigured + * {@code DataSource}. {@code writeDataSource} is {@code @Primary} so Spring Data JPA's + * autoconfiguration (which asks for "the" {@code DataSource} by type) wires to it without + * any further configuration; {@code readDataSource} is reachable only by the explicit + * {@code @Qualifier("readDataSource")} used in {@link com.ankurm.cqrsdemo.read.ReadSideConfig} + * and {@link com.ankurm.cqrsdemo.query.OrderQueryController}. There is no JPA + * {@code EntityManagerFactory} pointed at {@code readDataSource} at all — the read side + * physically has no Hibernate session to accidentally reuse. + */ +@Configuration +public class DataSourceConfig { + + @Bean + @Primary + @ConfigurationProperties("spring.datasource.write") + public DataSourceProperties writeDataSourceProperties() { + return new DataSourceProperties(); + } + + @Bean + @Primary + public DataSource writeDataSource(@Qualifier("writeDataSourceProperties") DataSourceProperties properties) { + return properties.initializeDataSourceBuilder().build(); + } + + @Bean + @ConfigurationProperties("spring.datasource.read") + public DataSourceProperties readDataSourceProperties() { + return new DataSourceProperties(); + } + + @Bean + public DataSource readDataSource(@Qualifier("readDataSourceProperties") DataSourceProperties properties) { + return properties.initializeDataSourceBuilder().build(); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderQueryController.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderQueryController.java new file mode 100644 index 0000000..b746856 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderQueryController.java @@ -0,0 +1,51 @@ +package com.ankurm.cqrsdemo.query; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.http.ResponseEntity; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.RestController; + +import java.util.List; + +/** + * Every method here is a {@code SELECT} against {@code order_summary} through the + * read-side {@link JdbcTemplate} -- nothing in this class can reach {@code OrderRepository} + * or an {@code EntityManager} even by accident, because it is never given one. The + * constructor below is the whole proof; the test for it is + * {@code OrderQueryControllerHasNoWriteSideAccessTest}, which inspects this exact + * constructor through reflection rather than trusting this comment. + */ +@RestController +public class OrderQueryController { + + private final JdbcTemplate readJdbcTemplate; + + public OrderQueryController(@Qualifier("readJdbcTemplate") JdbcTemplate readJdbcTemplate) { + this.readJdbcTemplate = readJdbcTemplate; + } + + @GetMapping("/order-summaries") + public List listAll() { + return readJdbcTemplate.query("SELECT * FROM order_summary ORDER BY updated_at DESC", this::toSummary); + } + + @GetMapping("/order-summaries/{orderId}") + public ResponseEntity findOne(@PathVariable String orderId) { + List rows = readJdbcTemplate.query( + "SELECT * FROM order_summary WHERE order_id = ?", + this::toSummary, orderId); + return rows.isEmpty() ? ResponseEntity.notFound().build() : ResponseEntity.ok(rows.get(0)); + } + + private OrderSummary toSummary(java.sql.ResultSet rs, int rowNum) throws java.sql.SQLException { + return new OrderSummary( + rs.getString("order_id"), + rs.getString("customer_name"), + rs.getInt("item_count"), + rs.getLong("total_cents"), + rs.getString("status"), + rs.getTimestamp("updated_at").toInstant()); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderSummary.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderSummary.java new file mode 100644 index 0000000..4663928 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/query/OrderSummary.java @@ -0,0 +1,8 @@ +package com.ankurm.cqrsdemo.query; + +import java.time.Instant; + +/** The read model's own shape -- denormalized, display-ready, nothing like the write side's {@code Order}. */ +public record OrderSummary(String orderId, String customerName, int itemCount, long totalCents, + String status, Instant updatedAt) { +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/read/AsyncOrderSummaryProjection.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/AsyncOrderSummaryProjection.java new file mode 100644 index 0000000..e405ef9 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/AsyncOrderSummaryProjection.java @@ -0,0 +1,48 @@ +package com.ankurm.cqrsdemo.read; + +import com.ankurm.cqrsdemo.command.OrderCancelled; +import com.ankurm.cqrsdemo.command.OrderPlaced; +import com.ankurm.cqrsdemo.command.OrderShipped; +import org.springframework.context.annotation.Profile; +import org.springframework.scheduling.annotation.Async; +import org.springframework.stereotype.Component; +import org.springframework.transaction.event.TransactionPhase; +import org.springframework.transaction.event.TransactionalEventListener; + +/** + * The {@code async-projection} profile. {@code @Async} on a {@code @TransactionalEventListener} + * still waits for the write-side transaction to commit, then hands the actual call off to the + * {@code @EnableAsync} executor instead of running it inline -- so the HTTP response for the + * command returns as soon as the commit does, without waiting for the read-model write. + * The read model is then consistent a little while later, not immediately. See + * {@code output/01-async-profile-staleness-window.txt} for a query that lands inside that + * window and one that lands after it. + */ +@Component +@Profile("async-projection") +class AsyncOrderSummaryProjection { + + private final OrderSummaryWriter writer; + + AsyncOrderSummaryProjection(OrderSummaryWriter writer) { + this.writer = writer; + } + + @Async + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + void on(OrderPlaced event) { + writer.applyOrderPlaced(event); + } + + @Async + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + void on(OrderShipped event) { + writer.applyStatusChange(event.orderId(), "SHIPPED"); + } + + @Async + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + void on(OrderCancelled event) { + writer.applyStatusChange(event.orderId(), "CANCELLED"); + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/read/OrderSummaryWriter.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/OrderSummaryWriter.java new file mode 100644 index 0000000..95c9a6d --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/OrderSummaryWriter.java @@ -0,0 +1,56 @@ +package com.ankurm.cqrsdemo.read; + +import com.ankurm.cqrsdemo.command.OrderPlaced; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Component; + +import java.sql.Timestamp; +import java.time.Instant; + +/** + * The actual read-model write, shared by the synchronous and {@code @Async} projection + * listeners so the two profiles differ only in how they are called, never in what gets + * written. {@link #simulatedWriteDelayMs} stands in for whatever a real read-model write + * costs -- a second database round trip, a Redis call, a search index update -- so the + * trade-off the two profiles demonstrate is visible on a stopwatch, not just asserted. + */ +@Component +class OrderSummaryWriter { + + private final JdbcTemplate readJdbcTemplate; + private final long simulatedWriteDelayMs; + + OrderSummaryWriter(JdbcTemplate readJdbcTemplate, + @Value("${cqrs.projection.simulated-write-delay-ms}") long simulatedWriteDelayMs) { + this.readJdbcTemplate = readJdbcTemplate; + this.simulatedWriteDelayMs = simulatedWriteDelayMs; + } + + void applyOrderPlaced(OrderPlaced event) { + sleep(); + readJdbcTemplate.update(""" + MERGE INTO order_summary (order_id, customer_name, item_count, total_cents, status, updated_at) + KEY (order_id) + VALUES (?, ?, ?, ?, 'PLACED', ?) + """, + event.orderId(), event.customerName(), event.lines().size(), event.totalCents(), + Timestamp.from(Instant.now())); + } + + void applyStatusChange(String orderId, String status) { + sleep(); + readJdbcTemplate.update( + "UPDATE order_summary SET status = ?, updated_at = ? WHERE order_id = ?", + status, Timestamp.from(Instant.now()), orderId); + } + + private void sleep() { + try { + Thread.sleep(simulatedWriteDelayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/read/ReadSideConfig.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/ReadSideConfig.java new file mode 100644 index 0000000..7f682ea --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/ReadSideConfig.java @@ -0,0 +1,37 @@ +package com.ankurm.cqrsdemo.read; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.jdbc.core.JdbcTemplate; + +import javax.sql.DataSource; + +/** + * Builds the one {@link JdbcTemplate} the read side is allowed to use, against the + * {@code readDataSource} bean from {@link com.ankurm.cqrsdemo.config.DataSourceConfig} -- + * never against the write side's {@code DataSource}, and creates the {@code order_summary} + * table it owns on startup. Boot's own {@code JdbcTemplateAutoConfiguration} would hand out + * a {@code JdbcTemplate} for the {@code @Primary} (write) {@code DataSource} if anyone asked + * for a plain, unqualified one -- which is exactly why + * {@link com.ankurm.cqrsdemo.query.OrderQueryController} always asks for this bean by name. + */ +@Configuration +public class ReadSideConfig { + + @Bean + public JdbcTemplate readJdbcTemplate(@Qualifier("readDataSource") DataSource readDataSource) { + JdbcTemplate jdbcTemplate = new JdbcTemplate(readDataSource); + jdbcTemplate.execute(""" + CREATE TABLE IF NOT EXISTS order_summary ( + order_id VARCHAR(64) PRIMARY KEY, + customer_name VARCHAR(200) NOT NULL, + item_count INT NOT NULL, + total_cents BIGINT NOT NULL, + status VARCHAR(20) NOT NULL, + updated_at TIMESTAMP NOT NULL + ) + """); + return jdbcTemplate; + } +} diff --git a/cqrs/src/main/java/com/ankurm/cqrsdemo/read/SyncOrderSummaryProjection.java b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/SyncOrderSummaryProjection.java new file mode 100644 index 0000000..2ccbbe8 --- /dev/null +++ b/cqrs/src/main/java/com/ankurm/cqrsdemo/read/SyncOrderSummaryProjection.java @@ -0,0 +1,43 @@ +package com.ankurm.cqrsdemo.read; + +import com.ankurm.cqrsdemo.command.OrderCancelled; +import com.ankurm.cqrsdemo.command.OrderPlaced; +import com.ankurm.cqrsdemo.command.OrderShipped; +import org.springframework.context.annotation.Profile; +import org.springframework.stereotype.Component; +import org.springframework.transaction.event.TransactionPhase; +import org.springframework.transaction.event.TransactionalEventListener; + +/** + * The default profile. {@code @TransactionalEventListener} with no {@code @Async} runs on + * the same thread that is about to return the HTTP response, after the write-side + * transaction has committed. The read model is guaranteed consistent by the time the + * command's HTTP response goes out -- at the cost of the command's latency including + * the read-model write. See {@code output/00-sync-profile-latency.txt} for what that + * actually measures to. + */ +@Component +@Profile("!async-projection") +class SyncOrderSummaryProjection { + + private final OrderSummaryWriter writer; + + SyncOrderSummaryProjection(OrderSummaryWriter writer) { + this.writer = writer; + } + + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + void on(OrderPlaced event) { + writer.applyOrderPlaced(event); + } + + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + void on(OrderShipped event) { + writer.applyStatusChange(event.orderId(), "SHIPPED"); + } + + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + void on(OrderCancelled event) { + writer.applyStatusChange(event.orderId(), "CANCELLED"); + } +} diff --git a/cqrs/src/main/resources/application.properties b/cqrs/src/main/resources/application.properties new file mode 100644 index 0000000..47bc637 --- /dev/null +++ b/cqrs/src/main/resources/application.properties @@ -0,0 +1,23 @@ +spring.application.name=cqrs + +# The write side: a normal Spring Data JPA setup. @Primary on this bean (see +# DataSourceConfig) is what lets JPA autoconfiguration find it with no further wiring. +spring.datasource.write.url=jdbc:h2:mem:cqrs-write;DB_CLOSE_DELAY=-1 +spring.datasource.write.driver-class-name=org.h2.Driver +spring.datasource.write.username=sa +spring.datasource.write.password= + +spring.jpa.hibernate.ddl-auto=update +spring.jpa.open-in-view=false + +# The read side: a second, independent database. Reached only through the +# readJdbcTemplate bean in ReadSideConfig -- there is no EntityManagerFactory +# pointed at this one. +spring.datasource.read.url=jdbc:h2:mem:cqrs-read;DB_CLOSE_DELAY=-1 +spring.datasource.read.driver-class-name=org.h2.Driver +spring.datasource.read.username=sa +spring.datasource.read.password= + +# How long the projection pretends the read-model write takes. Large enough to +# measure reliably in a test, small enough not to make the test suite slow. +cqrs.projection.simulated-write-delay-ms=300 diff --git a/cqrs/src/test/java/com/ankurm/cqrsdemo/AsyncProfileStalenessWindowTest.java b/cqrs/src/test/java/com/ankurm/cqrsdemo/AsyncProfileStalenessWindowTest.java new file mode 100644 index 0000000..cd86d13 --- /dev/null +++ b/cqrs/src/test/java/com/ankurm/cqrsdemo/AsyncProfileStalenessWindowTest.java @@ -0,0 +1,76 @@ +package com.ankurm.cqrsdemo; + +import com.ankurm.cqrsdemo.command.OrderCommandController.LineRequest; +import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderRequest; +import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderResponse; +import com.ankurm.cqrsdemo.query.OrderSummary; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.resttestclient.TestRestTemplate; +import org.springframework.boot.resttestclient.autoconfigure.AutoConfigureTestRestTemplate; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.test.context.ActiveProfiles; + +import java.time.Duration; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +/** + * {@code async-projection} profile: {@link com.ankurm.cqrsdemo.read.AsyncOrderSummaryProjection} + * is active instead. {@code POST /orders} now returns as soon as the write-side transaction + * commits, without waiting for the simulated 300ms read-model write -- but that means a query + * that lands inside that window sees the order not yet in the read model at all. This test + * deliberately queries in both places: immediately (inside the window) and after waiting for + * it to close (outside the window), and both results are the real point of the test, not the + * final consistent one alone. + */ +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) +@AutoConfigureTestRestTemplate +@ActiveProfiles("async-projection") +class AsyncProfileStalenessWindowTest { + + @Autowired + private TestRestTemplate rest; + + @Test + void placingAnOrderReturnsFastAndTheReadModelIsMomentarilyStale() { + // Warm-up request, deliberately not measured or asserted on: the very first HTTP + // request against a freshly started test context pays for connection-pool startup + // and Hibernate/JPA first-use costs that have nothing to do with this profile's + // async projection. Measuring the first request directly made the baseline ~370ms + // before the projection was even involved -- comfortably past the 300ms this test + // is trying to prove the command path *doesn't* pay. Everything after this call + // runs against an already-warm context. + rest.postForEntity("/orders", new PlaceOrderRequest("Warmup", List.of(new LineRequest("WARMUP", 1, 1))), + PlaceOrderResponse.class); + + PlaceOrderRequest request = new PlaceOrderRequest("Dev", + List.of(new LineRequest("MECH-KEYBOARD", 1, 8999))); + + long start = System.currentTimeMillis(); + ResponseEntity placed = rest.postForEntity("/orders", request, PlaceOrderResponse.class); + long elapsedMs = System.currentTimeMillis() - start; + + String orderId = placed.getBody().orderId(); + System.out.println("POST /orders (async profile) took " + elapsedMs + " ms, orderId=" + orderId); + assertThat(elapsedMs).isLessThan(300L); + + ResponseEntity immediately = rest.getForEntity("/order-summaries/" + orderId, OrderSummary.class); + System.out.println("Immediately after POST returned (inside the staleness window), " + + "GET /order-summaries/" + orderId + " -> status=" + immediately.getStatusCode() + + ", body=" + immediately.getBody()); + assertThat(immediately.getStatusCode()).isEqualTo(HttpStatus.NOT_FOUND); + + await().atMost(Duration.ofSeconds(2)).pollInterval(Duration.ofMillis(50)).untilAsserted(() -> { + ResponseEntity eventually = rest.getForEntity("/order-summaries/" + orderId, OrderSummary.class); + assertThat(eventually.getStatusCode()).isEqualTo(HttpStatus.OK); + System.out.println("After waiting for the projection (outside the staleness window), " + + "GET /order-summaries/" + orderId + " -> status=" + eventually.getStatusCode() + + ", body=" + eventually.getBody()); + }); + } +} diff --git a/cqrs/src/test/java/com/ankurm/cqrsdemo/QuerySideArchitectureTest.java b/cqrs/src/test/java/com/ankurm/cqrsdemo/QuerySideArchitectureTest.java new file mode 100644 index 0000000..fe30562 --- /dev/null +++ b/cqrs/src/test/java/com/ankurm/cqrsdemo/QuerySideArchitectureTest.java @@ -0,0 +1,43 @@ +package com.ankurm.cqrsdemo; + +import com.ankurm.cqrsdemo.query.OrderQueryController; +import org.junit.jupiter.api.Test; +import org.springframework.jdbc.core.JdbcTemplate; + +import java.lang.reflect.Constructor; +import java.lang.reflect.Field; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * The query side's version of the hexagonal-architecture post's + * {@code mvn dependency:tree}: instead of trusting the Javadoc comment on + * {@link OrderQueryController} that says it cannot reach the write side, this inspects + * the actual compiled class through reflection. No Spring context needed -- this is a + * claim about the class file, not about runtime wiring. + */ +class QuerySideArchitectureTest { + + @Test + void queryControllersOnlyDependencyIsTheReadJdbcTemplate() { + Constructor[] constructors = OrderQueryController.class.getDeclaredConstructors(); + assertThat(constructors).hasSize(1); + + Class[] paramTypes = constructors[0].getParameterTypes(); + System.out.println("OrderQueryController constructor parameter types: " + + java.util.Arrays.toString(paramTypes)); + assertThat(paramTypes).containsExactly(JdbcTemplate.class); + } + + @Test + void queryControllerDeclaresNoFieldFromTheWriteSidePackage() { + Field[] fields = OrderQueryController.class.getDeclaredFields(); + for (Field field : fields) { + String packageName = field.getType().getPackageName(); + System.out.println("field " + field.getName() + " : " + field.getType().getName()); + assertThat(packageName) + .as("field %s should not be able to reference the write side", field.getName()) + .doesNotStartWith("com.ankurm.cqrsdemo.command"); + } + } +} diff --git a/cqrs/src/test/java/com/ankurm/cqrsdemo/SyncProfileLatencyTest.java b/cqrs/src/test/java/com/ankurm/cqrsdemo/SyncProfileLatencyTest.java new file mode 100644 index 0000000..a6ee4f4 --- /dev/null +++ b/cqrs/src/test/java/com/ankurm/cqrsdemo/SyncProfileLatencyTest.java @@ -0,0 +1,51 @@ +package com.ankurm.cqrsdemo; + +import com.ankurm.cqrsdemo.command.OrderCommandController.LineRequest; +import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderRequest; +import com.ankurm.cqrsdemo.command.OrderCommandController.PlaceOrderResponse; +import com.ankurm.cqrsdemo.query.OrderSummary; +import org.junit.jupiter.api.Test; +import org.springframework.boot.resttestclient.autoconfigure.AutoConfigureTestRestTemplate; +import org.springframework.boot.resttestclient.TestRestTemplate; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.http.ResponseEntity; + +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Default profile: no {@code async-projection}, so {@link com.ankurm.cqrsdemo.read.SyncOrderSummaryProjection} + * is active. The projection's simulated write delay (300ms, see {@code application.properties}) + * runs on the same thread that is about to answer the command's HTTP request, so it shows up + * directly in how long {@code POST /orders} takes -- and because of that, the read model is + * already correct the instant the command returns, with no waiting required. + */ +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) +@AutoConfigureTestRestTemplate +class SyncProfileLatencyTest { + + @org.springframework.beans.factory.annotation.Autowired + private TestRestTemplate rest; + + @Test + void placingAnOrderTakesAtLeastTheSimulatedProjectionDelayAndTheReadModelIsImmediatelyConsistent() { + PlaceOrderRequest request = new PlaceOrderRequest("Priya", + List.of(new LineRequest("USB-C-CABLE", 2, 1299), new LineRequest("USB-C-HUB", 1, 4999))); + + long start = System.currentTimeMillis(); + ResponseEntity placed = rest.postForEntity("/orders", request, PlaceOrderResponse.class); + long elapsedMs = System.currentTimeMillis() - start; + + String orderId = placed.getBody().orderId(); + System.out.println("POST /orders (sync profile) took " + elapsedMs + " ms, orderId=" + orderId); + assertThat(elapsedMs).isGreaterThanOrEqualTo(300L); + + ResponseEntity summary = rest.getForEntity("/order-summaries/" + orderId, OrderSummary.class); + System.out.println("Immediately after POST returned, GET /order-summaries/" + orderId + + " -> status=" + summary.getStatusCode() + ", body=" + summary.getBody()); + assertThat(summary.getBody()).isNotNull(); + assertThat(summary.getBody().totalCents()).isEqualTo(2 * 1299 + 4999); + assertThat(summary.getBody().status()).isEqualTo("PLACED"); + } +} diff --git a/pom.xml b/pom.xml index 9d32aad..775893b 100644 --- a/pom.xml +++ b/pom.xml @@ -15,6 +15,7 @@ Each Maven module backs one post; modules are added over time, one commit each. - order-fulfillment : Spring Modulith 2.1: Enforcing Module Boundaries Inside a Spring Boot Monolith - outbox : Transactional Outbox with the Spring Modulith Event Publication Registry + - cqrs : CQRS in Spring Boot Without a Framework: Separate Read Models and Projections @@ -27,5 +28,6 @@ order-fulfillment outbox + cqrs