From ff3e3fa4d625a4ad14c7ea3ff0ab6ff4986a26c0 Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Mon, 12 Jun 2023 19:12:58 -0700 Subject: [PATCH] [fix][broker] new load balancer system topic should not be auto-created now --- .../apache/pulsar/broker/service/BrokerService.java | 12 +++++++++++- .../pulsar/broker/service/BrokerServiceTest.java | 11 ++++++++++- 2 files changed, 21 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 39635fa673554..2aba9adbb8935 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -100,6 +100,7 @@ import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.broker.intercept.ManagedLedgerInterceptorImpl; import org.apache.pulsar.broker.loadbalance.LoadManager; +import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitStateChannelImpl; import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.broker.resources.DynamicConfigurationResources; import org.apache.pulsar.broker.resources.LocalPoliciesResources; @@ -3304,10 +3305,19 @@ private CompletableFuture isAllowAutoTopicCreationAsync(final TopicName topicName.getNamespaceObject()); return CompletableFuture.completedFuture(false); } - //System topic can always be created automatically + + // ServiceUnitStateChannelImpl.TOPIC expects to be a non-partitioned-topic now. + // We don't allow the auto-creation here. + // ServiceUnitStateChannelImpl.start() is responsible to create the topic. + if (ServiceUnitStateChannelImpl.TOPIC.equals(topicName.toString())) { + return CompletableFuture.completedFuture(false); + } + + //Other system topics can be created automatically if (pulsar.getConfiguration().isSystemTopicEnabled() && isSystemTopic(topicName)) { return CompletableFuture.completedFuture(true); } + final boolean allowed; AutoTopicCreationOverride autoTopicCreationOverride = getAutoTopicCreationOverride(topicName, policies); if (autoTopicCreationOverride != null) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index d0c128eb89942..9bc78de83379f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -74,6 +74,7 @@ import org.apache.http.client.methods.HttpGet; import org.apache.http.impl.client.HttpClientBuilder; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitStateChannelImpl; import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.broker.service.BrokerServiceException.PersistenceException; import org.apache.pulsar.broker.service.persistent.PersistentTopic; @@ -1579,7 +1580,6 @@ public void testDynamicConfigurationsForceDeleteTenantAllowed() throws Exception }); } - // this test is disabled since it is flaky @Test(enabled = false) public void testBrokerStatsTopicLoadFailed() throws Exception { @@ -1657,4 +1657,13 @@ public void testBrokerStatsTopicLoadFailed() throws Exception { return flag.get(); }); } + + @Test + public void testIsSystemTopicAllowAutoTopicCreationAsync() throws Exception { + BrokerService brokerService = pulsar.getBrokerService(); + assertFalse(brokerService.isAllowAutoTopicCreationAsync( + ServiceUnitStateChannelImpl.TOPIC).get()); + assertTrue(brokerService.isAllowAutoTopicCreationAsync( + "persistent://pulsar/system/my-system-topic").get()); + } }