[fix][broker]Namespaces can be created with may empty replication_clusters policy (#25551)
diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java
index ce950c9..429f7f5 100644
--- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java
+++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java
@@ -491,7 +491,7 @@
// There are still more than one clusters configured for the global namespace
throw new RestException(Status.PRECONDITION_FAILED,
"Cannot delete the global namespace " + nsName + ". There are still more than "
- + "one replication clusters configured.");
+ + "one replication clusters configured or replication clusters is empty.");
}
if (!cluster.equals(config().getClusterName())) {
// the only replication cluster is other cluster, redirect
@@ -2392,7 +2392,7 @@
log.info(msg);
return FutureUtil.failedFuture(new RestException(Status.BAD_REQUEST, msg));
}
- pulsar().getBrokerService().setCurrentClusterAllowedIfNoClusterIsAllowed(ns, policies);
+ pulsar().getBrokerService().setCurrentClusterAllowedWhenCreating(ns, policies);
// Validate cluster names and permissions
return Stream.concat(policies.replication_clusters.stream(), policies.allowed_clusters.stream())
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 d544a78..366e19d 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
@@ -4153,14 +4153,14 @@
return nsPolicies.replication_clusters.contains(pulsar.getConfig().getClusterName());
}
- public void setCurrentClusterAllowedIfNoClusterIsAllowed(NamespaceName nsName, Policies nsPolicies) {
- if (nsPolicies.replication_clusters.contains(pulsar.getConfig().getClusterName())
- || nsPolicies.allowed_clusters.contains(pulsar.getConfig().getClusterName())) {
- return;
- }
+ public void setCurrentClusterAllowedWhenCreating(NamespaceName nsName, Policies nsPolicies) {
if (nsPolicies.replication_clusters.isEmpty()) {
nsPolicies.replication_clusters.add(pulsar.getConfig().getClusterName());
- } else {
+ }
+ if (nsPolicies.allowed_clusters.isEmpty()) {
+ return;
+ }
+ if (!nsPolicies.allowed_clusters.contains(pulsar.getConfig().getClusterName())) {
nsPolicies.allowed_clusters.add(pulsar.getConfig().getClusterName());
}
}
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java
index 8650cd6..fd323fa 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java
@@ -41,7 +41,11 @@
import com.google.common.collect.Lists;
import com.google.common.collect.Sets;
import java.lang.reflect.Field;
+import java.net.URI;
import java.net.URL;
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;
import java.time.Clock;
import java.util.ArrayList;
@@ -137,6 +141,7 @@
import org.apache.pulsar.common.policies.data.PartitionedTopicStats;
import org.apache.pulsar.common.policies.data.PersistencePolicies;
import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats;
+import org.apache.pulsar.common.policies.data.Policies;
import org.apache.pulsar.common.policies.data.RetentionPolicies;
import org.apache.pulsar.common.policies.data.SubscriptionStats;
import org.apache.pulsar.common.policies.data.TenantInfoImpl;
@@ -1739,6 +1744,36 @@
Collections.singletonList(localCluster));
}
+ @Test
+ public void testCreateNamespaceWithEmptyReplicationClustersByHttp() throws Exception {
+ String localCluster = pulsar.getConfiguration().getClusterName();
+ String namespacePart = newUniqueName("ns");
+ String namespace = defaultTenant + "/" + namespacePart;
+
+ // Create namespace with "allowed_cluster", and the param "replication_clusters" is empty.
+ HttpClient httpClient = HttpClient.newHttpClient();
+ URI adminV2Uri = URI.create(brokerUrl.toString()).resolve("/admin/v2/");
+ String namespaceRequestBody = "{\"allowed_clusters\": [\"" + localCluster + "\"]}";
+ HttpRequest createNamespaceRequest =
+ HttpRequest.newBuilder(adminV2Uri.resolve("namespaces/" + namespace))
+ .header("Content-Type", "application/json")
+ .PUT(HttpRequest.BodyPublishers.ofString(namespaceRequestBody))
+ .build();
+ HttpResponse<String> createNamespaceResponse = httpClient.send(createNamespaceRequest,
+ HttpResponse.BodyHandlers.ofString());
+ assertEquals(createNamespaceResponse.statusCode(), Status.NO_CONTENT.getStatusCode(),
+ "Failed to create namespace by HTTP: " + createNamespaceResponse.body());
+
+ // Verify: replication_clusters is not empty.
+ Awaitility.await().untilAsserted(() -> {
+ Policies policies = admin.namespaces().getPolicies(namespace);
+ assertEquals(policies.replication_clusters.size(), 1);
+ assertEquals(policies.allowed_clusters.size(), 1);
+ assertTrue(policies.replication_clusters.contains(localCluster));
+ assertTrue(policies.allowed_clusters.contains(localCluster));
+ });
+ }
+
@Test(timeOut = 30000)
public void testConsumerStatsLastTimestamp() throws PulsarClientException, PulsarAdminException,
InterruptedException {