'
+ spring-boot-kafka
+ 4.1.1
+ spring-boot-kafka
+--
+ spring-boot
+ 4.1.1
+ compile
+--
+ spring-boot-transaction
+ 4.1.1
+ compile
+--
+ spring-kafka
+ 4.1.1
+ compile
diff --git a/awaitility/pom.xml b/awaitility/pom.xml
new file mode 100644
index 0000000..82e3ae3
--- /dev/null
+++ b/awaitility/pom.xml
@@ -0,0 +1,83 @@
+
+
+ 4.0.0
+
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ 4.1.1
+
+
+
+ com.ankurm
+ awaitility
+ 1.0
+ jar
+
+
+ 25
+ UTF-8
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-kafka
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-kafka-test
+ test
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-maven-plugin
+
+
+ org.apache.maven.plugins
+ maven-surefire-plugin
+ 3.6.0
+
+
+
+
diff --git a/awaitility/scripts/run-all.sh b/awaitility/scripts/run-all.sh
new file mode 100755
index 0000000..a1e2639
--- /dev/null
+++ b/awaitility/scripts/run-all.sh
@@ -0,0 +1,13 @@
+#!/usr/bin/env bash
+# Regenerates every file under docs/output/. Nothing in this module's documentation, or in the
+# article it accompanies, is typed by hand.
+#
+# scripts/run-all.sh
+set -euo pipefail
+cd "$(dirname "$0")/.."
+mkdir -p docs/output
+
+mvn -B test > /tmp/awaitility-run.log 2>&1 || { cat /tmp/awaitility-run.log; exit 1; }
+grep -E 'Running |Tests run:|BUILD ' /tmp/awaitility-run.log > docs/output/00-full-test-run.txt
+echo "wrote docs/output/00-full-test-run.txt"
+echo "the rest of docs/output/*.txt is written by the tests themselves (see Capture.java)"
diff --git a/awaitility/src/main/java/com/ankurm/awaitility/AsyncGreetingService.java b/awaitility/src/main/java/com/ankurm/awaitility/AsyncGreetingService.java
new file mode 100644
index 0000000..ce13925
--- /dev/null
+++ b/awaitility/src/main/java/com/ankurm/awaitility/AsyncGreetingService.java
@@ -0,0 +1,37 @@
+package com.ankurm.awaitility;
+
+import org.springframework.scheduling.annotation.Async;
+import org.springframework.stereotype.Service;
+
+/**
+ * A deliberately {@code void} {@code @Async} method: this is the case from
+ * docs/01-what-async-actually-does.md (in the sibling {@code async} module) where there is no
+ * {@code Future} the caller can block on, by design — the method is meant to be
+ * fire-and-forget. Fire-and-forget is also the exact shape of code that pushes people toward
+ * {@code Thread.sleep} in tests, because there is nothing else to wait on.
+ */
+@Service
+public class AsyncGreetingService {
+
+ private final Mailbox mailbox;
+
+ public AsyncGreetingService(Mailbox mailbox) {
+ this.mailbox = mailbox;
+ }
+
+ @Async
+ public void sendConfirmation(String orderId) {
+ sleepFor(220);
+ mailbox.record(orderId, "order " + orderId + " confirmed");
+ }
+
+ private static void sleepFor(long millis) {
+ try {
+ Thread.sleep(millis);
+ }
+ catch (InterruptedException ex) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException(ex);
+ }
+ }
+}
diff --git a/awaitility/src/main/java/com/ankurm/awaitility/AwaitilityDemoApplication.java b/awaitility/src/main/java/com/ankurm/awaitility/AwaitilityDemoApplication.java
new file mode 100644
index 0000000..2a1fa25
--- /dev/null
+++ b/awaitility/src/main/java/com/ankurm/awaitility/AwaitilityDemoApplication.java
@@ -0,0 +1,25 @@
+package com.ankurm.awaitility;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.scheduling.annotation.EnableAsync;
+import org.springframework.scheduling.annotation.EnableScheduling;
+
+/**
+ * Entry point for the Awaitility demonstrations.
+ *
+ * {@code @EnableAsync} and {@code @EnableScheduling} turn on the two Spring mechanisms this
+ * module tests against: a fire-and-forget {@code @Async} method with no {@code Future} to block
+ * on, and a {@code @Scheduled} job running on its own thread. Kafka needs no equivalent
+ * annotation here — {@code spring-kafka} on the classpath is enough for Spring Boot's
+ * {@code KafkaAutoConfiguration} to wire a listener container factory.
+ */
+@SpringBootApplication
+@EnableAsync
+@EnableScheduling
+public class AwaitilityDemoApplication {
+
+ public static void main(String[] args) {
+ SpringApplication.run(AwaitilityDemoApplication.class, args);
+ }
+}
diff --git a/awaitility/src/main/java/com/ankurm/awaitility/FlakyStartupResource.java b/awaitility/src/main/java/com/ankurm/awaitility/FlakyStartupResource.java
new file mode 100644
index 0000000..7d8ada6
--- /dev/null
+++ b/awaitility/src/main/java/com/ankurm/awaitility/FlakyStartupResource.java
@@ -0,0 +1,30 @@
+package com.ankurm.awaitility;
+
+import java.time.Duration;
+import java.time.Instant;
+import java.util.concurrent.atomic.AtomicReference;
+import org.springframework.stereotype.Component;
+
+/**
+ * Stands in for anything that is genuinely not ready yet right after it starts warming up
+ * — a connection pool, a cache still loading, a client still completing its first
+ * handshake — and throws rather than returning a sentinel while that is true. Call
+ * {@link #startWarmingUp()} to begin a 300ms window during which {@link #value()} throws
+ * {@link IllegalStateException}; after that it returns normally.
+ */
+@Component
+public class FlakyStartupResource {
+
+ private final AtomicReference readyAt = new AtomicReference<>(Instant.EPOCH);
+
+ public void startWarmingUp() {
+ readyAt.set(Instant.now().plus(Duration.ofMillis(300)));
+ }
+
+ public String value() {
+ if (Instant.now().isBefore(readyAt.get())) {
+ throw new IllegalStateException("resource is still warming up");
+ }
+ return "ready";
+ }
+}
diff --git a/awaitility/src/main/java/com/ankurm/awaitility/Mailbox.java b/awaitility/src/main/java/com/ankurm/awaitility/Mailbox.java
new file mode 100644
index 0000000..fda8b05
--- /dev/null
+++ b/awaitility/src/main/java/com/ankurm/awaitility/Mailbox.java
@@ -0,0 +1,31 @@
+package com.ankurm.awaitility;
+
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import org.springframework.stereotype.Component;
+
+/**
+ * Stands in for "the side effect of an asynchronous operation" — a row written by a
+ * background thread that a test (or another part of the application) has no handle to wait on
+ * directly. {@link AsyncGreetingService#sendConfirmation(String)} writes here; nothing returns a
+ * {@code Future} for it to block on, which is the whole reason the test for it needs Awaitility
+ * rather than a blocking {@code get()}.
+ */
+@Component
+public class Mailbox {
+
+ private final ConcurrentMap sent = new ConcurrentHashMap<>();
+
+ void record(String orderId, String message) {
+ sent.put(orderId, message);
+ }
+
+ /** Returns {@code null} if nothing has been recorded for this order yet. */
+ public String confirmationFor(String orderId) {
+ return sent.get(orderId);
+ }
+
+ public int size() {
+ return sent.size();
+ }
+}
diff --git a/awaitility/src/main/java/com/ankurm/awaitility/OrderEventListener.java b/awaitility/src/main/java/com/ankurm/awaitility/OrderEventListener.java
new file mode 100644
index 0000000..2863f88
--- /dev/null
+++ b/awaitility/src/main/java/com/ankurm/awaitility/OrderEventListener.java
@@ -0,0 +1,34 @@
+package com.ankurm.awaitility;
+
+import java.util.concurrent.ConcurrentLinkedQueue;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.kafka.annotation.KafkaListener;
+import org.springframework.stereotype.Component;
+
+/**
+ * Consumes from the {@code orders} topic on its own listener container thread. A test that sends
+ * a record with {@code KafkaTemplate} and then immediately inspects {@link #received()} is racing
+ * the consumer's poll loop — there is no return value, no {@code Future}, and no fixed
+ * delay that is both fast enough to keep the test quick and slow enough to never flake under CI
+ * load. See docs/output/05-kafka-listener-await.txt for why this module's own test does not use
+ * {@code Thread.sleep} either.
+ *
+ * Gated behind {@code app.kafka-demo.enabled} (off by default, turned on only by the one test
+ * that needs it) so every other test's {@code @SpringBootTest} context does not also try to open
+ * a consumer against a broker that, for them, does not exist.
+ */
+@Component
+@ConditionalOnProperty(name = "app.kafka-demo.enabled", havingValue = "true")
+public class OrderEventListener {
+
+ private final ConcurrentLinkedQueue received = new ConcurrentLinkedQueue<>();
+
+ @KafkaListener(topics = "orders", groupId = "awaitility-demo")
+ public void onOrderEvent(String payload) {
+ received.add(payload);
+ }
+
+ public ConcurrentLinkedQueue received() {
+ return received;
+ }
+}
diff --git a/awaitility/src/main/java/com/ankurm/awaitility/ReportJob.java b/awaitility/src/main/java/com/ankurm/awaitility/ReportJob.java
new file mode 100644
index 0000000..525f4cc
--- /dev/null
+++ b/awaitility/src/main/java/com/ankurm/awaitility/ReportJob.java
@@ -0,0 +1,26 @@
+package com.ankurm.awaitility;
+
+import java.util.concurrent.atomic.AtomicInteger;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+/**
+ * A trivial {@code @Scheduled} job: it exists only so a test has an asynchronous, externally
+ * triggered counter to wait on instead of sleeping for a guessed multiple of the fixed rate. See
+ * the sibling {@code scheduling} module for what happens to a job like this across three
+ * replicas; this one runs standalone and is not locked.
+ */
+@Component
+public class ReportJob {
+
+ private final AtomicInteger runs = new AtomicInteger();
+
+ @Scheduled(fixedRate = 150)
+ public void run() {
+ runs.incrementAndGet();
+ }
+
+ public int runCount() {
+ return runs.get();
+ }
+}
diff --git a/awaitility/src/main/resources/application.yaml b/awaitility/src/main/resources/application.yaml
new file mode 100644
index 0000000..ddd68e2
--- /dev/null
+++ b/awaitility/src/main/resources/application.yaml
@@ -0,0 +1,10 @@
+spring:
+ application:
+ name: awaitility-demo
+ # No spring.kafka.bootstrap-servers here. src/test/resources/application-test.yaml points it at
+ # ${spring.embedded.kafka.brokers}, the system property EmbeddedKafkaBroker sets for the
+ # duration of each test that carries @EmbeddedKafka.
+
+app:
+ kafka-demo:
+ enabled: false # flipped on with a single @TestPropertySource, only by the Kafka test
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/AwaitAsyncConfirmationTest.java b/awaitility/src/test/java/com/ankurm/awaitility/AwaitAsyncConfirmationTest.java
new file mode 100644
index 0000000..c0c547e
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/AwaitAsyncConfirmationTest.java
@@ -0,0 +1,58 @@
+package com.ankurm.awaitility;
+
+import java.time.Duration;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * {@link AsyncGreetingService#sendConfirmation(String)} is {@code void}: there is no
+ * {@code Future} to call {@code get()} on, so the only way to find out when the background
+ * thread has written to {@link Mailbox} is to look — repeatedly, until it's there.
+ *
+ * docs/output/01-sleep-guesses-wrong.txt is the real failure from the version of this test
+ * that guessed a fixed delay instead: {@code Thread.sleep(100)} against a confirmation that
+ * actually takes 220ms. It failed on every run, not just an unlucky one, which is the more
+ * common shape of a bad guess than the intermittent kind "flaky" usually implies.
+ */
+@SpringBootTest
+class AwaitAsyncConfirmationTest {
+
+ @Autowired
+ private AsyncGreetingService service;
+
+ @Autowired
+ private Mailbox mailbox;
+
+ @Test
+ void confirmationArrivesEventually() {
+ service.sendConfirmation("ORD-1");
+
+ await().atMost(Duration.ofSeconds(2))
+ .untilAsserted(() -> assertThat(mailbox.confirmationFor("ORD-1")).isNotNull());
+
+ assertThat(mailbox.confirmationFor("ORD-1")).isEqualTo("order ORD-1 confirmed");
+ }
+
+ @Test
+ void capture() {
+ String before = mailbox.confirmationFor("ORD-CAPTURE");
+ service.sendConfirmation("ORD-CAPTURE");
+ long start = System.nanoTime();
+
+ await().atMost(Duration.ofSeconds(2))
+ .untilAsserted(() -> assertThat(mailbox.confirmationFor("ORD-CAPTURE")).isNotNull());
+
+ long elapsedMs = (System.nanoTime() - start) / 1_000_000;
+ Capture.write("02-await-finds-it.txt",
+ "await().untilAsserted() against a 220ms-delayed void @Async method",
+ """
+ mailbox.confirmationFor("ORD-CAPTURE") immediately before sendConfirmation(): %s
+ await().atMost(2s).untilAsserted(...) returned after approximately: %dms
+ mailbox.confirmationFor("ORD-CAPTURE") once await() returns: %s
+ """.formatted(before, elapsedMs, mailbox.confirmationFor("ORD-CAPTURE")));
+ }
+}
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/Capture.java b/awaitility/src/test/java/com/ankurm/awaitility/Capture.java
new file mode 100644
index 0000000..4b670d5
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/Capture.java
@@ -0,0 +1,23 @@
+package com.ankurm.awaitility;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+/** Writes a transcript under docs/output/ so every number in the article has a file behind it. */
+final class Capture {
+
+ private Capture() {
+ }
+
+ static void write(String fileName, String heading, String body) {
+ Path dir = Path.of(System.getProperty("user.dir"), "docs", "output");
+ try {
+ Files.createDirectories(dir);
+ Files.writeString(dir.resolve(fileName), "== " + heading + " ==\n\n" + body + "\n");
+ }
+ catch (IOException ex) {
+ throw new IllegalStateException("could not write " + fileName, ex);
+ }
+ }
+}
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/DefaultTimingTest.java b/awaitility/src/test/java/com/ankurm/awaitility/DefaultTimingTest.java
new file mode 100644
index 0000000..5ff03c5
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/DefaultTimingTest.java
@@ -0,0 +1,73 @@
+package com.ankurm.awaitility;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * Proves two numbers from the real jar instead of quoting them from the reference guide:
+ * {@code org.awaitility.Awaitility}'s {@code defaultWaitConstraint} is
+ * {@code AtMostWaitConstraint.TEN_SECONDS} and its {@code defaultPollInterval} is a fixed 100ms
+ * — both read directly from {@code Awaitility.java} and {@code AtMostWaitConstraint.java}
+ * in the 4.3.0 sources jar. Neither test overrides {@code atMost(...)} or
+ * {@code pollInterval(...)}.
+ */
+class DefaultTimingTest {
+
+ @Test
+ void theDefaultTimeoutIsTenSeconds() {
+ long start = System.nanoTime();
+
+ assertThatTimeoutIsThrownBy(() -> await().until(() -> false));
+
+ long elapsedMs = (System.nanoTime() - start) / 1_000_000;
+ // Allow scheduling jitter; the point is "about 10 seconds", not exactly 10,000ms.
+ assertThat(elapsedMs).isBetween(9_000L, 11_000L);
+
+ Capture.write("03-default-timeout-is-ten-seconds.txt",
+ "await().until(() -> false) with no atMost(...) override",
+ "elapsed before ConditionTimeoutException: %dms (expected: ~10,000ms)".formatted(elapsedMs));
+ }
+
+ @Test
+ void theDefaultPollIntervalIsOneHundredMillis() {
+ List pollTimestampsMs = new ArrayList<>();
+ long start = System.nanoTime();
+
+ assertThatTimeoutIsThrownBy(() -> await()
+ .atMost(java.time.Duration.ofMillis(550))
+ .until(() -> {
+ pollTimestampsMs.add((System.nanoTime() - start) / 1_000_000);
+ return false;
+ }));
+
+ List gapsMs = new ArrayList<>();
+ for (int i = 1; i < pollTimestampsMs.size(); i++) {
+ gapsMs.add(pollTimestampsMs.get(i) - pollTimestampsMs.get(i - 1));
+ }
+ double averageGapMs = gapsMs.stream().mapToLong(Long::longValue).average().orElse(0);
+
+ assertThat(averageGapMs).isBetween(80.0, 120.0);
+
+ Capture.write("04-default-poll-interval-is-100ms.txt",
+ "await().atMost(550ms).until(...) with no pollInterval(...) override",
+ """
+ poll timestamps (ms since start): %s
+ gaps between consecutive polls (ms): %s
+ average gap: %.1fms (expected: ~100ms)
+ """.formatted(pollTimestampsMs, gapsMs, averageGapMs));
+ }
+
+ private static void assertThatTimeoutIsThrownBy(Runnable awaitCall) {
+ try {
+ awaitCall.run();
+ throw new AssertionError("expected a ConditionTimeoutException but none was thrown");
+ }
+ catch (org.awaitility.core.ConditionTimeoutException expected) {
+ // the default timeout firing is the point of this test
+ }
+ }
+}
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/IgnoreExceptionsTest.java b/awaitility/src/test/java/com/ankurm/awaitility/IgnoreExceptionsTest.java
new file mode 100644
index 0000000..6e3d9b5
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/IgnoreExceptionsTest.java
@@ -0,0 +1,50 @@
+package com.ankurm.awaitility;
+
+import java.time.Duration;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * By default Awaitility does not treat an exception thrown inside the polled condition
+ * as "still false" — {@code defaultExceptionIgnorer} in {@code Awaitility.java} is a
+ * predicate that returns {@code false} for every exception, meaning nothing is ignored. An
+ * exception from the very first poll propagates immediately and fails the test, without waiting
+ * out the {@code atMost(...)} window at all. docs/output/08-exception-propagates-immediately.txt
+ * is the real, immediate failure from the version of this test that left
+ * {@code ignoreExceptionsInstanceOf(...)} out.
+ */
+@SpringBootTest
+class IgnoreExceptionsTest {
+
+ @Autowired
+ private FlakyStartupResource resource;
+
+ @Test
+ void ignoredExceptionsAreTreatedAsNotYetTrue() {
+ resource.startWarmingUp();
+
+ await().atMost(Duration.ofSeconds(1))
+ .ignoreExceptionsInstanceOf(IllegalStateException.class)
+ .untilAsserted(() -> assertThat(resource.value()).isEqualTo("ready"));
+ }
+
+ @Test
+ void capture() {
+ resource.startWarmingUp();
+ long start = System.nanoTime();
+
+ await().atMost(Duration.ofSeconds(1))
+ .ignoreExceptionsInstanceOf(IllegalStateException.class)
+ .untilAsserted(() -> assertThat(resource.value()).isEqualTo("ready"));
+
+ long elapsedMs = (System.nanoTime() - start) / 1_000_000;
+ Capture.write("09-ignore-exceptions-waits-it-out.txt",
+ "await().ignoreExceptionsInstanceOf(IllegalStateException.class) against a resource that throws for its first 300ms",
+ "approximate elapsed time before resource.value() finally returned \"ready\": %dms (expected: a little over 300ms)"
+ .formatted(elapsedMs));
+ }
+}
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/KafkaListenerAwaitTest.java b/awaitility/src/test/java/com/ankurm/awaitility/KafkaListenerAwaitTest.java
new file mode 100644
index 0000000..17ee412
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/KafkaListenerAwaitTest.java
@@ -0,0 +1,78 @@
+package com.ankurm.awaitility;
+
+import java.time.Duration;
+import java.util.Map;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.kafka.core.DefaultKafkaProducerFactory;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.kafka.core.ProducerFactory;
+import org.springframework.kafka.test.EmbeddedKafkaBroker;
+import org.springframework.kafka.test.context.EmbeddedKafka;
+import org.springframework.kafka.test.utils.KafkaTestUtils;
+import org.springframework.test.context.TestPropertySource;
+import org.springframework.test.context.ActiveProfiles;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * {@code @EmbeddedKafka} starts a real, in-process Kafka broker for the duration of this test
+ * class — no Docker, no Testcontainers. {@link OrderEventListener#onOrderEvent(String)}
+ * consumes on the listener container's own thread, so sending a record with
+ * {@link KafkaTemplate#send} and then immediately reading {@code received()} is racing the
+ * consumer's poll loop exactly the way docs/02-the-self-invocation-trap.md's proxy call races a
+ * background thread, just with a network hop (loopback, but still a hop) in between.
+ */
+@SpringBootTest
+@ActiveProfiles("test")
+@TestPropertySource(properties = "app.kafka-demo.enabled=true")
+@EmbeddedKafka(partitions = 1, topics = "orders")
+class KafkaListenerAwaitTest {
+
+ @Autowired
+ private OrderEventListener listener;
+
+ @Autowired
+ private EmbeddedKafkaBroker embeddedKafkaBroker;
+
+ @Test
+ void messageSentIsEventuallyConsumed() {
+ listener.received().clear();
+ KafkaTemplate producer = newProducer();
+
+ producer.send("orders", "order-42-created");
+
+ await().atMost(Duration.ofSeconds(5))
+ .untilAsserted(() -> assertThat(listener.received()).contains("order-42-created"));
+ }
+
+ @Test
+ void capture() {
+ listener.received().clear();
+ KafkaTemplate producer = newProducer();
+ long start = System.nanoTime();
+
+ producer.send("orders", "order-77-created");
+ await().atMost(Duration.ofSeconds(5))
+ .untilAsserted(() -> assertThat(listener.received()).contains("order-77-created"));
+
+ long elapsedMs = (System.nanoTime() - start) / 1_000_000;
+ Capture.write("05-kafka-listener-await.txt",
+ "await() for an @KafkaListener to consume a record a KafkaTemplate just sent, against an @EmbeddedKafka broker",
+ """
+ listener.received() immediately after send() returned: (not checked -- that is the point)
+ listener.received() once await() returned: %s
+ approximate elapsed time: %dms
+ """.formatted(listener.received(), elapsedMs));
+ }
+
+ private KafkaTemplate newProducer() {
+ Map props = KafkaTestUtils.producerProps(embeddedKafkaBroker);
+ props.put(ProducerConfig.ACKS_CONFIG, "all");
+ ProducerFactory producerFactory = new DefaultKafkaProducerFactory<>(props);
+ return new KafkaTemplate<>(producerFactory);
+ }
+}
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/ScheduledJobAwaitTest.java b/awaitility/src/test/java/com/ankurm/awaitility/ScheduledJobAwaitTest.java
new file mode 100644
index 0000000..efb3c2e
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/ScheduledJobAwaitTest.java
@@ -0,0 +1,53 @@
+package com.ankurm.awaitility;
+
+import java.time.Duration;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * {@link ReportJob} runs on its own {@code fixedRate = 150} schedule, starting the moment the
+ * context comes up — including however long context startup itself took, which varies by
+ * machine and by whether this is the first test in the module or the tenth reusing a cached
+ * context. "Sleep for 3 * 150ms and assert runCount() == 3" is exactly as wrong here as it was
+ * for the {@code @Async} case, for the same reason: the guess has to be right about a timing
+ * detail the test does not control. This test instead waits for at least three more runs
+ * than there were when the test started, which is true regardless of how long the job has
+ * already been running.
+ */
+@SpringBootTest
+class ScheduledJobAwaitTest {
+
+ @Autowired
+ private ReportJob reportJob;
+
+ @Test
+ void atLeastThreeMoreRunsHappenWithinAwaitWindow() {
+ int baseline = reportJob.runCount();
+
+ await().atMost(Duration.ofSeconds(2))
+ .untilAsserted(() -> assertThat(reportJob.runCount()).isGreaterThanOrEqualTo(baseline + 3));
+ }
+
+ @Test
+ void capture() {
+ int baseline = reportJob.runCount();
+ long start = System.nanoTime();
+
+ await().atMost(Duration.ofSeconds(2))
+ .untilAsserted(() -> assertThat(reportJob.runCount()).isGreaterThanOrEqualTo(baseline + 3));
+
+ long elapsedMs = (System.nanoTime() - start) / 1_000_000;
+ Capture.write("06-scheduled-job-await.txt",
+ "await() for 3 more @Scheduled(fixedRate = 150) executions, regardless of how long the job had already been running",
+ """
+ runCount() when the test started: %d
+ target (baseline + 3): %d
+ runCount() once await() returned: %d
+ approximate elapsed time: %dms
+ """.formatted(baseline, baseline + 3, reportJob.runCount(), elapsedMs));
+ }
+}
diff --git a/awaitility/src/test/java/com/ankurm/awaitility/UncaughtExceptionHandlerSwapTest.java b/awaitility/src/test/java/com/ankurm/awaitility/UncaughtExceptionHandlerSwapTest.java
new file mode 100644
index 0000000..acc0d9e
--- /dev/null
+++ b/awaitility/src/test/java/com/ankurm/awaitility/UncaughtExceptionHandlerSwapTest.java
@@ -0,0 +1,49 @@
+package com.ankurm.awaitility;
+
+import java.time.Duration;
+import java.util.concurrent.atomic.AtomicReference;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+
+/**
+ * A fact from reading {@code ConditionAwaiter.java}, not from the user guide: with
+ * {@code catchUncaughtExceptions} at its default of {@code true}, every {@code await(...)} call
+ * installs itself as the JVM's default {@link Thread#setDefaultUncaughtExceptionHandler}
+ * for as long as that one call is polling, then restores whatever handler was there before.
+ * That is a process-wide setting, not a per-thread one — a test that installs its own
+ * default handler and then calls {@code await()} will have it silently replaced until the
+ * {@code await()} call returns.
+ */
+class UncaughtExceptionHandlerSwapTest {
+
+ @Test
+ void awaitTemporarilyInstallsItsOwnDefaultHandlerThenRestoresTheOriginal() {
+ Thread.UncaughtExceptionHandler original = (t, e) -> {
+ };
+ Thread.setDefaultUncaughtExceptionHandler(original);
+
+ AtomicReference handlerDuringAwait = new AtomicReference<>();
+
+ await().atMost(Duration.ofMillis(300)).until(() -> {
+ handlerDuringAwait.set(Thread.getDefaultUncaughtExceptionHandler());
+ return true;
+ });
+
+ Thread.UncaughtExceptionHandler handlerAfterAwait = Thread.getDefaultUncaughtExceptionHandler();
+
+ assertThat(handlerDuringAwait.get()).isNotSameAs(original);
+ assertThat(handlerAfterAwait).isSameAs(original);
+
+ Capture.write("07-uncaught-exception-handler-swap.txt",
+ "Thread.getDefaultUncaughtExceptionHandler() before, during and after one await() call",
+ """
+ handler during await() is the one this test installed: %s (expected: false)
+ handler during await() is Awaitility's own (class %s)
+ handler after await() is the one this test installed: %s (expected: true)
+ """.formatted(handlerDuringAwait.get() == original,
+ handlerDuringAwait.get().getClass().getName(),
+ handlerAfterAwait == original));
+ }
+}
diff --git a/awaitility/src/test/resources/application-test.yaml b/awaitility/src/test/resources/application-test.yaml
new file mode 100644
index 0000000..2ff6605
--- /dev/null
+++ b/awaitility/src/test/resources/application-test.yaml
@@ -0,0 +1,10 @@
+spring:
+ kafka:
+ bootstrap-servers: ${spring.embedded.kafka.brokers}
+ consumer:
+ auto-offset-reset: earliest
+ key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
+ value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
+ producer:
+ key-serializer: org.apache.kafka.common.serialization.StringSerializer
+ value-serializer: org.apache.kafka.common.serialization.StringSerializer