From 2007e92257548627ecbba3459b6765d1142496f8 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Mon, 7 Sep 2026 12:07:40 +0800 Subject: [PATCH] [fix][broker]Incorrect source side topic permission deletion when disabling replication --- .../service/persistent/PersistentTopic.java | 8 +++- .../OneWayReplicatorUsingGlobalZKTest.java | 38 +++++++++++++++++++ 2 files changed, 45 insertions(+), 1 deletion(-) 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 e8c3bd7ca89e7..0817776fb17db 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 @@ -1680,7 +1680,13 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, closeClientFuture.thenAccept(__ -> { CompletableFuture deleteTopicAuthenticationFuture = new CompletableFuture<>(); - brokerService.deleteTopicAuthenticationWithRetry(topic, deleteTopicAuthenticationFuture, 5); + checkAllowedCluster(brokerService.getPulsar().getConfig().getClusterName()).thenAccept(allow -> { + if (!allow) { + deleteTopicAuthenticationFuture.complete(null); + } else { + brokerService.deleteTopicAuthenticationWithRetry(topic, deleteTopicAuthenticationFuture, 5); + } + }); deleteTopicAuthenticationFuture.thenCompose(ignore -> deleteSchema()) .thenCompose(ignore -> deleteTopicPolicies()) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorUsingGlobalZKTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorUsingGlobalZKTest.java index 1d19f361e9b66..53bea469d366d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorUsingGlobalZKTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/OneWayReplicatorUsingGlobalZKTest.java @@ -20,6 +20,7 @@ import static org.apache.pulsar.broker.service.TopicPoliciesService.GetType.GLOBAL_ONLY; import static org.apache.pulsar.broker.service.TopicPoliciesService.GetType.LOCAL_ONLY; +import static org.assertj.core.api.Assertions.assertThat; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; @@ -28,6 +29,7 @@ import java.time.Duration; import java.util.Arrays; import java.util.Collections; +import java.util.EnumSet; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -51,6 +53,7 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.AutoFailoverPolicyData; import org.apache.pulsar.common.policies.data.AutoFailoverPolicyType; import org.apache.pulsar.common.policies.data.AutoTopicCreationOverride; @@ -120,6 +123,41 @@ public void testReplicatorCreateTopic(boolean isPartitioned) throws Exception { super.testReplicatorCreateTopic(isPartitioned); } + @Test + public void testDeletingTopicOnOneClusterKeepsSharedTopicPermissions() throws Exception { + String topic = BrokerTestUtil.newUniqueName("persistent://" + replicatedNamespace + "/topic-permissions-"); + String role = "test-role"; + Set permissions = EnumSet.of(AuthAction.produce, AuthAction.consume); + + admin1.topics().createNonPartitionedTopic(topic); + waitReplicatorStarted(topic); + admin1.topics().grantPermission(topic, role, permissions); + Awaitility.await().untilAsserted(() -> { + assertThat(admin1.topics().getPermissions(topic)).containsEntry(role, permissions); + assertThat(admin2.topics().getPermissions(topic)).containsEntry(role, permissions); + }); + + admin1.topicPolicies(true).setReplicationClusters(topic, Arrays.asList(cluster1)); + Awaitility.await().untilAsserted(() -> { + assertFalse(admin2.topics().getList(replicatedNamespace).contains(topic)); + assertThat(admin1.topics().getPermissions(topic)).containsEntry(role, permissions); + }); + + admin1.topicPolicies(true).setReplicationClusters(topic, Arrays.asList(cluster1, cluster2)); + Awaitility.await().untilAsserted(() -> { + assertThat(admin1.topics().getPermissions(topic)).containsEntry(role, permissions); + assertThat(admin2.topics().getPermissions(topic)).containsEntry(role, permissions); + }); + + // clear topic. + admin1.topicPolicies(true).setReplicationClusters(topic, Arrays.asList(cluster1)); + Awaitility.await().untilAsserted(() -> { + assertFalse(admin2.topics().getList(replicatedNamespace).contains(topic)); + assertThat(admin1.topics().getPermissions(topic)).containsEntry(role, permissions); + }); + admin1.topics().delete(topic); + } + @Override @Test public void testDeleteRemoteTopicByGlobalPolicy() throws Exception {