This is an automated email from the ASF dual-hosted git repository.

lhotari pushed a commit to branch branch-4.0
in repository https://gitbox.apache.org/repos/asf/pulsar.git

commit 3c4948b613ca90d3819edcfca8fab43ae7eb6f05
Author: Lari Hotari <[email protected]>
AuthorDate: Wed Jul 1 06:32:42 2026 +0300

    [fix][test] Run makeReadEntryProbFail's errorOrNot on a caller-provided 
executor (#26123)
    
    (cherry picked from commit 12d8aa9a9773f2546f58aca4678c65f6d506dd95)
---
 .../bookkeeper/mledger/impl/ManagedLedgerTest.java  | 21 ++++++++++++---------
 .../pulsar/broker/service/OneWayReplicatorTest.java |  4 +++-
 .../PersistentReplicatorInflightTaskTest.java       | 21 ++++++++++++++++++---
 3 files changed, 33 insertions(+), 13 deletions(-)

diff --git 
a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java
 
b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java
index dff78eab2f3..6d5503244db 100644
--- 
a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java
+++ 
b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java
@@ -77,6 +77,7 @@ import java.util.concurrent.ConcurrentLinkedDeque;
 import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.Executor;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.FutureTask;
@@ -187,19 +188,21 @@ public class ManagedLedgerTest extends 
MockedBookKeeperTestCase {
         ml.currentLedger = spyLedgerHandle;
     }
 
-    public static void makeReadEntryProbFail(ManagedLedgerImpl ml, 
Supplier<ManagedLedgerException> errorOrNot)
-            throws Exception {
+    public static void makeReadEntryProbFail(ManagedLedgerImpl ml, 
Supplier<ManagedLedgerException> errorOrNot,
+                                             Executor errorSupplierExecutor) 
throws Exception {
         ml.entryCache.clear();
         LedgerHandle currentLedger = ml.currentLedger;
         final LedgerHandle spyLedgerHandle = spy(currentLedger);
         doAnswer(invocation -> {
-            long ledgerId = (long) invocation.getArguments()[0];
-            long entryId = (long) invocation.getArguments()[1];
-            ManagedLedgerException mightError = errorOrNot.get();
-            if (mightError != null) {
-                return CompletableFuture.failedFuture(mightError);
-            }
-            return currentLedger.readUnconfirmedAsync(ledgerId, entryId);
+            long ledgerId = invocation.getArgument(0);
+            long entryId = invocation.getArgument(1);
+            // Evaluate errorOrNot on errorSupplierExecutor. Pass a 
single-threaded executor when errorOrNot may
+            // block (e.g. it waits on a CountDownLatch) so it doesn't block 
the calling read thread; pass
+            // MoreExecutors.directExecutor() to evaluate it inline on the 
calling thread.
+            return CompletableFuture.supplyAsync(errorOrNot, 
errorSupplierExecutor)
+                    .thenCompose(mightError -> mightError != null
+                            ? 
CompletableFuture.<LedgerEntries>failedFuture(mightError)
+                            : currentLedger.readUnconfirmedAsync(ledgerId, 
entryId));
         }).when(spyLedgerHandle).readUnconfirmedAsync(anyLong(), anyLong());
         ml.currentLedger = spyLedgerHandle;
     }
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java
index 5e6a78305da..99f608a200c 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorTest.java
@@ -35,6 +35,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
 import com.github.benmanes.caffeine.cache.AsyncLoadingCache;
 import com.google.common.collect.Multimap;
 import com.google.common.collect.Sets;
+import com.google.common.util.concurrent.MoreExecutors;
 import io.netty.channel.Channel;
 import io.netty.util.concurrent.FastThreadLocalThread;
 import java.lang.reflect.Field;
@@ -723,7 +724,8 @@ public class OneWayReplicatorTest extends 
OneWayReplicatorTestBase {
             }
             return new ManagedLedgerException.TooManyRequestsException("mocked 
error");
         };
-        ManagedLedgerTest.makeReadEntryProbFail(ml1, bkErrorOrNot);
+        // bkErrorOrNot doesn't block, so evaluate it inline on the calling 
read thread via directExecutor().
+        ManagedLedgerTest.makeReadEntryProbFail(ml1, bkErrorOrNot, 
MoreExecutors.directExecutor());
 
         // Verify: the replication will finish even though received 
ManagedLedgerException.TooManyRequestsException.
         pulsar1.getConfig().setReplicationStartAt("earliest");
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java
index 63f7abbd640..4bf96d37b8a 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentReplicatorInflightTaskTest.java
@@ -46,9 +46,12 @@ import java.util.List;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
+import lombok.Cleanup;
 import lombok.extern.slf4j.Slf4j;
 import org.apache.bookkeeper.mledger.Entry;
 import org.apache.bookkeeper.mledger.ManagedCursor;
@@ -162,6 +165,10 @@ public class PersistentReplicatorInflightTaskTest extends 
OneWayReplicatorTestBa
             PersistentTopic topic = (PersistentTopic) 
pulsar1.getBrokerService().getTopic(topicName, false)
                     .join().get();
             ManagedLedgerImpl ml = (ManagedLedgerImpl) 
topic.getManagedLedger();
+            // errorOrNot blocks on failRead, so run it on a dedicated 
single-threaded executor instead of the
+            // calling read thread; otherwise the blocked read thread would 
stall replicator.terminate().
+            @Cleanup("shutdownNow")
+            ExecutorService readFailExecutor = 
Executors.newSingleThreadExecutor();
             ManagedLedgerTest.makeReadEntryProbFail(ml, () -> {
                 readStarted.countDown();
                 try {
@@ -173,7 +180,7 @@ public class PersistentReplicatorInflightTaskTest extends 
OneWayReplicatorTestBa
                     return new ManagedLedgerException(e);
                 }
                 return new 
ManagedLedgerException.TooManyRequestsException("mocked read failure");
-            });
+            }, readFailExecutor);
 
             pulsar1.getConfig().setReplicationStartAt("earliest");
             admin1.topics().setReplicationClusters(topicName, 
Arrays.asList(cluster1, cluster2));
@@ -324,6 +331,10 @@ public class PersistentReplicatorInflightTaskTest extends 
OneWayReplicatorTestBa
             ManagedLedgerImpl ml = (ManagedLedgerImpl) 
