diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index 8377162a7feab..c7bad62171378 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -187,13 +187,13 @@ public CompletableFuture initialize() { }, getPoliciesNotifyThread()) // Load the topic's initial policies (global and local) and register the policy listener, so a // non-persistent topic applies its own policies on load, the same as a persistent topic does. - .thenCompose(ignore -> initTopicPolicy()) - // a failure to load the initial topic policies must not fail topic loading. - .exceptionally(ex -> { - log.warn().attr("topic", topic).exception(ex) - .log("Error loading topic policies during initialization. Ignoring the failure."); - return null; - }); + .thenCompose(ignore -> initTopicPolicy() + .whenComplete((__, ex) -> { + if (ex != null) { + log.error().attr("topic", topic).exception(ex) + .log("Failed to initialize topic policies"); + } + })); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 3c4a404c32c12..b9e4b15da4a04 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -505,17 +505,17 @@ public CompletableFuture initialize() { isAllowAutoUpdateSchema = policies.is_allow_auto_update_schema; isAllowAutoUpdateSchemaWithReplicator = policies.is_allow_auto_update_schema_with_replicator; }, getPoliciesNotifyThread()) - .thenCompose(ignore -> initTopicPolicy()) - .thenCompose(ignore -> removeOrphanReplicationCursors()) - .exceptionally(ex -> { - log.warn() - .attr("topic", topic) - .exceptionMessage(ex) - .log("Error loading topic policies during initialization. Ignoring the failure. " - + "isEncryptionRequired will be set to false."); - isEncryptionRequired = false; - return null; - })); + .thenCompose(ignore -> initTopicPolicy() + .whenComplete((ignoredPolicy, ex) -> { + if (ex != null) { + log.error() + .attr("topic", topic) + .exceptionMessage(ex) + .log("Failed to initialize topic policies"); + isEncryptionRequired = false; + } + })) + .thenCompose(ignore -> removeOrphanReplicationCursors())); } private void initializeDispatchRateLimiterIfNeeded() { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index 141e3b72897c3..979bef565e0d4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -261,6 +261,41 @@ public void teardown() throws Exception { } } + @Test + public void testInitializeFailsWhenTopicPolicyLoadingFails() { + RuntimeException failure = + new RuntimeException("Failed to load topic policies"); + + TopicPoliciesService topicPoliciesService = + mock(TopicPoliciesService.class); + + doReturn(CompletableFuture.completedFuture(true)) + .when(topicPoliciesService) + .registerListenerAsync(any(), any()); + + doReturn(FutureUtil.failedFuture(failure)) + .when(topicPoliciesService) + .getTopicPoliciesAsync(any(), any()); + + doReturn(topicPoliciesService) + .when(pulsarTestContext.getPulsarService()) + .getTopicPoliciesService(); + + PersistentTopic topic = + new PersistentTopic( + successTopicName, + ledgerMock, + brokerService); + + ExecutionException exception = + Assert.expectThrows( + ExecutionException.class, + () -> topic.initialize() + .get(5, TimeUnit.SECONDS)); + + assertSame(exception.getCause(), failure); + } + @Test @SuppressWarnings("unchecked") public void testCreateTopic() { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java index d09f3068c2f09..1c09ba487eff5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java @@ -20,6 +20,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertSame; import static org.testng.Assert.assertThrows; import static org.testng.Assert.assertTrue; import java.lang.reflect.Field; @@ -27,6 +28,7 @@ import java.util.Optional; import java.util.UUID; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import lombok.Cleanup; import org.apache.pulsar.broker.service.AbstractTopic; @@ -34,6 +36,7 @@ import org.apache.pulsar.broker.service.PulsarCommandSender; import org.apache.pulsar.broker.service.SubscriptionOption; import org.apache.pulsar.broker.service.TransportCnx; +import org.apache.pulsar.broker.service.TopicPoliciesService; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; @@ -46,8 +49,10 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterPolicies.ClusterUrl; import org.apache.pulsar.common.policies.data.TopicStats; +import org.apache.pulsar.common.util.FutureUtil; import org.awaitility.Awaitility; import org.mockito.Mockito; +import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -67,6 +72,45 @@ protected void cleanup() throws Exception { super.internalCleanup(); } + @Test + public void testInitializeFailsWhenTopicPolicyLoadingFails() { + RuntimeException failure = + new RuntimeException("Failed to load topic policies"); + + TopicPoliciesService topicPoliciesService = + Mockito.mock(TopicPoliciesService.class); + + Mockito.doReturn(CompletableFuture.completedFuture(true)) + .when(topicPoliciesService) + .registerListenerAsync( + Mockito.any(), + Mockito.any()); + + Mockito.doReturn(FutureUtil.failedFuture(failure)) + .when(topicPoliciesService) + .getTopicPoliciesAsync( + Mockito.any(), + Mockito.any()); + + Mockito.doReturn(topicPoliciesService) + .when(pulsar) + .getTopicPoliciesService(); + + NonPersistentTopic topic = + new NonPersistentTopic( + "non-persistent://prop/ns-abc/" + + "policy-load-failure", + pulsar.getBrokerService()); + + ExecutionException exception = + Assert.expectThrows( + ExecutionException.class, + () -> topic.initialize() + .get(5, TimeUnit.SECONDS)); + + assertSame(exception.getCause(), failure); + } + @Test public void testAccumulativeStats() throws Exception { final String topicName = "non-persistent://prop/ns-abc/aTopic";