This is an automated email from the ASF dual-hosted git repository.

lhotari 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 b14524eb672 [fix][broker] Forward topic policy updates after init 
failures (#26110)
b14524eb672 is described below

commit b14524eb672181124f8834a2048b8ed6a278385d
Author: sinan liu <[email protected]>
AuthorDate: Tue Jun 30 00:58:18 2026 +0800

    [fix][broker] Forward topic policy updates after init failures (#26110)
---
 .../broker/service/persistent/PersistentTopic.java | 22 ++++++++---
 .../service/persistent/PersistentTopicTest.java    | 46 ++++++++++++++++++++++
 2 files changed, 62 insertions(+), 6 deletions(-)

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 89c34e417a1..c2f656c776d 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
@@ -4903,12 +4903,22 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
                     CompletableFuture<Optional<TopicPolicies>> 
localPoliciesFuture =
                             
topicPoliciesService.getTopicPoliciesAsync(partitionedTopicName,
                                     TopicPoliciesService.GetType.LOCAL_ONLY);
-                    return 
globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> {
-                        // finally update the topic policies with the latest 
value or loaded value
-                        return CompletableFuture.runAsync(() ->
-                                
topicPolicyListener.completeInitialization(global.orElse(null), 
local.orElse(null)),
-                                getPoliciesNotifyThread());
-                    }).thenCompose(Function.identity());
+                    CompletableFuture<Void> initialPoliciesFuture =
+                            
globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> {
+                                // finally update the topic policies with the 
latest value or loaded value
+                                return CompletableFuture.runAsync(() ->
+                                        
topicPolicyListener.completeInitialization(global.orElse(null),
+                                                local.orElse(null)),
+                                        getPoliciesNotifyThread());
+                            }).thenCompose(Function.identity());
+                    return initialPoliciesFuture.exceptionallyCompose(ex ->
+                            // The topic load path logs and continues when 
initial policy loading fails. Make sure the
+                            // already-registered wrapper is not left 
buffering future live updates forever.
+                            CompletableFuture.runAsync(
+                                    () -> 
topicPolicyListener.completeInitialization(null, null),
+                                    getPoliciesNotifyThread())
+                                    .thenCompose(__ -> FutureUtil.failedFuture(
+                                            
FutureUtil.unwrapCompletionException(ex))));
                 });
     }
 
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
index 1b258ed53ad..535b0ae75ad 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java
@@ -65,6 +65,7 @@ import org.apache.bookkeeper.client.LedgerHandle;
 import org.apache.bookkeeper.mledger.AsyncCallbacks;
 import org.apache.bookkeeper.mledger.ManagedCursor;
 import org.apache.bookkeeper.mledger.ManagedLedger;
+import org.apache.bookkeeper.mledger.ManagedLedgerConfig;
 import org.apache.bookkeeper.mledger.Position;
 import org.apache.bookkeeper.mledger.PositionBound;
 import org.apache.bookkeeper.mledger.PositionFactory;
@@ -76,6 +77,7 @@ import org.apache.pulsar.broker.service.BrokerService;
 import org.apache.pulsar.broker.service.BrokerTestBase;
 import org.apache.pulsar.broker.service.Topic;
 import org.apache.pulsar.broker.service.TopicPoliciesService;
+import org.apache.pulsar.broker.service.TopicPolicyListener;
 import 
org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsClient.Metric;
 import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.client.api.Consumer;
@@ -100,6 +102,7 @@ import org.apache.pulsar.common.policies.data.TenantInfo;
 import org.apache.pulsar.common.policies.data.TopicPolicies;
 import org.apache.pulsar.common.policies.data.TopicStats;
 import org.awaitility.Awaitility;
+import org.mockito.ArgumentCaptor;
 import org.testng.Assert;
 import org.testng.annotations.AfterMethod;
 import org.testng.annotations.BeforeMethod;
@@ -860,6 +863,49 @@ public class PersistentTopicTest extends BrokerTestBase {
                 TimeUnit.MINUTES.toMillis(1));
     }
 
+    @Test
+    public void 
testTopicPolicyListenerForwardsLiveUpdatesAfterInitialLoadFailure() throws 
Exception {
+        class RecordingPersistentTopic extends PersistentTopic {
+            final List<TopicPolicies> receivedUpdates = new ArrayList<>();
+
+            RecordingPersistentTopic(String topic, ManagedLedger ledger, 
BrokerService brokerService) {
+                super(topic, ledger, brokerService);
+            }
+
+            @Override
+            public void onUpdate(TopicPolicies policies) {
+                receivedUpdates.add(policies);
+            }
+        }
+
+        final String topic = 
"persistent://prop/ns-abc/testTopicPolicyInitFailure-" + UUID.randomUUID();
+        ManagedLedger ledger = mock(ManagedLedger.class);
+        doReturn(new ManagedLedgerConfig()).when(ledger).getConfig();
+        doReturn(Collections.emptyMap()).when(ledger).getProperties();
+
+        TopicPoliciesService policiesService = 
mock(TopicPoliciesService.class);
+        doReturn(policiesService).when(pulsar).getTopicPoliciesService();
+        doReturn(CompletableFuture.completedFuture(true)).when(policiesService)
+                .registerListenerAsync(any(TopicName.class), 
any(TopicPolicyListener.class));
+        doReturn(CompletableFuture.failedFuture(new RuntimeException("initial 
topic policy load failed")))
+                
.when(policiesService).getTopicPoliciesAsync(any(TopicName.class),
+                        any(TopicPoliciesService.GetType.class));
+
+        RecordingPersistentTopic persistentTopic =
+                new RecordingPersistentTopic(topic, ledger, 
pulsar.getBrokerService());
+        persistentTopic.initTopicPolicy().handle((ignored, ex) -> null).get(3, 
TimeUnit.SECONDS);
+
+        ArgumentCaptor<TopicPolicyListener> listenerCaptor = 
ArgumentCaptor.forClass(TopicPolicyListener.class);
+        verify(policiesService).registerListenerAsync(any(TopicName.class), 
listenerCaptor.capture());
+
+        TopicPolicies livePolicies = new TopicPolicies();
+        livePolicies.setIsGlobal(false);
+        livePolicies.setMaxConsumerPerTopic(10);
+        listenerCaptor.getValue().onUpdate(livePolicies);
+
+        assertEquals(persistentTopic.receivedUpdates, 
Collections.singletonList(livePolicies));
+    }
+
     @Test
     public void testDynamicConfigurationAutoSkipNonRecoverableData() throws 
Exception {
         pulsar.getConfiguration().setAutoSkipNonRecoverableData(false);

Reply via email to