This is an automated email from the ASF dual-hosted git repository.
merlimat pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new e0c58b0f9a9 [fix][test] Await reader reconnect after seek() to fix
flaky TopicReaderTest assertions (#26141)
e0c58b0f9a9 is described below
commit e0c58b0f9a94f7ba25551255eb61e9938f334edd
Author: Lari Hotari <[email protected]>
AuthorDate: Fri Jul 3 00:54:49 2026 +0300
[fix][test] Await reader reconnect after seek() to fix flaky
TopicReaderTest assertions (#26141)
---
.../org/apache/pulsar/client/api/TopicReaderTest.java | 16 ++++++++++------
1 file changed, 10 insertions(+), 6 deletions(-)
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java
index b6ff59fa0c8..bc2dafa299c 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/TopicReaderTest.java
@@ -56,6 +56,7 @@ import org.apache.pulsar.client.impl.TopicMessageIdImpl;
import org.apache.pulsar.client.impl.TopicMessageImpl;
import org.apache.pulsar.common.policies.data.TopicStats;
import org.apache.pulsar.common.util.RelativeTimeUtil;
+import org.awaitility.Awaitility;
import org.testng.Assert;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.DataProvider;
@@ -1291,8 +1292,9 @@ public class TopicReaderTest extends SharedPulsarBaseTest
{
testMessageOrderAndDuplicates(messageSetB, receivedMessage,
expectedMessage);
}
- // Reader should be finished
- assertTrue(reader.isConnected());
+ // Reader should be finished. seek() triggers an asynchronous
reconnect of the underlying consumer(s), so
+ // isConnected() can be transiently false right after the post-seek
reads; await it to avoid flakiness.
+ Awaitility.await().untilAsserted(() ->
assertTrue(reader.isConnected()));
assertFalse(reader.hasMessageAvailable());
assertEquals(((ReaderImpl) reader).getConsumer().numMessagesInQueue(),
0);
@@ -1337,8 +1339,9 @@ public class TopicReaderTest extends SharedPulsarBaseTest
{
Assert.assertTrue(messageSetB.add(receivedMessage), "Received
duplicate message " + receivedMessage);
}
- // Reader should be finished
- assertTrue(reader.isConnected());
+ // Reader should be finished. seek() triggers an asynchronous
reconnect of the underlying consumer(s), so
+ // isConnected() can be transiently false right after the post-seek
reads; await it to avoid flakiness.
+ Awaitility.await().untilAsserted(() ->
assertTrue(reader.isConnected()));
assertFalse(reader.hasMessageAvailable());
assertEquals(((MultiTopicsReaderImpl)
reader).getMultiTopicsConsumer().numMessagesInQueue(), 0);
@@ -1392,8 +1395,9 @@ public class TopicReaderTest extends SharedPulsarBaseTest
{
testMessageOrderAndDuplicates(messageSetB, receivedMessage,
expectedMessage);
}
- // Reader should be finished
- assertTrue(reader.isConnected());
+ // Reader should be finished. seek() triggers an asynchronous
reconnect of the underlying consumer(s), so
+ // isConnected() can be transiently false right after the post-seek
reads; await it to avoid flakiness.
+ Awaitility.await().untilAsserted(() ->
assertTrue(reader.isConnected()));
assertFalse(reader.hasMessageAvailable());
assertEquals(((ReaderImpl) reader).getConsumer().numMessagesInQueue(),
0);