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);
 

Reply via email to