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 3c47859cafecccd262abe6ca7c6e1b69b75b8f18 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) (cherry picked from commit aa2d1bc61f0fe82b50848ebcbb569992ef48918d) --- .../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 40df6b0b5ff..626805a7032 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; @@ -48,6 +47,7 @@ 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 javax.servlet.ServletContext; import javax.ws.rs.InternalServerErrorException; import javax.ws.rs.WebApplicationException; @@ -58,6 +58,7 @@ import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; 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.ExtPersistentTopics; import org.apache.pulsar.broker.admin.v2.NonPersistentTopics; @@ -129,7 +130,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 { @@ -172,7 +174,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(); @@ -180,10 +181,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()); @@ -199,6 +217,8 @@ public class PersistentTopicsTest extends MockedPulsarServiceBaseTest { @Override @AfterMethod(alwaysRun = true) protected void cleanup() throws Exception { + namespaceResourcesOverride = null; + listPersistentTopicsAsyncHandler = null; super.internalCleanup(); } @@ -606,8 +626,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); @@ -619,7 +639,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); @@ -653,12 +673,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))
