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 aa2d1bc61f0 [fix][test] Fix flaky PersistentTopicsTest setup caused by
concurrent Mockito stubbing (#26083)
aa2d1bc61f0 is described below
commit aa2d1bc61f0fe82b50848ebcbb569992ef48918d
Author: void-ptr974 <[email protected]>
AuthorDate: Wed Jul 1 08:53:02 2026 +0800
[fix][test] Fix flaky PersistentTopicsTest setup caused by concurrent
Mockito stubbing (#26083)
---
.../pulsar/broker/admin/PersistentTopicsTest.java | 49 +++++++++++++++-------
1 file changed, 35 insertions(+), 14 deletions(-)
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
index b8945f63cc7..1cf3c4abd36 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
@@ -30,7 +30,6 @@ import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
-import static org.mockito.Mockito.when;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertNotNull;
import static org.testng.Assert.assertSame;
@@ -54,10 +53,12 @@ import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.function.Function;
import lombok.Cleanup;
import lombok.CustomLog;
import org.apache.bookkeeper.mledger.ManagedCursor;
import org.apache.commons.collections4.MapUtils;
+import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.pulsar.broker.BrokerTestUtil;
import org.apache.pulsar.broker.admin.v2.NonPersistentTopics;
import org.apache.pulsar.broker.admin.v2.PersistentTopics;
@@ -127,7 +128,8 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
protected Field uriField;
protected UriInfo uriInfo;
private NonPersistentTopics nonPersistentTopic;
- private NamespaceResources namespaceResources;
+ private volatile NamespaceResources namespaceResourcesOverride;
+ private volatile Function<NamespaceName, CompletableFuture<List<String>>>
listPersistentTopicsAsyncHandler;
@BeforeClass
public void initPersistentTopics() throws Exception {
@@ -161,7 +163,6 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
nonPersistentTopic = spy(NonPersistentTopics.class);
nonPersistentTopic.setServletContext(mock(ServletContext.class));
nonPersistentTopic.setPulsar(pulsar);
- namespaceResources = mock(NamespaceResources.class);
doReturn(false).when(nonPersistentTopic).isRequestHttps();
doReturn(null).when(nonPersistentTopic).originalPrincipal();
doReturn("test").when(nonPersistentTopic).clientAppId();
@@ -169,10 +170,27 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
doNothing().when(nonPersistentTopic).validateAdminAccessForTenant(this.testTenant);
doReturn(mock(AuthenticationDataHttps.class)).when(nonPersistentTopic).clientAuthData();
- PulsarResources resources =
- spy(new PulsarResources(pulsar.getLocalMetadataStore(),
pulsar.getConfigurationMetadataStore()));
- doReturn(spy(new
TopicResources(pulsar.getLocalMetadataStore()))).when(resources).getTopicResources();
- doReturn(resources).when(pulsar).getPulsarResources();
+ TopicResources topicResources =
BrokerTestUtil.spyWithoutRecordingInvocations(
+ new TopicResources(pulsar.getLocalMetadataStore()));
+ PulsarResources resources =
BrokerTestUtil.spyWithoutRecordingInvocations(
+ new PulsarResources(pulsar.getLocalMetadataStore(),
pulsar.getConfigurationMetadataStore()));
+ NamespaceResources namespaceResources =
resources.getNamespaceResources();
+ doAnswer(invocation -> {
+ NamespaceResources override = namespaceResourcesOverride;
+ return override != null ? override : namespaceResources;
+ }).when(resources).getNamespaceResources();
+ doReturn(topicResources).when(resources).getTopicResources();
+ doAnswer(invocation -> {
+ Function<NamespaceName, CompletableFuture<List<String>>> handler =
listPersistentTopicsAsyncHandler;
+ if (handler != null) {
+ CompletableFuture<List<String>> result =
handler.apply(invocation.getArgument(0));
+ if (result != null) {
+ return result;
+ }
+ }
+ return invocation.callRealMethod();
+ }).when(topicResources).listPersistentTopicsAsync(any());
+ FieldUtils.writeField(pulsar, "pulsarResources", resources, true);
admin.clusters().createCluster("use",
ClusterData.builder().serviceUrl("http://127.0.0.3:8082").build());
admin.clusters().createCluster("test",
ClusterData.builder().serviceUrl(pulsar.getWebServiceAddress()).build());
@@ -188,6 +206,8 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
@Override
@AfterMethod(alwaysRun = true)
protected void cleanup() throws Exception {
+ namespaceResourcesOverride = null;
+ listPersistentTopicsAsyncHandler = null;
super.internalCleanup();
}
@@ -598,8 +618,8 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
CompletableFuture<Optional<Policies>> policyFuture = new
CompletableFuture<>();
Policies policies = new Policies();
policyFuture.complete(Optional.of(policies));
-
when(pulsar.getPulsarResources().getNamespaceResources()).thenReturn(namespaceResources);
-
doReturn(policyFuture).when(namespaceResources).getPoliciesAsync(namespaceName);
+ namespaceResourcesOverride = mock(NamespaceResources.class);
+
doReturn(policyFuture).when(namespaceResourcesOverride).getPoliciesAsync(namespaceName);
AsyncResponse response = mock(AsyncResponse.class);
ArgumentCaptor<RestException> errCaptor =
ArgumentCaptor.forClass(RestException.class);
persistentTopics.createPartitionedTopic(response, testTenant,
testNamespace, topicName, 2, true);
@@ -611,7 +631,7 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
// Test policy not exist and return 'Namespace not found'
CompletableFuture<Optional<Policies>> policyFuture2 = new
CompletableFuture<>();
policyFuture2.complete(Optional.empty());
-
doReturn(policyFuture2).when(namespaceResources).getPoliciesAsync(namespaceName);
+
doReturn(policyFuture2).when(namespaceResourcesOverride).getPoliciesAsync(namespaceName);
response = mock(AsyncResponse.class);
errCaptor = ArgumentCaptor.forClass(RestException.class);
persistentTopics.createPartitionedTopic(response, testTenant,
testNamespace, topicName, 2, true);
@@ -646,12 +666,13 @@ public class PersistentTopicsTest extends
MockedPulsarServiceBaseTest {
final String nonPartitionTopicName2 = "special-topic-partition-123";
final String partitionedTopicName = "special-topic";
- when(pulsar.getPulsarResources().getTopicResources()
-
.listPersistentTopicsAsync(NamespaceName.get("my-tenant/my-namespace")))
- .thenReturn(CompletableFuture.completedFuture(List.of(
+ NamespaceName namespaceName =
NamespaceName.get("my-tenant/my-namespace");
+ listPersistentTopicsAsyncHandler = namespace ->
namespaceName.equals(namespace)
+ ? CompletableFuture.completedFuture(List.of(
"persistent://my-tenant/my-namespace/" +
nonPartitionTopicName1,
"persistent://my-tenant/my-namespace/" +
nonPartitionTopicName2
- )));
+ ))
+ : null;
// doReturn(ImmutableSet.of(nonPartitionTopicName1,
nonPartitionTopicName2)).when(mockZooKeeperChildrenCache)
// .get(anyString());
//
doReturn(CompletableFuture.completedFuture(ImmutableSet.of(nonPartitionTopicName1,
nonPartitionTopicName2))