Add awaitility module: testing async @Async/@Scheduled/@KafkaListener code without Thread.sleep

Covers replacing Thread.sleep with await() across a void @Async method, an
@EmbeddedKafka-backed @KafkaListener, and an already-running @Scheduled job;
verifies Awaitility 4.3.0's real defaults (10s timeout, 100ms poll interval)
and its default-uncaught-exception-handler swap directly against the jar;
and documents a real pom.xml trap where Boot 4.1.1 split Kafka's
autoconfiguration (spring-boot-kafka) out of spring-kafka itself, which
silently leaves @KafkaListener beans with no running container.
This commit is contained in:
Claude
2026-10-08 15:18:58 +00:00
parent 09631dcaab
commit a7b76243fa
31 changed files with 900 additions and 0 deletions
+1
View File
@@ -10,6 +10,7 @@ by that module's `scripts/run-all.sh`, never typed by hand.
|---|---|---|
| [`async/`](async/README.md) | [@Async in Spring Boot 4: Executors, Virtual Threads and the Self-Invocation Trap](https://ankurm.com/spring-boot-4-async-executors-virtual-threads/) | Which thread a method actually ran on, in every case where the answer is not the one you expect |
| [`scheduling/`](scheduling/README.md) | [@Scheduled, ShedLock and Distributed Cron: Scheduling That Survives Three Replicas](https://ankurm.com/spring-scheduled-shedlock-distributed-cron/) | Three replicas against one database running the same job three times, then one row and one conditional UPDATE fixing it |
| [`awaitility/`](awaitility/README.md) | Testing Asynchronous Code with Awaitility (@Async, Kafka Listeners, Schedulers) | Replacing `Thread.sleep` with `await()` across a void `@Async` method, an `@KafkaListener`, and a `@Scheduled` job — plus the real pom.xml trap where Boot 4.1.1 split Kafka's autoconfiguration out of `spring-boot-autoconfigure` |
| [`virtual-threads-benchmark/`](virtual-threads-benchmark/README.md) | [Virtual Threads on Spring Boot 4.1: The Benchmarks, Re-Run, and the Pinning Advice That Expired](https://ankurm.com/leveraging-virtual-threads-in-spring-boot-3-4-building-high-throughput-services/) | Platform threads vs virtual threads, re-benchmarked on Boot 4.1.1 / JDK 25, plus JEP 491's fix to `synchronized` pinning proven against a real JDK |
| [`virtual-threads-benchmark-webflux/`](virtual-threads-benchmark-webflux/README.md) | [Virtual Threads vs Reactive (WebFlux) vs Platform Threads: Benchmarks and a Decision Framework](https://ankurm.com/virtual-threads-vs-webflux-vs-platform-threads-spring-boot-benchmarks/) | The WebFlux leg of the three-way comparison, plus the event-loop-starvation failure mode an isolated CPU benchmark can't show |
+84
View File
@@ -0,0 +1,84 @@
# awaitility
Companion code for the Awaitility article on [ankurm.com](https://ankurm.com). All intermediate
and reference-grade depth lives in the post itself (accordions / "going deeper" paragraphs), not
in a `docs/NN-topic.md` chapter folder — the one exception, as in the rest of this repository, is
`docs/output/`, which holds real captured transcripts and nothing else.
## What this is
Six real, runnable demonstrations of the same idea: a fixed `Thread.sleep(...)` in a test is a
guess about timing the test does not control, and `Awaitility.await()` replaces the guess with a
bounded poll loop. Each demonstration targets a different source of asynchrony:
| Test | What it waits for |
|---|---|
| `AwaitAsyncConfirmationTest` | A `void` `@Async` method with no `Future` to block on |
| `DefaultTimingTest` | Nothing — it proves Awaitility's own defaults (10s timeout, 100ms poll interval) against the real jar |
| `KafkaListenerAwaitTest` | An `@KafkaListener` consuming a record a `KafkaTemplate` just sent, against a real in-process `@EmbeddedKafka` broker |
| `ScheduledJobAwaitTest` | Three more executions of an already-running `@Scheduled(fixedRate = 150)` job |
| `UncaughtExceptionHandlerSwapTest` | Nothing — it proves `await()` temporarily installs its own `Thread.setDefaultUncaughtExceptionHandler` and restores the original afterward |
| `IgnoreExceptionsTest` | A resource that throws for its first 300ms, using `ignoreExceptionsInstanceOf(...)` to treat that as "not yet", not a failure |
Two more transcripts in `docs/output/` are real failures from tests that are *not* in the suite
above: `01-sleep-guesses-wrong.txt` (a fixed `Thread.sleep(100)` against a 220ms operation) and
`08-exception-propagates-immediately.txt` (the same resource as `IgnoreExceptionsTest`, minus
`ignoreExceptionsInstanceOf(...)`). Both were real `mvn test` runs, captured once, then the
failing test class was deleted — the mistake is preserved as a transcript, not as a permanently
red test.
## The pom.xml trap this module exists to document
The first version of this module's `pom.xml` depended on `org.springframework.kafka:spring-kafka`
directly, the way every pre-Boot-4.1 tutorial does. It compiled. Every `@SpringBootTest` using
Kafka then failed with an empty `ConcurrentLinkedQueue` and no error at context startup, because
in Boot 4.1.1 the Kafka autoconfiguration classes (`KafkaTemplate`, `ConsumerFactory`, the
`KafkaListenerEndpointRegistry` that `@KafkaListener` needs) moved out of the monolithic
`spring-boot-autoconfigure` jar into their own module, `org.springframework.boot:spring-boot-kafka`
— pulled in by the new `org.springframework.boot:spring-boot-starter-kafka`, not by `spring-kafka`
alone. See `pom.xml`'s own comments and the post for the full diagnosis.
## Versions
Read from `spring-boot-dependencies-4.1.1.pom` and `spring-boot-starter-test-4.1.1.pom` on Maven
Central, not from release notes: **JDK 25** (Temurin 25.0.4.1+1), **Spring Boot 4.1.1**,
**Spring Kafka 4.1.1**, **Awaitility 4.3.0** (already on the classpath via
`spring-boot-starter-test` — no explicit `<dependency>` for it anywhere in `pom.xml`).
## Quickstart
```
mvn test
```
11 tests, 0 failures. `KafkaListenerAwaitTest` starts a real in-process Kafka broker
(`@EmbeddedKafka`) and takes a few seconds; `DefaultTimingTest` deliberately waits out Awaitility's
real 10-second default timeout once, so the suite as a whole takes about 30 seconds.
## Captured output
| File | What it's from |
|---|---|
| `00-full-test-run.txt` | The full `mvn test` run, 11/11 green |
| `01-sleep-guesses-wrong.txt` | Real failure: a guessed `Thread.sleep(100)` against a 220ms operation (test since removed) |
| `02-await-finds-it.txt` | The fixed version: `await().untilAsserted(...)` against the same operation |
| `03-default-timeout-is-ten-seconds.txt` | `await().until(() -> false)` timing out at ~10,000ms with no override |
| `04-default-poll-interval-is-100ms.txt` | Real poll timestamps, ~100ms apart, with no override |
| `05-kafka-listener-await.txt` | `await()` for an `@KafkaListener` to consume a record just sent |
| `06-scheduled-job-await.txt` | `await()` for 3 more executions of an already-running `@Scheduled` job |
| `07-uncaught-exception-handler-swap.txt` | Proof that `await()` swaps and restores the JVM's default uncaught-exception handler |
| `08-exception-propagates-immediately.txt` | Real failure: an exception from the polled condition, with no `ignoreExceptionsInstanceOf(...)`, failing on the first poll (test since removed) |
| `09-ignore-exceptions-waits-it-out.txt` | The fixed version: the same resource, with `ignoreExceptionsInstanceOf(...)` |
| `10-missing-kafka-starter-failure.txt` | Real failure: the Kafka test with `spring-kafka` as a direct dependency instead of `spring-boot-starter-kafka` (pom.xml since fixed) |
| `11-pom-diagnosis.txt` | The real `curl`/`grep` transcript that found the Boot 4.1 Kafka module split against Maven Central's own POMs |
## What's sourced from documentation, not run here
Awaitility's own `@since` Javadoc tags, Spring Boot's own migration notes for the Kafka module
split, and the general shape of `@KafkaListener`/`@EmbeddedKafka` wiring are cited from the real
jars and POMs on Maven Central (see the post for exact artifact coordinates), not re-derived from
scratch in this module.
## Licence
MIT — see [LICENSE](../LICENSE).
@@ -0,0 +1,14 @@
[INFO] Running com.ankurm.awaitility.KafkaListenerAwaitTest
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 8.012 s -- in com.ankurm.awaitility.KafkaListenerAwaitTest
[INFO] Running com.ankurm.awaitility.ScheduledJobAwaitTest
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 1.841 s -- in com.ankurm.awaitility.ScheduledJobAwaitTest
[INFO] Running com.ankurm.awaitility.UncaughtExceptionHandlerSwapTest
[INFO] Tests run: 1, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 0.115 s -- in com.ankurm.awaitility.UncaughtExceptionHandlerSwapTest
[INFO] Running com.ankurm.awaitility.AwaitAsyncConfirmationTest
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 0.654 s -- in com.ankurm.awaitility.AwaitAsyncConfirmationTest
[INFO] Running com.ankurm.awaitility.IgnoreExceptionsTest
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 0.636 s -- in com.ankurm.awaitility.IgnoreExceptionsTest
[INFO] Running com.ankurm.awaitility.DefaultTimingTest
[INFO] Tests run: 2, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 10.70 s -- in com.ankurm.awaitility.DefaultTimingTest
[INFO] Tests run: 11, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS
@@ -0,0 +1,9 @@
== Thread.sleep(100) against a confirmation that actually takes 220ms (real failure, test since removed) ==
$ mvn -B -Dtest=SleepGuessesWrongTest test
[ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0, Time elapsed: 3.568 s <<< FAILURE! -- in com.ankurm.awaitility.SleepGuessesWrongTest
com.ankurm.awaitility.SleepGuessesWrongTest.confirmationArrivesWithin100ms -- Time elapsed: 1.084 s <<< FAILURE!
java.lang.AssertionError:
Expecting actual not to be null
at com.ankurm.awaitility.SleepGuessesWrongTest.confirmationArrivesWithin100ms(SleepGuessesWrongTest.java:27)
@@ -0,0 +1,6 @@
== await().untilAsserted() against a 220ms-delayed void @Async method ==
mailbox.confirmationFor("ORD-CAPTURE") immediately before sendConfirmation(): null
await().atMost(2s).untilAsserted(...) returned after approximately: 307ms
mailbox.confirmationFor("ORD-CAPTURE") once await() returns: order ORD-CAPTURE confirmed
@@ -0,0 +1,3 @@
== await().until(() -> false) with no atMost(...) override ==
elapsed before ConditionTimeoutException: 10057ms (expected: ~10,000ms)
@@ -0,0 +1,6 @@
== await().atMost(550ms).until(...) with no pollInterval(...) override ==
poll timestamps (ms since start): [107, 207, 308, 409, 509]
gaps between consecutive polls (ms): [100, 101, 101, 100]
average gap: 100.5ms (expected: ~100ms)
@@ -0,0 +1,6 @@
== 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: [order-77-created]
approximate elapsed time: 378ms
@@ -0,0 +1,7 @@
== await() for 3 more @Scheduled(fixedRate = 150) executions, regardless of how long the job had already been running ==
runCount() when the test started: 1
target (baseline + 3): 4
runCount() once await() returned: 4
approximate elapsed time: 510ms
@@ -0,0 +1,6 @@
== Thread.getDefaultUncaughtExceptionHandler() before, during and after one await() call ==
handler during await() is the one this test installed: false (expected: false)
handler during await() is Awaitility's own (class org.awaitility.core.CallableCondition$1)
handler after await() is the one this test installed: true (expected: true)
@@ -0,0 +1,11 @@
== await().untilAsserted(...) with no ignoreExceptionsInstanceOf(...), against a resource that throws for its first 300ms (real failure, test since removed) ==
$ mvn -B -Dtest=ExceptionPropagatesImmediatelyTest test
[ERROR] Tests run: 1, Failures: 0, Errors: 1, Skipped: 0, Time elapsed: 2.940 s <<< FAILURE! -- in com.ankurm.awaitility.ExceptionPropagatesImmediatelyTest
com.ankurm.awaitility.ExceptionPropagatesImmediatelyTest.failsOnTheFirstPollInsteadOfWaiting -- Time elapsed: 0.946 s <<< ERROR!
java.lang.IllegalStateException: resource is still warming up
at com.ankurm.awaitility.FlakyStartupResource.value(FlakyStartupResource.java:26)
at com.ankurm.awaitility.ExceptionPropagatesImmediatelyTest.lambda$failsOnTheFirstPollInsteadOfWaiting$0(ExceptionPropagatesImmediatelyTest.java:28)
Note "Errors: 1", not a ConditionTimeoutException after the full atMost(1s) window: the exception
from the very first poll propagated straight out instead of being retried.
@@ -0,0 +1,3 @@
== await().ignoreExceptionsInstanceOf(IllegalStateException.class) against a resource that throws for its first 300ms ==
approximate elapsed time before resource.value() finally returned "ready": 302ms (expected: a little over 300ms)
@@ -0,0 +1,23 @@
== Real failure with org.springframework.kafka:spring-kafka as a direct dependency instead of org.springframework.boot:spring-boot-starter-kafka (pom.xml since fixed) ==
$ mvn -B test
[ERROR] Errors:
[ERROR] KafkaListenerAwaitTest.messageSentIsEventuallyConsumed:49 ? ConditionTimeout Assertion condition defined as a Lambda expression in com.ankurm.awaitility.KafkaListenerAwaitTest
Expecting ConcurrentLinkedQueue:
[]
to contain:
["order-42-created"]
but could not find the following element(s):
["order-42-created"]
within 5 seconds.
[INFO]
[ERROR] Tests run: 11, Failures: 0, Errors: 2, Skipped: 0
No exception at context startup, and the application context came up without a single warning.
Autowiring KafkaListenerEndpointRegistry directly in a scratch test and calling
getListenerContainers() on it confirmed the real cause: the bean did not exist. Neither did
KafkaTemplate, ConsumerFactory, or kafkaListenerContainerFactory. Grepping the full log for the
Spring Kafka logger prefix also confirmed it:
$ grep -c "o\.s\.k\." full-test-run.log
0
@@ -0,0 +1,38 @@
== Confirming the Boot 4.1 Kafka module split directly against Maven Central's real POMs ==
$ curl -sS https://repo1.maven.org/maven2/org/springframework/boot/spring-boot-dependencies/4.1.1/spring-boot-dependencies-4.1.1.pom -o spring-boot-dependencies-4.1.1.pom
$ grep -B3 '<artifactId>spring-boot-kafka</artifactId>' spring-boot-dependencies-4.1.1.pom
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-kafka</artifactId>
$ curl -sS https://repo1.maven.org/maven2/org/springframework/boot/spring-boot-starter-kafka/4.1.1/spring-boot-starter-kafka-4.1.1.pom | grep -A2 '<artifactId>'
<artifactId>spring-boot-starter-kafka</artifactId>
<version>4.1.1</version>
<name>spring-boot-starter-kafka</name>
--
<artifactId>spring-boot-starter</artifactId>
<version>4.1.1</version>
<scope>compile</scope>
--
<artifactId>spring-boot-kafka</artifactId>
<version>4.1.1</version>
<scope>compile</scope>
$ curl -sS https://repo1.maven.org/maven2/org/springframework/boot/spring-boot-kafka/4.1.1/spring-boot-kafka-4.1.1.pom | grep -A2 '<artifactId>'
<artifactId>spring-boot-kafka</artifactId>
<version>4.1.1</version>
<name>spring-boot-kafka</name>
--
<artifactId>spring-boot</artifactId>
<version>4.1.1</version>
<scope>compile</scope>
--
<artifactId>spring-boot-transaction</artifactId>
<version>4.1.1</version>
<scope>compile</scope>
--
<artifactId>spring-kafka</artifactId>
<version>4.1.1</version>
<scope>compile</scope>
+83
View File
@@ -0,0 +1,83 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<!-- Boot 4.1.1 manages Spring Framework 7.0.9, Spring Kafka 4.1.1 and Awaitility 4.3.0.
Versions were read from spring-boot-dependencies-4.1.1.pom and
spring-boot-starter-test-4.1.1.pom, not from release notes. -->
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>4.1.1</version>
<relativePath/>
</parent>
<groupId>com.ankurm</groupId>
<artifactId>awaitility</artifactId>
<version>1.0</version>
<packaging>jar</packaging>
<properties>
<java.version>25</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<!-- NOT org.springframework.kafka:spring-kafka directly, that was Boot 3.x's shape.
In Boot 4.1.1 the Kafka autoconfiguration classes (KafkaTemplate, ConsumerFactory,
the KafkaListenerEndpointRegistry that @KafkaListener needs) were extracted out of the
monolithic spring-boot-autoconfigure jar into their own module, org.springframework.boot
:spring-boot-kafka. Depending on spring-kafka alone now gets you the annotations and
KafkaTemplate class with none of the autoconfiguration wiring: every @KafkaListener bean
is created but no container ever starts, silently, with no error at context startup.
This starter is what actually pulls spring-boot-kafka (confirmed via its real POM on
Maven Central) in on top of spring-kafka itself. -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-kafka</artifactId>
</dependency>
<!-- No explicit <awaitility> dependency anywhere in this file. spring-boot-starter-test
4.1.1 already declares org.awaitility:awaitility:4.3.0 as a compile-scope dependency
of its own, confirmed against the real POM on Maven Central, not assumed. Every
project using this starter already has Awaitility on its test classpath today. -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!-- The test-scope sibling of spring-boot-starter-kafka above: pulls spring-kafka-test
(the @EmbeddedKafka JUnit 5 extension and KafkaTestUtils) on the same Boot-4.1.1-managed
version as everything else here. @EmbeddedKafka starts a real, in-process Kafka broker
for the duration of a test class, no Docker, no Testcontainers, no network port opened
outside this JVM, which is also why it is the only honest way to demonstrate Kafka
listener testing in a sandbox with no Docker daemon. -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-kafka-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>3.6.0</version>
</plugin>
</plugins>
</build>
</project>
+13
View File
@@ -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)"
@@ -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 &mdash; 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);
}
}
}
@@ -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.
*
* <p>{@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 &mdash; {@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);
}
}
@@ -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
* &mdash; a connection pool, a cache still loading, a client still completing its first
* handshake &mdash; 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<Instant> 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";
}
}
@@ -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" &mdash; 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<String, String> 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();
}
}
@@ -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 &mdash; 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.
*
* <p>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<String> received = new ConcurrentLinkedQueue<>();
@KafkaListener(topics = "orders", groupId = "awaitility-demo")
public void onOrderEvent(String payload) {
received.add(payload);
}
public ConcurrentLinkedQueue<String> received() {
return received;
}
}
@@ -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();
}
}
@@ -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
@@ -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 &mdash; repeatedly, until it's there.
*
* <p>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")));
}
}
@@ -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);
}
}
}
@@ -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
* &mdash; 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<Long> 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<Long> 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
}
}
}
@@ -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 <em>not</em> treat an exception thrown inside the polled condition
* as "still false" &mdash; {@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));
}
}
@@ -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 &mdash; 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<String, String> 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<String, String> 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<String, String> newProducer() {
Map<String, Object> props = KafkaTestUtils.producerProps(embeddedKafkaBroker);
props.put(ProducerConfig.ACKS_CONFIG, "all");
ProducerFactory<String, String> producerFactory = new DefaultKafkaProducerFactory<>(props);
return new KafkaTemplate<>(producerFactory);
}
}
@@ -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 &mdash; 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 <em>at least three more runs
* than there were when the test started</em>, 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));
}
}
@@ -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 <em>itself</em> 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 &mdash; 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<Thread.UncaughtExceptionHandler> 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));
}
}
@@ -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