diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java index 09f8de818db0a..a465d7485888b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesService.java @@ -24,9 +24,11 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.BlockingDeque; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.LinkedBlockingDeque; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import javax.annotation.Nonnull; @@ -84,6 +86,8 @@ public class SystemTopicBasedTopicPoliciesService implements TopicPoliciesServic @VisibleForTesting final Map>> listeners = new ConcurrentHashMap<>(); + final Map> results = new ConcurrentHashMap<>(); + public SystemTopicBasedTopicPoliciesService(PulsarService pulsarService) { this.pulsarService = pulsarService; this.clusterName = pulsarService.getConfiguration().getClusterName(); @@ -97,7 +101,10 @@ public CompletableFuture deleteTopicPoliciesAsync(TopicName topicName) { @Override public CompletableFuture updateTopicPoliciesAsync(TopicName topicName, TopicPolicies policies) { - return sendTopicPolicyEvent(topicName, ActionType.UPDATE, policies); + CompletableFuture result = new CompletableFuture<>(); + results.computeIfAbsent(topicName, k -> new LinkedBlockingDeque<>()).add(result); + return sendTopicPolicyEvent(topicName, ActionType.UPDATE, policies) + .thenCombine(result, (__, ___) -> null); } private CompletableFuture sendTopicPolicyEvent(TopicName topicName, ActionType actionType, @@ -420,11 +427,37 @@ private void cleanCacheAndCloseReader(@Nonnull NamespaceName namespace, boolean }); } + private void completeResults(TopicName topicName) { + if (results.get(topicName) != null) { + BlockingDeque futures = results.get(topicName); + CompletableFuture future; + while (true) { + future = futures.poll(); + if (future == null) { + break; + } + future.complete(null); + } + } + } + + private void completeResults(BlockingDeque futures) { + CompletableFuture future; + while (true) { + future = futures.poll(); + if (future == null) { + break; + } + future.complete(null); + } + } + private void readMorePolicies(SystemTopicClient.Reader reader) { reader.readNextAsync() .thenAccept(msg -> { refreshTopicPoliciesCache(msg); notifyListener(msg); + completeResults(TopicName.get(TopicName.get(msg.getKey()).getPartitionedTopicName())); }) .whenComplete((__, ex) -> { if (ex == null) { @@ -627,6 +660,11 @@ public void unregisterListener(TopicName topicName, TopicPolicyListener futures = results.remove(topicName); + if (futures != null) { + completeResults(futures); + } } } return topicListeners; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java index 31b5bcb23cd98..09d72235f8424 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SystemTopicBasedTopicPoliciesServiceTest.java @@ -239,6 +239,74 @@ public void testGetPolicy() throws ExecutionException, InterruptedException, Top Assert.assertEquals(policies1, policiesGet1); } + @Test + public void testUpdateAndGetPolicy() throws ExecutionException, InterruptedException, TopicPoliciesCacheNotInitException { + // Init topic policies + TopicPolicies initPolicy = TopicPolicies.builder() + .maxConsumerPerTopic(10) + .build(); + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy).get(); + + // Wait for all topic policies updated. + Awaitility.await().untilAsserted(() -> + Assert.assertTrue(systemTopicBasedTopicPoliciesService + .getPoliciesCacheInit(TOPIC1.getNamespaceObject()))); + + // Assert broker is cache all topic policies + Awaitility.await().untilAsserted(() -> + Assert.assertEquals(systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1) + .getMaxConsumerPerTopic().intValue(), 10)); + + // Update policy for TOPIC1 + TopicPolicies policies1 = TopicPolicies.builder() + .maxProducerPerTopic(1) + .build(); + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, policies1).get(); + + // Ensure that the cache is updated before updateTopicPoliciesAsync returns. + TopicPolicies policiesGet1 = systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1); + Assert.assertEquals(policiesGet1, policies1); + } + + @Test + public void testGetPolicyAndRemoveTopic() throws ExecutionException, InterruptedException { + class TopicPolicyListenerImpl implements TopicPolicyListener { + + @Override + public void onUpdate(TopicPolicies data) { + //no op. + } + } + TopicPolicyListener listener = new TopicPolicyListenerImpl(); + systemTopicBasedTopicPoliciesService.registerListener(TOPIC1, listener); + Assert.assertTrue(systemTopicBasedTopicPoliciesService.listeners.get(TOPIC1).size() >= 1); + + // Init topic policies + TopicPolicies initPolicy = TopicPolicies.builder() + .maxConsumerPerTopic(10) + .build(); + systemTopicBasedTopicPoliciesService.updateTopicPoliciesAsync(TOPIC1, initPolicy).get(); + + // Wait for all topic policies updated. + Awaitility.await().untilAsserted(() -> + Assert.assertTrue(systemTopicBasedTopicPoliciesService + .getPoliciesCacheInit(TOPIC1.getNamespaceObject()))); + + // Assert broker cache all topic policies. + Awaitility.await().untilAsserted(() -> + Assert.assertEquals(systemTopicBasedTopicPoliciesService.getTopicPolicies(TOPIC1) + .getMaxConsumerPerTopic().intValue(), 10)); + + // Assert policy setting result queue is created. + Assert.assertTrue(systemTopicBasedTopicPoliciesService.results.get(TOPIC1) != null); + + // unregister listener and remove the policy setting result queue. + systemTopicBasedTopicPoliciesService.unregisterListener(TOPIC1, listener); + Assert.assertFalse(systemTopicBasedTopicPoliciesService.listeners.containsKey(TOPIC1)); + Assert.assertFalse(systemTopicBasedTopicPoliciesService.results.containsKey(TOPIC1)); + } + + @Test public void testCacheCleanup() throws Exception { final String topic = "persistent://" + NAMESPACE1 + "/test" + UUID.randomUUID();