topic.getManagedLedger();
             // Clear the entry cache and intercept ledger reads: block the 
first read so it stays in flight
             // (the in-flight task keeps entries == null), then let it and all 
later reads succeed.
+            // errorOrNot blocks the first read on releaseRead, so run it on a 
dedicated single-threaded executor
+            // instead of the calling read thread.
+            @Cleanup("shutdownNow")
+            ExecutorService readFailExecutor = 
Executors.newSingleThreadExecutor();
             ManagedLedgerTest.makeReadEntryProbFail(ml, () -> {
                 if (blockNextRead.compareAndSet(true, false)) {
                     readBlocked.countDown();
@@ -335,7 +346,7 @@ public class PersistentReplicatorInflightTaskTest extends 
OneWayReplicatorTestBa
                 }
                 // Never fail the read; it must complete successfully to reach 
the skip branch.
                 return null;
-            });
+            }, readFailExecutor);
 
             // Start replication from earliest; the first read blocks inside 
the ledger read (in flight).
             pulsar1.getConfig().setReplicationStartAt("earliest");
@@ -467,6 +478,10 @@ public class PersistentReplicatorInflightTaskTest extends 
OneWayReplicatorTestBa
             ManagedLedgerImpl ml = (ManagedLedgerImpl) 
topic.getManagedLedger();
             // Clear the entry cache and intercept ledger reads: block the 
first read so it stays in flight
             // (the in-flight task keeps entries == null), then let it and all 
later reads succeed.
+            // errorOrNot blocks the first read on releaseRead, so run it on a 
dedicated single-threaded executor
+            // instead of the calling read thread.
+            @Cleanup("shutdownNow")
+            ExecutorService readFailExecutor = 
Executors.newSingleThreadExecutor();
             ManagedLedgerTest.makeReadEntryProbFail(ml, () -> {
                 if (blockNextRead.compareAndSet(true, false)) {
                     readBlocked.countDown();
@@ -478,7 +493,7 @@ public class PersistentReplicatorInflightTaskTest extends 
OneWayReplicatorTestBa
                 }
                 // Never fail the read; it must complete successfully to reach 
the skip branch.
                 return null;
-            });
+            }, readFailExecutor);
 
             // Start replication from earliest; the first read blocks inside 
the ledger read (in flight).
             pulsar1.getConfig().setReplicationStartAt("earliest");

Reply via email to