Assert that messages are acknowledged in KafkaBrokerApiTest Checking that messages are received is not enough for asserting the full expected behaviour from the KafkaBrokerApi. After message consumption, the receiver should not receive the same message again and therefore the receiver should fail for timeout. Change-Id: I5fa3019dad90a945d8d70ea394acb3885a73edfb
diff --git a/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java b/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java index f224fe6..c2bb692 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java
@@ -111,9 +111,13 @@ public static class TestConsumer implements Consumer<EventMessage> { public final List<EventMessage> messages = new ArrayList<>(); - private final CountDownLatch lock; + private CountDownLatch lock; public TestConsumer(int numMessagesExpected) { + resetExpectedMessages(numMessagesExpected); + } + + public void resetExpectedMessages(int numMessagesExpected) { lock = new CountDownLatch(numMessagesExpected); } @@ -190,6 +194,8 @@ assertThat(testConsumer.await()).isTrue(); assertThat(testConsumer.messages).hasSize(1); assertThat(gson.toJson(testConsumer.messages.get(0))).isEqualTo(gson.toJson(testEventMessage)); + + assertNoMoreExpectedMessages(testConsumer); } @Test @@ -206,5 +212,12 @@ assertThat(testConsumer.await()).isTrue(); assertThat(testConsumer.messages).hasSize(1); assertThat(gson.toJson(testConsumer.messages.get(0))).isEqualTo(gson.toJson(testEventMessage)); + + assertNoMoreExpectedMessages(testConsumer); + } + + private void assertNoMoreExpectedMessages(TestConsumer testConsumer) { + testConsumer.resetExpectedMessages(1); + assertThat(testConsumer.await()).isFalse(); } }