From d815a37f2ea0cf0a59717b06c5e2fcc3c67d5639 Mon Sep 17 00:00:00 2001 From: asmhatre Date: Sat, 3 Oct 2026 19:27:32 +0000 Subject: [PATCH] Add outbox module: Transactional Outbox with the Spring Modulith Event Publication Registry --- outbox/README.md | 54 +++++++++++ outbox/output/00-dual-write-problem.txt | 35 +++++++ outbox/output/01-registry-safety-net.txt | 20 ++++ .../02-replay-incomplete-publications.txt | 23 +++++ outbox/pom.xml | 94 +++++++++++++++++++ .../com/ankurm/outboxdemo/Application.java | 13 +++ .../outboxdemo/billing/BillingManagement.java | 35 +++++++ .../outboxdemo/billing/InvoiceIssued.java | 16 ++++ .../billing/NaiveBillingService.java | 62 ++++++++++++ .../outboxdemo/billing/internal/Invoice.java | 50 ++++++++++ .../billing/internal/InvoiceRepository.java | 6 ++ .../src/main/resources/application.properties | 19 ++++ .../billing/DualWriteProblemTest.java | 80 ++++++++++++++++ .../billing/FlakyAuditListener.java | 31 ++++++ .../billing/RegistrySafetyNetTest.java | 91 ++++++++++++++++++ .../ReplayIncompletePublicationsTest.java | 90 ++++++++++++++++++ pom.xml | 2 + 17 files changed, 721 insertions(+) create mode 100644 outbox/README.md create mode 100644 outbox/output/00-dual-write-problem.txt create mode 100644 outbox/output/01-registry-safety-net.txt create mode 100644 outbox/output/02-replay-incomplete-publications.txt create mode 100644 outbox/pom.xml create mode 100644 outbox/src/main/java/com/ankurm/outboxdemo/Application.java create mode 100644 outbox/src/main/java/com/ankurm/outboxdemo/billing/BillingManagement.java create mode 100644 outbox/src/main/java/com/ankurm/outboxdemo/billing/InvoiceIssued.java create mode 100644 outbox/src/main/java/com/ankurm/outboxdemo/billing/NaiveBillingService.java create mode 100644 outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/Invoice.java create mode 100644 outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/InvoiceRepository.java create mode 100644 outbox/src/main/resources/application.properties create mode 100644 outbox/src/test/java/com/ankurm/outboxdemo/billing/DualWriteProblemTest.java create mode 100644 outbox/src/test/java/com/ankurm/outboxdemo/billing/FlakyAuditListener.java create mode 100644 outbox/src/test/java/com/ankurm/outboxdemo/billing/RegistrySafetyNetTest.java create mode 100644 outbox/src/test/java/com/ankurm/outboxdemo/billing/ReplayIncompletePublicationsTest.java diff --git a/outbox/README.md b/outbox/README.md new file mode 100644 index 0000000..0045d1b --- /dev/null +++ b/outbox/README.md @@ -0,0 +1,54 @@ +# outbox + +Companion code for *[Transactional Outbox with the Spring Modulith Event Publication +Registry](https://ankurm.com)* on [ankurm.com](https://ankurm.com). + +A single `billing` module demonstrating the dual-write problem, Spring Modulith's event +publication registry acting as a real transactional outbox, externalizing events to a +real (embedded) Kafka broker via `spring-modulith-events-kafka`, and replaying an +incomplete publication with `IncompleteEventPublications`. Every number in the +companion post traces to `output/`, in order. + +## Versions + +| Artifact | Version | +|---|---| +| Spring Boot | 4.1.1 | +| Spring Modulith | 2.1.1 (GA; verified against Maven Central's `maven-metadata.xml`) | +| spring-modulith-events-kafka | 2.1.1 | +| spring-kafka / spring-kafka-test | 4.1.1 (managed by the Spring Boot 4.1.1 BOM; `kafka-clients` 4.2.1) | +| JDK | 25 (LTS) | +| H2 | runtime, in-memory, for the demo only | + +## Quickstart + +``` +mvn -pl outbox -am test -Dtest=DualWriteProblemTest +mvn -pl outbox -am test -Dtest=RegistrySafetyNetTest +mvn -pl outbox -am test -Dtest=ReplayIncompletePublicationsTest +``` + +Run separately on purpose -- `RegistrySafetyNetTest` starts a real embedded Kafka +broker (`@EmbeddedKafka`), which takes several seconds and is unrelated to what the +other two tests are checking. + +## What each class is + +| Class | Role | +|---|---| +| `billing.BillingManagement` | The outbox-safe way to issue an invoice: one local transaction, one `Invoice` row, one event publication registry row | +| `billing.NaiveBillingService` | The "before" picture -- a separate, uncoordinated `kafka.send()` call after the database write already committed. Never wired as the module's real API; only ever called directly from `DualWriteProblemTest` | +| `billing.InvoiceIssued` | `@Externalized("invoices::#{#this.invoiceId}")` -- routes to the `invoices` Kafka topic, keyed by the invoice id | +| `billing.FlakyAuditListener` (test-only) | Throws on its first invocation, succeeds after -- stands in for any downstream dependency being briefly unavailable, without needing to actually break a running broker mid-test | + +## Captured output (`output/`) + +| File | What it captures | +|---|---| +| `00-dual-write-problem.txt` | The invoice committed, Kafka never received either message, and a real stack trace showing `KafkaTemplate.send()` throwing *synchronously* when the broker is totally unreachable -- narrower than, and a correction to, the usual "fire-and-forget send() fails silently" claim | +| `01-registry-safety-net.txt` | The real production path: a real embedded Kafka broker receiving the real message, and the registry moving both listeners' rows into `EVENT_PUBLICATION_ARCHIVE` (`completion-mode=ARCHIVE` is turned on for this module -- the previous post in this series left it commented out) | +| `02-replay-incomplete-publications.txt` | A listener failing once, the registry leaving its row incomplete, and `IncompleteEventPublications.resubmitIncompletePublicationsOlderThan(Duration)` -- a real method, verified against the compiled `events-api` jar -- redelivering it successfully | + +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/outbox/output/00-dual-write-problem.txt b/outbox/output/00-dual-write-problem.txt new file mode 100644 index 0000000..9061c43 --- /dev/null +++ b/outbox/output/00-dual-write-problem.txt @@ -0,0 +1,35 @@ +$ mvn -pl outbox -am test -Dtest=DualWriteProblemTest +(NaiveBillingService: save the Invoice, then separately call kafka.send() -- no + coordinator between the two. spring.kafka.bootstrap-servers points at localhost:19999, + nothing listening there) + +2026-10-04T00:52:56.429+05:30 ERROR 5271 --- [outbox] [ main] o.s.k.support.LoggingProducerListener : Exception thrown when sending a message with key='1' and payload='WIDGET-1:1999' to topic invoices-naive: + +org.apache.kafka.common.errors.TimeoutException: Topic invoices-naive not present in metadata after 2000 ms. + +Exception thrown synchronously from issueInvoiceTheNaiveWay(): org.springframework.kafka.KafkaException: Send failed +Invoice count regardless: 1 (was 0). The database write already committed inside saveInvoiceOnly() before this exception was ever thrown -- the exception came from the Kafka call, several lines later, on an already-committed row. + +2026-10-04T00:52:58.495+05:30 ERROR 5271 --- [outbox] [ main] o.s.k.support.LoggingProducerListener : Exception thrown when sending a message with key='2' and payload='WIDGET-1:2999' to topic invoices-naive: + +org.apache.kafka.common.errors.TimeoutException: Topic invoices-naive not present in metadata after 2000 ms. + +Invoice count after the blocked, failed send: 2 (was 1). The row committed before the send was even attempted. + +[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0 + +Two real findings here, not one. First, the obvious one: the invoice exists in the +database in both tests, and Kafka never received either message -- that's the dual-write +problem, reproduced with nothing faked beyond pointing Kafka at a port with no listener. + +Second, a more specific one that contradicts the usual "fire-and-forget send() fails +silently" framing: when the broker can't be reached to fetch topic metadata at all, +Spring's KafkaTemplate.send() throws synchronously, wrapping the Kafka client's own +TimeoutException in a KafkaException, even though the method's declared return type is +a CompletableFuture. The "silent" version of this bug is real, but it is a narrower +case than total broker unavailability -- it's what happens when the broker is reachable +and the topic's metadata is known, but the specific produce request times out waiting +for an acknowledgment afterward. That failure mode surfaces only through the returned +future or a callback, never as a synchronous exception, and is not reproduced in this +transcript (it needs a broker that accepts the connection and then stops responding +mid-request, not one that was never there to begin with). diff --git a/outbox/output/01-registry-safety-net.txt b/outbox/output/01-registry-safety-net.txt new file mode 100644 index 0000000..715123c --- /dev/null +++ b/outbox/output/01-registry-safety-net.txt @@ -0,0 +1,20 @@ +$ mvn -pl outbox -am test -Dtest=RegistrySafetyNetTest +(real embedded Kafka broker via @EmbeddedKafka; BillingManagement.issueInvoice() + publishes InvoiceIssued, @Externalized("invoices::#{#this.invoiceId}") routes it + through spring-modulith-events-kafka) + +Kafka record received -- topic=invoices key=1 value={"invoiceId":1,"sku":"GADGET-7","amountCents":4500} +EVENT_PUBLICATION_ARCHIVE rows for InvoiceIssued: 2 +EVENT_PUBLICATION rows still outstanding for InvoiceIssued: 0 + +[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0 + +The message really arrived on the real topic, with the real routing key (the invoice's +own id, exactly as the SpEL expression after "::" in @Externalized specifies) and the +real JSON payload. Two archive rows for InvoiceIssued, not one, because this module's +test sources also declare FlakyAuditListener (see output/02) -- it is a second, +independent listener on the same event, and with completion-mode=ARCHIVE configured in +application.properties, both listeners' completed publications moved out of +EVENT_PUBLICATION and into EVENT_PUBLICATION_ARCHIVE once they succeeded. Zero rows +left outstanding in EVENT_PUBLICATION for this event type confirms neither listener is +still pending. diff --git a/outbox/output/02-replay-incomplete-publications.txt b/outbox/output/02-replay-incomplete-publications.txt new file mode 100644 index 0000000..a49f9eb --- /dev/null +++ b/outbox/output/02-replay-incomplete-publications.txt @@ -0,0 +1,23 @@ +$ mvn -pl outbox -am test -Dtest=ReplayIncompletePublicationsTest +(FlakyAuditListener throws on its first invocation of InvoiceIssued, succeeds on every + one after -- standing in for any downstream dependency that is briefly unavailable) + +Incomplete FlakyAuditListener publications after the simulated failure: 1 +Incomplete FlakyAuditListener publications after resubmission: 0 +Invoice 1 never changed -- the retry was purely about redelivering the event, the invoice row was correct the whole time. + +[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0 + +The sequence that actually ran: BillingManagement.issueInvoice() commits the Invoice +row and the EVENT_PUBLICATION row for FlakyAuditListener in one local transaction. +FlakyAuditListener.on(InvoiceIssued) is then invoked asynchronously, throws on purpose, +and the registry leaves that row with COMPLETION_DATE still null -- "incomplete" is +exactly that: a row with no completion date, regardless of whether the listener ever +ran at all or ran and failed. Calling +IncompleteEventPublications.resubmitIncompletePublicationsOlderThan(Duration.ZERO) -- +a real method on a real interface, verified against the compiled events-api 2.1.1 jar, +not assumed from documentation prose -- re-invokes the listener for every such row. +This time failNextAttempt is already false, so it succeeds, and the row's completion +date is set. The invoice itself was never touched by any of this: the retry is about +redelivering a notification, not about redoing a database write that already +succeeded the first time. diff --git a/outbox/pom.xml b/outbox/pom.xml new file mode 100644 index 0000000..3c3771a --- /dev/null +++ b/outbox/pom.xml @@ -0,0 +1,94 @@ + + + 4.0.0 + + + com.ankurm + spring-modulith-demo + 1.0.0 + + + outbox + jar + + outbox + + Companion code for "Transactional Outbox with the Spring Modulith Event Publication + Registry" on ankurm.com. + + + + + + org.springframework.boot + spring-boot-dependencies + ${spring-boot.version} + pom + import + + + org.springframework.modulith + spring-modulith-bom + ${spring-modulith.version} + pom + import + + + + + + + org.springframework.boot + spring-boot-starter + + + org.springframework.modulith + spring-modulith-starter-jdbc + + + org.springframework.modulith + spring-modulith-events-kafka + + + org.springframework.boot + spring-boot-starter-data-jpa + + + org.springframework.kafka + spring-kafka + + + com.h2database + h2 + runtime + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.springframework.modulith + spring-modulith-starter-test + test + + + org.springframework.kafka + spring-kafka-test + test + + + + + + + org.springframework.boot + spring-boot-maven-plugin + ${spring-boot.version} + + + + diff --git a/outbox/src/main/java/com/ankurm/outboxdemo/Application.java b/outbox/src/main/java/com/ankurm/outboxdemo/Application.java new file mode 100644 index 0000000..f46a8c8 --- /dev/null +++ b/outbox/src/main/java/com/ankurm/outboxdemo/Application.java @@ -0,0 +1,13 @@ +package com.ankurm.outboxdemo; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.modulith.Modulithic; + +@SpringBootApplication +@Modulithic(systemName = "Outbox Demo") +public class Application { + public static void main(String[] args) { + SpringApplication.run(Application.class, args); + } +} diff --git a/outbox/src/main/java/com/ankurm/outboxdemo/billing/BillingManagement.java b/outbox/src/main/java/com/ankurm/outboxdemo/billing/BillingManagement.java new file mode 100644 index 0000000..a871911 --- /dev/null +++ b/outbox/src/main/java/com/ankurm/outboxdemo/billing/BillingManagement.java @@ -0,0 +1,35 @@ +package com.ankurm.outboxdemo.billing; + +import com.ankurm.outboxdemo.billing.internal.Invoice; +import com.ankurm.outboxdemo.billing.internal.InvoiceRepository; +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +/** + * The outbox-safe way to issue an invoice. The {@code Invoice} row and the event + * publication registry row for {@link InvoiceIssued} are written in the same local + * database transaction -- that insert is what makes this an outbox, not the event + * itself. Whatever happens to the event afterwards (externalized to Kafka, picked up + * by an in-process listener, or nothing at all because the process dies one instruction + * later) cannot un-happen the fact that a row recording "this needs to go out" already + * made it to disk, atomically with the invoice itself. + */ +@Service +public class BillingManagement { + + private final InvoiceRepository invoices; + private final ApplicationEventPublisher events; + + public BillingManagement(InvoiceRepository invoices, ApplicationEventPublisher events) { + this.invoices = invoices; + this.events = events; + } + + @Transactional + public Invoice issueInvoice(String sku, int amountCents) { + var invoice = invoices.save(new Invoice(sku, amountCents, "ISSUED")); + events.publishEvent(new InvoiceIssued(invoice.getId(), invoice.getSku(), invoice.getAmountCents())); + return invoice; + } +} diff --git a/outbox/src/main/java/com/ankurm/outboxdemo/billing/InvoiceIssued.java b/outbox/src/main/java/com/ankurm/outboxdemo/billing/InvoiceIssued.java new file mode 100644 index 0000000..5047566 --- /dev/null +++ b/outbox/src/main/java/com/ankurm/outboxdemo/billing/InvoiceIssued.java @@ -0,0 +1,16 @@ +package com.ankurm.outboxdemo.billing; + +import org.springframework.modulith.events.Externalized; + +/** + * {@code @Externalized}'s routing syntax is {@code "target::key"} -- everything before + * {@code ::} is the Kafka topic, everything after is a SpEL expression (root object + * {@code #this}) used as the message key. Verified against the compiled annotation + * (org.springframework.modulith.events.Externalized, 2.1.1): it is {@code @Target(TYPE)}, + * one aliased {@code value()/target()} String attribute, nothing else -- the routing + * syntax itself is parsed at runtime by the externalization machinery, not encoded in + * the annotation's own shape. + */ +@Externalized("invoices::#{#this.invoiceId}") +public record InvoiceIssued(Long invoiceId, String sku, int amountCents) { +} diff --git a/outbox/src/main/java/com/ankurm/outboxdemo/billing/NaiveBillingService.java b/outbox/src/main/java/com/ankurm/outboxdemo/billing/NaiveBillingService.java new file mode 100644 index 0000000..b3b46ef --- /dev/null +++ b/outbox/src/main/java/com/ankurm/outboxdemo/billing/NaiveBillingService.java @@ -0,0 +1,62 @@ +package com.ankurm.outboxdemo.billing; + +import com.ankurm.outboxdemo.billing.internal.Invoice; +import com.ankurm.outboxdemo.billing.internal.InvoiceRepository; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +/** + * The two-systems-no-coordinator version, kept here deliberately as the "before" + * picture -- it is never wired as the module's real API, only invoked directly from + * {@code DualWriteProblemTest}. The database write commits and returns before this + * method ever reaches the Kafka call: there is no transaction, local or distributed, + * spanning both. If the process dies, or the broker is unreachable, between those two + * lines, the invoice exists and nothing downstream will ever hear about it. + */ +@Service +public class NaiveBillingService { + + private final InvoiceRepository invoices; + private final KafkaTemplate kafka; + + public NaiveBillingService(InvoiceRepository invoices, KafkaTemplate kafka) { + this.invoices = invoices; + this.kafka = kafka; + } + + @Transactional + public Invoice saveInvoiceOnly(String sku, int amountCents) { + return invoices.save(new Invoice(sku, amountCents, "ISSUED")); + } + + /** + * Deliberately two separate steps: {@link #saveInvoiceOnly} commits and returns + * on its own transaction boundary before this method ever calls {@code kafka.send}. + * Fire-and-forget: the returned {@code CompletableFuture} is never inspected, which + * is the common case at most real call sites. If the broker is unreachable, this + * method returns successfully anyway -- the failure happens on a Kafka client + * thread, after this method has already returned, and nothing here ever sees it. + */ + public Invoice issueInvoiceTheNaiveWay(String sku, int amountCents) { + var invoice = saveInvoiceOnly(sku, amountCents); + kafka.send("invoices-naive", invoice.getId().toString(), + invoice.getSku() + ":" + invoice.getAmountCents()); + return invoice; + } + + /** + * Same two-step shortcut, except this call site does the "responsible" thing and + * blocks on the send future so it can react to a failure. It still doesn't help: + * the database write already committed in {@link #saveInvoiceOnly} by the time this + * line runs, so catching the exception here can log it, but cannot undo the invoice + * that is already sitting in the database with nothing downstream ever told about it. + */ + public Invoice issueInvoiceAndBlockOnTheSend(String sku, int amountCents) throws Exception { + var invoice = saveInvoiceOnly(sku, amountCents); + kafka.send("invoices-naive", invoice.getId().toString(), + invoice.getSku() + ":" + invoice.getAmountCents()) + .get(3, java.util.concurrent.TimeUnit.SECONDS); + return invoice; + } +} diff --git a/outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/Invoice.java b/outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/Invoice.java new file mode 100644 index 0000000..c2186ab --- /dev/null +++ b/outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/Invoice.java @@ -0,0 +1,50 @@ +package com.ankurm.outboxdemo.billing.internal; + +import jakarta.persistence.Entity; +import jakarta.persistence.GeneratedValue; +import jakarta.persistence.GenerationType; +import jakarta.persistence.Id; +import jakarta.persistence.Table; + +/** + * Deliberately named Invoice, not Order -- the previous post in this series + * (output/00-table-name-collision.txt in the order-fulfillment module) already + * spent a transcript on why ORDER is a reserved SQL keyword. This module picks a + * domain noun that doesn't collide, on purpose. + */ +@Entity +@Table(name = "invoices") +public class Invoice { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + private String sku; + private int amountCents; + private String status; + + protected Invoice() { + } + + public Invoice(String sku, int amountCents, String status) { + this.sku = sku; + this.amountCents = amountCents; + this.status = status; + } + + public Long getId() { + return id; + } + + public String getSku() { + return sku; + } + + public int getAmountCents() { + return amountCents; + } + + public String getStatus() { + return status; + } +} diff --git a/outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/InvoiceRepository.java b/outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/InvoiceRepository.java new file mode 100644 index 0000000..8cb6927 --- /dev/null +++ b/outbox/src/main/java/com/ankurm/outboxdemo/billing/internal/InvoiceRepository.java @@ -0,0 +1,6 @@ +package com.ankurm.outboxdemo.billing.internal; + +import org.springframework.data.jpa.repository.JpaRepository; + +public interface InvoiceRepository extends JpaRepository { +} diff --git a/outbox/src/main/resources/application.properties b/outbox/src/main/resources/application.properties new file mode 100644 index 0000000..962a7ba --- /dev/null +++ b/outbox/src/main/resources/application.properties @@ -0,0 +1,19 @@ +spring.application.name=outbox + +spring.datasource.url=jdbc:h2:mem:outbox;DB_CLOSE_DELAY=-1 +spring.datasource.driver-class-name=org.h2.Driver +spring.jpa.hibernate.ddl-auto=update +spring.jpa.open-in-view=false + +# The event publication registry is the outbox table. JDBC-backed here (not JPA) so +# the registry's own schema is independent of whatever ORM the business entities use. +spring.modulith.events.jdbc.schema-initialization.enabled=true + +# ARCHIVE, this time -- the previous post in this series left this commented out so +# every transcript there ran against the library's bare UPDATE default. This module is +# about the outbox actually being operated correctly, so it is turned on here instead. +spring.modulith.events.completion-mode=ARCHIVE +spring.modulith.events.republish-outstanding-events-on-restart=true + +logging.level.com.ankurm.outboxdemo=INFO +logging.level.org.springframework.modulith.events=INFO diff --git a/outbox/src/test/java/com/ankurm/outboxdemo/billing/DualWriteProblemTest.java b/outbox/src/test/java/com/ankurm/outboxdemo/billing/DualWriteProblemTest.java new file mode 100644 index 0000000..37169ea --- /dev/null +++ b/outbox/src/test/java/com/ankurm/outboxdemo/billing/DualWriteProblemTest.java @@ -0,0 +1,80 @@ +package com.ankurm.outboxdemo.billing; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import com.ankurm.outboxdemo.billing.internal.InvoiceRepository; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; + +/** + * No embedded broker here, on purpose -- spring.kafka.bootstrap-servers points at a + * port nothing is listening on, which is the simplest honest way to reproduce "the + * broker is unreachable" without faking anything. Two real, independent systems, + * called one after the other, is already enough to produce the dual-write problem; + * nothing here is simulated beyond pointing Kafka at an address with no listener. + */ +@SpringBootTest +class DualWriteProblemTest { + + @DynamicPropertySource + static void unreachableKafka(DynamicPropertyRegistry registry) { + registry.add("spring.kafka.bootstrap-servers", () -> "localhost:19999"); + registry.add("spring.kafka.producer.properties.max.block.ms", () -> "2000"); + } + + @Autowired + NaiveBillingService naiveBillingService; + + @Autowired + InvoiceRepository invoices; + + @Test + void fireAndForgetSendStillThrowsWhenTheBrokerIsTotallyUnreachable() { + long before = invoices.count(); + + // Expected this call to return cleanly and lose the failure silently -- + // KafkaTemplate.send() returns a CompletableFuture, so a caller that ignores it + // "shouldn't" see an exception. Running it for real shows otherwise: when the + // broker can't even be reached to fetch topic metadata, Spring's KafkaTemplate + // wraps that specific failure and rethrows it synchronously from send() itself + // (see the real stack trace in the captured transcript). The database write is + // already unaffected either way -- it's a separate, already-committed transaction + // by the time this line runs -- but the exception itself is not silent here. + Exception thrown = null; + try { + naiveBillingService.issueInvoiceTheNaiveWay("WIDGET-1", 1999); + } catch (Exception e) { + thrown = e; + } + + assertThat(thrown).isNotNull(); + assertThat(invoices.count()).isEqualTo(before + 1); + System.out.println("Exception thrown synchronously from issueInvoiceTheNaiveWay(): " + + thrown.getClass().getName() + ": " + thrown.getMessage()); + System.out.println("Invoice count regardless: " + invoices.count() + " (was " + before + + "). The database write already committed inside saveInvoiceOnly() before this" + + " exception was ever thrown -- the exception came from the Kafka call, several" + + " lines later, on an already-committed row."); + } + + @Test + void blockingOnTheFutureSurfacesTheSameFailureAndStillCannotUndoTheWrite() { + long before = invoices.count(); + + assertThatThrownBy(() -> naiveBillingService.issueInvoiceAndBlockOnTheSend("WIDGET-1", 2999)) + .isInstanceOf(org.springframework.kafka.KafkaException.class) + .hasCauseInstanceOf(org.apache.kafka.common.errors.TimeoutException.class); + + // The invoice is still there. Blocking on the future told us about the Kafka + // failure a little sooner, but by the time we could react to it, + // saveInvoiceOnly() had already committed on its own transaction boundary -- + // there was nothing left to roll back. + assertThat(invoices.count()).isEqualTo(before + 1); + System.out.println("Invoice count after the blocked, failed send: " + invoices.count() + + " (was " + before + "). The row committed before the send was even attempted."); + } +} diff --git a/outbox/src/test/java/com/ankurm/outboxdemo/billing/FlakyAuditListener.java b/outbox/src/test/java/com/ankurm/outboxdemo/billing/FlakyAuditListener.java new file mode 100644 index 0000000..69c9880 --- /dev/null +++ b/outbox/src/test/java/com/ankurm/outboxdemo/billing/FlakyAuditListener.java @@ -0,0 +1,31 @@ +package com.ankurm.outboxdemo.billing; + +import java.util.concurrent.atomic.AtomicInteger; +import org.springframework.modulith.events.ApplicationModuleListener; +import org.springframework.stereotype.Component; + +/** + * Test-only listener, never part of the packaged application (it lives under + * src/test, so component scanning only ever picks it up inside a test's + * ApplicationContext). Stands in for "the thing that can transiently fail" -- + * in a real system this is the Kafka broker being briefly unreachable, a REST + * call to a partner timing out, or any other downstream dependency with its own + * availability. The event publication registry doesn't know or care which; + * it only knows whether THIS listener, for THIS event, has completed yet. + */ +@Component +public class FlakyAuditListener { + + public static final AtomicInteger attempts = new AtomicInteger(0); + public static volatile boolean failNextAttempt = false; + + @ApplicationModuleListener + public void on(InvoiceIssued event) { + int attempt = attempts.incrementAndGet(); + if (failNextAttempt) { + failNextAttempt = false; + throw new IllegalStateException( + "simulated downstream failure on attempt " + attempt + " for invoice " + event.invoiceId()); + } + } +} diff --git a/outbox/src/test/java/com/ankurm/outboxdemo/billing/RegistrySafetyNetTest.java b/outbox/src/test/java/com/ankurm/outboxdemo/billing/RegistrySafetyNetTest.java new file mode 100644 index 0000000..f7b4353 --- /dev/null +++ b/outbox/src/test/java/com/ankurm/outboxdemo/billing/RegistrySafetyNetTest.java @@ -0,0 +1,91 @@ +package com.ankurm.outboxdemo.billing; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +import com.ankurm.outboxdemo.billing.internal.InvoiceRepository; +import java.time.Duration; +import java.util.Map; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.common.serialization.StringDeserializer; +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.jdbc.core.JdbcTemplate; +import org.springframework.kafka.test.context.EmbeddedKafka; + +/** + * The real production path: BillingManagement.issueInvoice() publishes InvoiceIssued, + * which is @Externalized, which spring-modulith-events-kafka routes to a real broker -- + * a genuine embedded Kafka instance here, not a mock. Proves two things at once: the + * message really arrives on the topic, and the registry really marks that listener's + * row complete (and, with completion-mode=ARCHIVE configured for this module, moves it + * into EVENT_PUBLICATION_ARCHIVE) once it does. + */ +@SpringBootTest +@EmbeddedKafka(partitions = 1, topics = "invoices", bootstrapServersProperty = "spring.kafka.bootstrap-servers") +class RegistrySafetyNetTest { + + @Autowired + BillingManagement billingManagement; + + @Autowired + InvoiceRepository invoices; + + @Autowired + JdbcTemplate jdbc; + + @Autowired + org.springframework.core.env.Environment env; + + KafkaConsumer consumer; + + @AfterEach + void closeConsumer() { + if (consumer != null) { + consumer.close(); + } + } + + @Test + void invoiceIssuedIsReallyPublishedToKafkaAndRegistryCompletes() { + String bootstrapServers = env.getProperty("spring.kafka.bootstrap-servers"); + consumer = new KafkaConsumer<>(Map.of( + "bootstrap.servers", bootstrapServers, + "group.id", "test-group", + "auto.offset.reset", "earliest", + "key.deserializer", StringDeserializer.class.getName(), + "value.deserializer", StringDeserializer.class.getName())); + consumer.subscribe(java.util.List.of("invoices")); + + var invoice = billingManagement.issueInvoice("GADGET-7", 4500); + assertThat(invoices.findById(invoice.getId())).isPresent(); + + var records = new java.util.ArrayList>(); + await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> { + consumer.poll(Duration.ofMillis(200)).forEach(records::add); + assertThat(records).isNotEmpty(); + }); + + ConsumerRecord received = records.get(0); + assertThat(received.topic()).isEqualTo("invoices"); + assertThat(received.key()).isEqualTo(invoice.getId().toString()); + System.out.println("Kafka record received -- topic=" + received.topic() + + " key=" + received.key() + " value=" + received.value()); + + await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> { + Integer archived = jdbc.queryForObject( + "select count(*) from EVENT_PUBLICATION_ARCHIVE where EVENT_TYPE like '%InvoiceIssued%'", + Integer.class); + System.out.println("EVENT_PUBLICATION_ARCHIVE rows for InvoiceIssued: " + archived); + assertThat(archived).isGreaterThanOrEqualTo(1); + }); + + Integer stillOutstanding = jdbc.queryForObject( + "select count(*) from EVENT_PUBLICATION where EVENT_TYPE like '%InvoiceIssued%'", + Integer.class); + System.out.println("EVENT_PUBLICATION rows still outstanding for InvoiceIssued: " + stillOutstanding); + } +} diff --git a/outbox/src/test/java/com/ankurm/outboxdemo/billing/ReplayIncompletePublicationsTest.java b/outbox/src/test/java/com/ankurm/outboxdemo/billing/ReplayIncompletePublicationsTest.java new file mode 100644 index 0000000..35a288f --- /dev/null +++ b/outbox/src/test/java/com/ankurm/outboxdemo/billing/ReplayIncompletePublicationsTest.java @@ -0,0 +1,90 @@ +package com.ankurm.outboxdemo.billing; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +import java.time.Duration; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.modulith.events.IncompleteEventPublications; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; + +/** + * FlakyAuditListener throws on its first invocation and succeeds on every one after + * that. This is standing in for "the downstream dependency was briefly unavailable" -- + * a Kafka broker, a partner API, anything -- without needing to actually break a real + * broker's availability mid-test. What is NOT simulated is the registry behavior: the + * incomplete publication really is left in EVENT_PUBLICATION with no completion date + * after the first failure, and resubmitIncompletePublicationsOlderThan really does + * re-invoke the listener and really does mark the row complete on success. Verified + * against the compiled IncompleteEventPublications interface (events-api 2.1.1): + * resubmitIncompletePublicationsOlderThan(Duration) is a real method, not a guess. + */ +@SpringBootTest +class ReplayIncompletePublicationsTest { + + // This test is about the FlakyAuditListener, not Kafka -- but spring-modulith-events-kafka + // is still on the classpath and still tries to externalize InvoiceIssued to a default + // localhost:9092 that doesn't exist here. Trimmed down purely so the JVM doesn't sit + // retrying metadata fetches for 30 seconds during shutdown; it does not affect what this + // test is actually asserting. + @DynamicPropertySource + static void fastFailKafka(DynamicPropertyRegistry registry) { + registry.add("spring.kafka.producer.properties.max.block.ms", () -> "500"); + } + + @Autowired + BillingManagement billingManagement; + + @Autowired + IncompleteEventPublications incompletePublications; + + @Autowired + JdbcTemplate jdbc; + + @BeforeEach + void reset() { + FlakyAuditListener.attempts.set(0); + FlakyAuditListener.failNextAttempt = false; + } + + @Test + void failedListenerLeavesAnIncompletePublicationThatResubmissionCompletes() { + FlakyAuditListener.failNextAttempt = true; + + var invoice = billingManagement.issueInvoice("GIZMO-3", 999); + + await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> + assertThat(FlakyAuditListener.attempts.get()).isEqualTo(1)); + + Integer incompleteAfterFailure = jdbc.queryForObject( + "select count(*) from EVENT_PUBLICATION where EVENT_TYPE like '%InvoiceIssued%' " + + "and LISTENER_ID like '%FlakyAuditListener%' and COMPLETION_DATE is null", + Integer.class); + System.out.println("Incomplete FlakyAuditListener publications after the simulated failure: " + + incompleteAfterFailure); + assertThat(incompleteAfterFailure).isEqualTo(1); + + incompletePublications.resubmitIncompletePublicationsOlderThan(Duration.ZERO); + + await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> + assertThat(FlakyAuditListener.attempts.get()).isEqualTo(2)); + + await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> { + Integer incompleteAfterResubmit = jdbc.queryForObject( + "select count(*) from EVENT_PUBLICATION where EVENT_TYPE like '%InvoiceIssued%' " + + "and LISTENER_ID like '%FlakyAuditListener%' and COMPLETION_DATE is null", + Integer.class); + System.out.println("Incomplete FlakyAuditListener publications after resubmission: " + + incompleteAfterResubmit); + assertThat(incompleteAfterResubmit).isEqualTo(0); + }); + + System.out.println("Invoice " + invoice.getId() + " never changed -- the retry was purely" + + " about redelivering the event, the invoice row was correct the whole time."); + } +} diff --git a/pom.xml b/pom.xml index 9cd5681..9d32aad 100644 --- a/pom.xml +++ b/pom.xml @@ -14,6 +14,7 @@ Companion repository for the ankurm.com Spring Modulith series. 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 @@ -25,5 +26,6 @@ order-fulfillment + outbox