From 4d6fce988cceb883ef127fdd34682313e1f29769 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Fri, 2 Aug 2024 18:11:29 +0530 Subject: [PATCH 01/10] [WIP] Introduce `unload` flag in `ns-isolation-policy set` call --- .../broker/admin/impl/ClustersBase.java | 50 ++++++++++++++++-- ...ApiNamespaceIsolationMultiBrokersTest.java | 52 ++++++++++++++++--- .../policies/data/NamespaceIsolationData.java | 4 ++ .../NamespaceIsolationPolicyUnloadType.java | 35 +++++++++++++ .../cli/CmdNamespaceIsolationPolicy.java | 16 +++++- .../data/NamespaceIsolationDataImpl.java | 16 +++++- 6 files changed, 160 insertions(+), 13 deletions(-) create mode 100644 pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 4fe8a01e679da..9412417207356 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -27,6 +27,7 @@ import io.swagger.annotations.ExampleProperty; import java.util.ArrayList; import java.util.Collections; +import java.util.HashSet; import java.util.LinkedHashSet; import java.util.List; import java.util.Map; @@ -34,6 +35,7 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.regex.Pattern; +import java.util.Set; import java.util.stream.Collectors; import javax.ws.rs.DELETE; import javax.ws.rs.GET; @@ -65,6 +67,7 @@ import org.apache.pulsar.common.policies.data.ClusterPoliciesImpl; import org.apache.pulsar.common.policies.data.FailureDomainImpl; import org.apache.pulsar.common.policies.data.NamespaceIsolationDataImpl; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadType; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicies; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicyImpl; import org.apache.pulsar.common.util.FutureUtil; @@ -721,10 +724,12 @@ public void setNamespaceIsolationPolicy( .setIsolationDataWithCreateAsync(cluster, (p) -> Collections.emptyMap()) .thenApply(__ -> new NamespaceIsolationPolicies())) ).thenCompose(nsIsolationPolicies -> { + NamespaceIsolationDataImpl oldPolicy = nsIsolationPolicies.getPolicies().getOrDefault(policyName, null); nsIsolationPolicies.setPolicy(policyName, policyData); return namespaceIsolationPolicies() - .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); - }).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(cluster, policyData)) + .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()) + .thenApply(__ -> oldPolicy); + }).thenCompose(oldPolicy -> filterAndUnloadMatchedNamespaceAsync(cluster, policyData, oldPolicy)) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); @@ -759,7 +764,13 @@ public void setNamespaceIsolationPolicy( * Get matched namespaces; call unload for each namespaces. */ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String cluster, - NamespaceIsolationDataImpl policyData) { + NamespaceIsolationDataImpl policyData, + NamespaceIsolationDataImpl oldPolicy) { + // exit early if none of the namespaces need to be unloaded + if (NamespaceIsolationPolicyUnloadType.none.equals(policyData.getUnload())) { + return CompletableFuture.completedFuture(null); + } + PulsarAdmin adminClient; try { adminClient = pulsar().getAdminClient(); @@ -768,6 +779,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus } // compile regex patterns once List namespacePatterns = policyData.getNamespaces().stream().map(Pattern::compile).toList(); + // TODO for 4.x, we should include both old and new namespace regex pattern for unload `all` option return adminClient.tenants().getTenantsAsync().thenCompose(tenants -> { List>> filteredNamespacesForEachTenant = tenants.stream() .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant).thenCompose(namespaces -> { @@ -793,6 +805,38 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus if (CollectionUtils.isEmpty(shouldUnloadNamespaces)) { return CompletableFuture.completedFuture(null); } + // If unload type is 'changed', we need to figure out a further subset of namespaces whose placement might + // actually have been changed. + + if (oldPolicy!= null && NamespaceIsolationPolicyUnloadType.changed.equals(policyData.getUnload())) { + // We also compare that the previous primary broker list is same as current, in case all namespaces need + // to be placed again anyway. + List oldBrokers = new ArrayList<>(policyData.getPrimary()); + oldBrokers.removeAll(policyData.getPrimary()); + + if (oldBrokers.isEmpty()) { + // list is same, so we continue finding the changed namespaces. + List oldNamespacePatterns = oldPolicy.getNamespaces().stream().map(Pattern::compile).toList(); + // We create a union regex list contains old + new regexes + Set combinedPatterns = new HashSet<>(oldNamespacePatterns); + combinedPatterns.addAll(namespacePatterns); + // We create a intersection of the old and new regexes. These won't need to be unloaded + Set commonPatterns = new HashSet<>(oldNamespacePatterns); + commonPatterns.retainAll(namespacePatterns); + + // Find the changed regexes (new - new ∩ old). TODO for 4.x, make this (new U old - new ∩ old) + combinedPatterns.removeAll(commonPatterns); + + // Now we further filter the filtered namespaces based on this combinedPatterns set + shouldUnloadNamespaces = shouldUnloadNamespaces.stream() + .filter(name -> combinedPatterns.stream() + .anyMatch(pattern -> pattern.matcher(name).matches()) + ).toList(); + + } + } + // unload type is either null or not in (changed, none), so we proceed to unload all namespaces + // TODO - default in 4.x should become `changed` List> futures = shouldUnloadNamespaces.stream() .map(namespaceName -> adminClient.namespaces().unloadAsync(namespaceName)) .collect(Collectors.toList()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java index da7d95d677af8..ac80066a46dea 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java @@ -29,11 +29,13 @@ import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.MultiBrokerBaseTest; import org.apache.pulsar.client.admin.PulsarAdmin; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.policies.data.AutoFailoverPolicyData; import org.apache.pulsar.common.policies.data.AutoFailoverPolicyType; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.NamespaceIsolationData; import org.apache.pulsar.common.policies.data.TenantInfoImpl; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadType; import org.testng.Assert; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; @@ -69,23 +71,43 @@ public void setupClusters() throws Exception { .createCluster("cluster-1", ClusterData.builder().serviceUrl(localBrokerWebService).build()); remoteAdmin.clusters() .createCluster("cluster-2", ClusterData.builder().serviceUrl(remoteBrokerWebService).build()); + setupForTenant("A"); + setupForTenant("B"); + } + + private void setupForTenant(String prefix) throws PulsarAdminException { + TenantInfoImpl tenantInfo = new TenantInfoImpl(Set.of(""), Set.of("test", "cluster-1", "cluster-2")); - localAdmin.tenants().createTenant("prop-ig", tenantInfo); - localAdmin.namespaces().createNamespace("prop-ig/ns1", Set.of("test", "cluster-1")); + localAdmin.tenants().createTenant(prefix + "prop-ig", tenantInfo); + localAdmin.namespaces().createNamespace(prefix + "prop-ig/ns1", Set.of("test", "cluster-1")); + localAdmin.namespaces().createNamespace(prefix + "prop-ig/n1", Set.of("test", "cluster-1")); + localAdmin.topics().createNonPartitionedTopic(prefix + "prop-ig/ns1/t1"); + } + + public void testNamespaceIsolationPolicyForReplNSWithAllUnload() throws Exception { + testNamespaceIsolationPolicyForReplNS("A", "policy-1", NamespaceIsolationPolicyUnloadType.all); + } + + public void testNamespaceIsolationPolicyForReplNSWithChangedUnload() throws Exception { + testNamespaceIsolationPolicyForReplNS("A", "policy-2", NamespaceIsolationPolicyUnloadType.changed); } - public void testNamespaceIsolationPolicyForReplNS() throws Exception { + public void testNamespaceIsolationPolicyForReplNSWithoutUnload() throws Exception { + testNamespaceIsolationPolicyForReplNS("B", "policy-3", NamespaceIsolationPolicyUnloadType.none); + } - // Verify that namespace is not present in cluster-2. - Set replicationClusters = localAdmin.namespaces().getPolicies("prop-ig/ns1").replication_clusters; + private void testNamespaceIsolationPolicyForReplNS(String prefix, String policyName1, NamespaceIsolationPolicyUnloadType unload) throws Exception { + // Verify that namespaces are not present in cluster-2. + Set replicationClusters = localAdmin.namespaces().getPolicies(prefix + "prop-ig/ns1").replication_clusters; + Assert.assertFalse(replicationClusters.contains("cluster-2")); + replicationClusters = localAdmin.namespaces().getPolicies(prefix + "prop-ig/n1").replication_clusters; Assert.assertFalse(replicationClusters.contains("cluster-2")); // setup ns-isolation-policy in both the clusters. - String policyName1 = "policy-1"; Map parameters1 = new HashMap<>(); parameters1.put("min_limit", "1"); parameters1.put("usage_threshold", "100"); - List nsRegexList = new ArrayList<>(Arrays.asList("prop-ig/.*")); + List nsRegexList = new ArrayList<>(Arrays.asList(prefix + "prop-ig/ns1.*")); NamespaceIsolationData nsPolicyData1 = NamespaceIsolationData.builder() // "prop-ig/ns1" is present in test cluster, policy set on test2 should work @@ -96,19 +118,35 @@ public void testNamespaceIsolationPolicyForReplNS() throws Exception { .policyType(AutoFailoverPolicyType.min_available) .parameters(parameters1) .build()) + .unload(unload) .build(); + // 1. Create policy should work in local cluster localAdmin.clusters().createNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); // verify policy is present in local cluster Map policiesMap = localAdmin.clusters().getNamespaceIsolationPolicies("test"); assertEquals(policiesMap.get(policyName1), nsPolicyData1); + // 2. Create policy should work in remote cluster remoteAdmin.clusters().createNamespaceIsolationPolicy("cluster-2", policyName1, nsPolicyData1); // verify policy is present in remote cluster policiesMap = remoteAdmin.clusters().getNamespaceIsolationPolicies("cluster-2"); assertEquals(policiesMap.get(policyName1), nsPolicyData1); + // 3. Update (add) policy should work in local cluster + nsPolicyData1.getNamespaces().add(prefix + "prop-ig/n1.*"); // this will add public/.* namespaces + localAdmin.clusters().updateNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); + // verify policy is present in local cluster + policiesMap = localAdmin.clusters().getNamespaceIsolationPolicies("test"); + assertEquals(policiesMap.get(policyName1), nsPolicyData1); + + // 4. Update (remove) policy should work in local cluster + nsPolicyData1.getNamespaces().remove(prefix + "prop-ig/n1.*"); // this will add public/.* namespaces + localAdmin.clusters().updateNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); + // verify policy is present in local cluster + policiesMap = localAdmin.clusters().getNamespaceIsolationPolicies("test"); + assertEquals(policiesMap.get(policyName1), nsPolicyData1); } } diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java index aa48e69c14571..c7d75d6856ad3 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java @@ -31,6 +31,8 @@ public interface NamespaceIsolationData { AutoFailoverPolicyData getAutoFailoverPolicy(); + NamespaceIsolationPolicyUnloadType getUnload(); + void validate(); interface Builder { @@ -42,6 +44,8 @@ interface Builder { Builder autoFailoverPolicy(AutoFailoverPolicyData autoFailoverPolicyData); + Builder unload(NamespaceIsolationPolicyUnloadType unload); + NamespaceIsolationData build(); } diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java new file mode 100644 index 0000000000000..f63192c574a09 --- /dev/null +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java @@ -0,0 +1,35 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.common.policies.data; + +/** + * The type of unload to perform while setting the isolation policy. + */ +public enum NamespaceIsolationPolicyUnloadType { + all, none, changed; + + public static NamespaceIsolationPolicyUnloadType fromString(String unloadTypeName) { + for (NamespaceIsolationPolicyUnloadType unloadType : NamespaceIsolationPolicyUnloadType.values()) { + if (unloadType.toString().equalsIgnoreCase(unloadTypeName)) { + return unloadType; + } + } + return null; + } +} diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java index e9896decd8c96..04e19bfdf1e12 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java @@ -32,6 +32,7 @@ import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationDataImpl; import org.apache.pulsar.common.policies.data.NamespaceIsolationData; import org.apache.pulsar.common.policies.data.NamespaceIsolationDataImpl; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadType; import picocli.CommandLine.Command; import picocli.CommandLine.Option; import picocli.CommandLine.Parameters; @@ -73,10 +74,18 @@ private class SetPolicy extends CliCommand { required = true, split = ",") private Map autoFailoverPolicyParams; + @Option(names = "--unload", description = "configure the type of unload to do - ['all', 'none', 'changed'] namespaces." + + " By default, all namespaces matching the namespaces regex will be unloaded and placed again. You can" + + " choose to not unload any namespace while setting this new policy by choosing `none` or choose to" + + " unload only the namespaces whose placement will actually change. If you chose 'none', you will need" + + " to manually unload the namespaces for them to be placed correctly, or wait till some namespaces get" + + " load balanced automatically based on load shedding configurations.") + private NamespaceIsolationPolicyUnloadType unload; + void run() throws PulsarAdminException { // validate and create the POJO NamespaceIsolationData namespaceIsolationData = createNamespaceIsolationData(namespaces, primary, secondary, - autoFailoverPolicyTypeName, autoFailoverPolicyParams); + autoFailoverPolicyTypeName, autoFailoverPolicyParams, unload); getAdmin().clusters().createNamespaceIsolationPolicy(clusterName, policyName, namespaceIsolationData); } @@ -167,7 +176,8 @@ private NamespaceIsolationData createNamespaceIsolationData(List namespa List primary, List secondary, String autoFailoverPolicyTypeName, - Map autoFailoverPolicyParams) { + Map autoFailoverPolicyParams, + NamespaceIsolationPolicyUnloadType unload) { // validate namespaces = validateList(namespaces); @@ -234,6 +244,8 @@ private NamespaceIsolationData createNamespaceIsolationData(List namespa throw new ParameterException("Unknown auto failover policy type specified : " + autoFailoverPolicyTypeName); } + nsIsolationDataBuilder.unload(unload); + return nsIsolationDataBuilder.build(); } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java index bdb51f63f89ed..4dd4b6a42bb48 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java @@ -75,6 +75,14 @@ public class NamespaceIsolationDataImpl implements NamespaceIsolationData { @JsonProperty("auto_failover_policy") private AutoFailoverPolicyData autoFailoverPolicy; + @ApiModelProperty( + name = "unload", + value = "The type of unload to perform while applying the new isolation policy.", + example = "'all' (default) for unloading all matching namespaces. 'none' for not unloading any namespace." + + " 'changed' for unloading only the namespaces whose placement is actually changing" + ) + private NamespaceIsolationPolicyUnloadType unload; + public static NamespaceIsolationDataImplBuilder builder() { return new NamespaceIsolationDataImplBuilder(); } @@ -106,6 +114,7 @@ public static class NamespaceIsolationDataImplBuilder implements NamespaceIsolat private List primary = new ArrayList<>(); private List secondary = new ArrayList<>(); private AutoFailoverPolicyData autoFailoverPolicy; + private NamespaceIsolationPolicyUnloadType unload; public NamespaceIsolationDataImplBuilder namespaces(List namespaces) { this.namespaces = namespaces; @@ -127,8 +136,13 @@ public NamespaceIsolationDataImplBuilder autoFailoverPolicy(AutoFailoverPolicyDa return this; } + public NamespaceIsolationDataImplBuilder unload(NamespaceIsolationPolicyUnloadType unload) { + this.unload = unload; + return this; + } + public NamespaceIsolationDataImpl build() { - return new NamespaceIsolationDataImpl(namespaces, primary, secondary, autoFailoverPolicy); + return new NamespaceIsolationDataImpl(namespaces, primary, secondary, autoFailoverPolicy, unload); } } } From 8384b271b9ee617bc1d5022ad71ca6adb8cb2992 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Mon, 12 Aug 2024 16:03:27 +0530 Subject: [PATCH 02/10] Changes based on PIP review --- .../broker/admin/impl/ClustersBase.java | 6 ++--- ...ApiNamespaceIsolationMultiBrokersTest.java | 12 +++++----- .../policies/data/NamespaceIsolationData.java | 4 ++-- ... NamespaceIsolationPolicyUnloadScope.java} | 14 ++++++----- .../cli/CmdNamespaceIsolationPolicy.java | 23 ++++++++++--------- .../policies/NamespaceIsolationPolicy.java | 6 +++++ .../data/NamespaceIsolationDataImpl.java | 11 +++++---- .../impl/NamespaceIsolationPolicyImpl.java | 8 +++++++ 8 files changed, 51 insertions(+), 33 deletions(-) rename pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/{NamespaceIsolationPolicyUnloadType.java => NamespaceIsolationPolicyUnloadScope.java} (62%) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 9412417207356..615a1178be933 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -67,7 +67,7 @@ import org.apache.pulsar.common.policies.data.ClusterPoliciesImpl; import org.apache.pulsar.common.policies.data.FailureDomainImpl; import org.apache.pulsar.common.policies.data.NamespaceIsolationDataImpl; -import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadType; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicies; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicyImpl; import org.apache.pulsar.common.util.FutureUtil; @@ -767,7 +767,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus NamespaceIsolationDataImpl policyData, NamespaceIsolationDataImpl oldPolicy) { // exit early if none of the namespaces need to be unloaded - if (NamespaceIsolationPolicyUnloadType.none.equals(policyData.getUnload())) { + if (NamespaceIsolationPolicyUnloadScope.none.equals(policyData.getUnloadScope())) { return CompletableFuture.completedFuture(null); } @@ -808,7 +808,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus // If unload type is 'changed', we need to figure out a further subset of namespaces whose placement might // actually have been changed. - if (oldPolicy!= null && NamespaceIsolationPolicyUnloadType.changed.equals(policyData.getUnload())) { + if (oldPolicy!= null && NamespaceIsolationPolicyUnloadScope.changed.equals(policyData.getUnloadScope())) { // We also compare that the previous primary broker list is same as current, in case all namespaces need // to be placed again anyway. List oldBrokers = new ArrayList<>(policyData.getPrimary()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java index ac80066a46dea..326f8cce34279 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java @@ -35,7 +35,7 @@ import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.NamespaceIsolationData; import org.apache.pulsar.common.policies.data.TenantInfoImpl; -import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadType; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope; import org.testng.Assert; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; @@ -85,18 +85,18 @@ private void setupForTenant(String prefix) throws PulsarAdminException { } public void testNamespaceIsolationPolicyForReplNSWithAllUnload() throws Exception { - testNamespaceIsolationPolicyForReplNS("A", "policy-1", NamespaceIsolationPolicyUnloadType.all); + testNamespaceIsolationPolicyForReplNS("A", "policy-1", NamespaceIsolationPolicyUnloadScope.all_matching); } public void testNamespaceIsolationPolicyForReplNSWithChangedUnload() throws Exception { - testNamespaceIsolationPolicyForReplNS("A", "policy-2", NamespaceIsolationPolicyUnloadType.changed); + testNamespaceIsolationPolicyForReplNS("A", "policy-2", NamespaceIsolationPolicyUnloadScope.changed); } public void testNamespaceIsolationPolicyForReplNSWithoutUnload() throws Exception { - testNamespaceIsolationPolicyForReplNS("B", "policy-3", NamespaceIsolationPolicyUnloadType.none); + testNamespaceIsolationPolicyForReplNS("B", "policy-3", NamespaceIsolationPolicyUnloadScope.none); } - private void testNamespaceIsolationPolicyForReplNS(String prefix, String policyName1, NamespaceIsolationPolicyUnloadType unload) throws Exception { + private void testNamespaceIsolationPolicyForReplNS(String prefix, String policyName1, NamespaceIsolationPolicyUnloadScope unload) throws Exception { // Verify that namespaces are not present in cluster-2. Set replicationClusters = localAdmin.namespaces().getPolicies(prefix + "prop-ig/ns1").replication_clusters; Assert.assertFalse(replicationClusters.contains("cluster-2")); @@ -118,7 +118,7 @@ private void testNamespaceIsolationPolicyForReplNS(String prefix, String policyN .policyType(AutoFailoverPolicyType.min_available) .parameters(parameters1) .build()) - .unload(unload) + .unloadScope(unload) .build(); // 1. Create policy should work in local cluster diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java index c7d75d6856ad3..4f367f72fda33 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationData.java @@ -31,7 +31,7 @@ public interface NamespaceIsolationData { AutoFailoverPolicyData getAutoFailoverPolicy(); - NamespaceIsolationPolicyUnloadType getUnload(); + NamespaceIsolationPolicyUnloadScope getUnloadScope(); void validate(); @@ -44,7 +44,7 @@ interface Builder { Builder autoFailoverPolicy(AutoFailoverPolicyData autoFailoverPolicyData); - Builder unload(NamespaceIsolationPolicyUnloadType unload); + Builder unloadScope(NamespaceIsolationPolicyUnloadScope unloadScope); NamespaceIsolationData build(); } diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadScope.java similarity index 62% rename from pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java rename to pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadScope.java index f63192c574a09..2edeac45630f5 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadType.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationPolicyUnloadScope.java @@ -21,13 +21,15 @@ /** * The type of unload to perform while setting the isolation policy. */ -public enum NamespaceIsolationPolicyUnloadType { - all, none, changed; +public enum NamespaceIsolationPolicyUnloadScope { + all_matching, // unloads all matching namespaces as per new regex + none, // unloads no namespaces + changed; // unloads only the namespaces which are newly added or removed from the regex list - public static NamespaceIsolationPolicyUnloadType fromString(String unloadTypeName) { - for (NamespaceIsolationPolicyUnloadType unloadType : NamespaceIsolationPolicyUnloadType.values()) { - if (unloadType.toString().equalsIgnoreCase(unloadTypeName)) { - return unloadType; + public static NamespaceIsolationPolicyUnloadScope fromString(String unloadScopeString) { + for (NamespaceIsolationPolicyUnloadScope unloadScope : NamespaceIsolationPolicyUnloadScope.values()) { + if (unloadScope.toString().equalsIgnoreCase(unloadScopeString)) { + return unloadScope; } } return null; diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java index 04e19bfdf1e12..0f5f6b211a544 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaceIsolationPolicy.java @@ -32,7 +32,7 @@ import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationDataImpl; import org.apache.pulsar.common.policies.data.NamespaceIsolationData; import org.apache.pulsar.common.policies.data.NamespaceIsolationDataImpl; -import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadType; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope; import picocli.CommandLine.Command; import picocli.CommandLine.Option; import picocli.CommandLine.Parameters; @@ -74,18 +74,19 @@ private class SetPolicy extends CliCommand { required = true, split = ",") private Map autoFailoverPolicyParams; - @Option(names = "--unload", description = "configure the type of unload to do - ['all', 'none', 'changed'] namespaces." + - " By default, all namespaces matching the namespaces regex will be unloaded and placed again. You can" + - " choose to not unload any namespace while setting this new policy by choosing `none` or choose to" + - " unload only the namespaces whose placement will actually change. If you chose 'none', you will need" + - " to manually unload the namespaces for them to be placed correctly, or wait till some namespaces get" + - " load balanced automatically based on load shedding configurations.") - private NamespaceIsolationPolicyUnloadType unload; + @Option(names = "--unload-scope", description = "configure the type of unload to do -" + + " ['all_matching', 'none', 'changed'] namespaces. By default, all namespaces matching the namespaces" + + " regex will be unloaded and placed again. You can choose to not unload any namespace while setting" + + " this new policy by choosing `none` or choose to unload only the namespaces whose placement will" + + " actually change. If you chose 'none', you will need to manually unload the namespaces for them to" + + " be placed correctly, or wait till some namespaces get load balanced automatically based on load" + + " shedding configurations.") + private NamespaceIsolationPolicyUnloadScope unloadScope; void run() throws PulsarAdminException { // validate and create the POJO NamespaceIsolationData namespaceIsolationData = createNamespaceIsolationData(namespaces, primary, secondary, - autoFailoverPolicyTypeName, autoFailoverPolicyParams, unload); + autoFailoverPolicyTypeName, autoFailoverPolicyParams, unloadScope); getAdmin().clusters().createNamespaceIsolationPolicy(clusterName, policyName, namespaceIsolationData); } @@ -177,7 +178,7 @@ private NamespaceIsolationData createNamespaceIsolationData(List namespa List secondary, String autoFailoverPolicyTypeName, Map autoFailoverPolicyParams, - NamespaceIsolationPolicyUnloadType unload) { + NamespaceIsolationPolicyUnloadScope unload) { // validate namespaces = validateList(namespaces); @@ -244,7 +245,7 @@ private NamespaceIsolationData createNamespaceIsolationData(List namespa throw new ParameterException("Unknown auto failover policy type specified : " + autoFailoverPolicyTypeName); } - nsIsolationDataBuilder.unload(unload); + nsIsolationDataBuilder.unloadScope(unload); return nsIsolationDataBuilder.build(); } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java index bd28d30d4cee9..1bdacb42e7825 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java @@ -23,6 +23,7 @@ import java.util.SortedSet; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.policies.data.BrokerStatus; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope; /** * Namespace isolation policy. @@ -43,6 +44,11 @@ public interface NamespaceIsolationPolicy { */ List getSecondaryBrokers(); + /** + * Get the unload scope for the policy set call + */ + NamespaceIsolationPolicyUnloadScope getUnloadScope(); + /** * Get the list of primary brokers for the namespace according to the policy. * diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java index 4dd4b6a42bb48..a46438b164032 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java @@ -81,7 +81,8 @@ public class NamespaceIsolationDataImpl implements NamespaceIsolationData { example = "'all' (default) for unloading all matching namespaces. 'none' for not unloading any namespace." + " 'changed' for unloading only the namespaces whose placement is actually changing" ) - private NamespaceIsolationPolicyUnloadType unload; + @JsonProperty("unload_scope") + private NamespaceIsolationPolicyUnloadScope unloadScope; public static NamespaceIsolationDataImplBuilder builder() { return new NamespaceIsolationDataImplBuilder(); @@ -114,7 +115,7 @@ public static class NamespaceIsolationDataImplBuilder implements NamespaceIsolat private List primary = new ArrayList<>(); private List secondary = new ArrayList<>(); private AutoFailoverPolicyData autoFailoverPolicy; - private NamespaceIsolationPolicyUnloadType unload; + private NamespaceIsolationPolicyUnloadScope unloadScope; public NamespaceIsolationDataImplBuilder namespaces(List namespaces) { this.namespaces = namespaces; @@ -136,13 +137,13 @@ public NamespaceIsolationDataImplBuilder autoFailoverPolicy(AutoFailoverPolicyDa return this; } - public NamespaceIsolationDataImplBuilder unload(NamespaceIsolationPolicyUnloadType unload) { - this.unload = unload; + public NamespaceIsolationDataImplBuilder unloadScope(NamespaceIsolationPolicyUnloadScope unloadScope) { + this.unloadScope = unloadScope; return this; } public NamespaceIsolationDataImpl build() { - return new NamespaceIsolationDataImpl(namespaces, primary, secondary, autoFailoverPolicy, unload); + return new NamespaceIsolationDataImpl(namespaces, primary, secondary, autoFailoverPolicy, unloadScope); } } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/impl/NamespaceIsolationPolicyImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/impl/NamespaceIsolationPolicyImpl.java index af3663869fa02..440282f29cb36 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/impl/NamespaceIsolationPolicyImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/impl/NamespaceIsolationPolicyImpl.java @@ -29,6 +29,7 @@ import org.apache.pulsar.common.policies.NamespaceIsolationPolicy; import org.apache.pulsar.common.policies.data.BrokerStatus; import org.apache.pulsar.common.policies.data.NamespaceIsolationData; +import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope; /** * Implementation of the namespace isolation policy. @@ -39,6 +40,7 @@ public class NamespaceIsolationPolicyImpl implements NamespaceIsolationPolicy { private List primary; private List secondary; private AutoFailoverPolicy autoFailoverPolicy; + private NamespaceIsolationPolicyUnloadScope unloadScope; private boolean matchNamespaces(String fqnn) { for (String nsRegex : namespaces) { @@ -64,6 +66,7 @@ public NamespaceIsolationPolicyImpl(NamespaceIsolationData policyData) { this.primary = policyData.getPrimary(); this.secondary = policyData.getSecondary(); this.autoFailoverPolicy = AutoFailoverPolicyFactory.create(policyData.getAutoFailoverPolicy()); + this.unloadScope = policyData.getUnloadScope(); } @Override @@ -76,6 +79,11 @@ public List getSecondaryBrokers() { return this.secondary; } + @Override + public NamespaceIsolationPolicyUnloadScope getUnloadScope() { + return this.unloadScope; + } + @Override public List findPrimaryBrokers(List availableBrokers, NamespaceName namespace) { if (!this.matchNamespaces(namespace.toString())) { From 976bc161ac869499aebfc8f012e9d0b45d82a323 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Mon, 12 Aug 2024 17:05:16 +0530 Subject: [PATCH 03/10] checkstyle fixes --- .../apache/pulsar/broker/admin/impl/ClustersBase.java | 10 ++++++---- .../common/policies/NamespaceIsolationPolicy.java | 2 +- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 615a1178be933..21ae837ba1175 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -33,9 +33,9 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.regex.Pattern; -import java.util.Set; import java.util.stream.Collectors; import javax.ws.rs.DELETE; import javax.ws.rs.GET; @@ -724,7 +724,8 @@ public void setNamespaceIsolationPolicy( .setIsolationDataWithCreateAsync(cluster, (p) -> Collections.emptyMap()) .thenApply(__ -> new NamespaceIsolationPolicies())) ).thenCompose(nsIsolationPolicies -> { - NamespaceIsolationDataImpl oldPolicy = nsIsolationPolicies.getPolicies().getOrDefault(policyName, null); + NamespaceIsolationDataImpl oldPolicy = nsIsolationPolicies + .getPolicies().getOrDefault(policyName, null); nsIsolationPolicies.setPolicy(policyName, policyData); return namespaceIsolationPolicies() .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()) @@ -808,7 +809,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus // If unload type is 'changed', we need to figure out a further subset of namespaces whose placement might // actually have been changed. - if (oldPolicy!= null && NamespaceIsolationPolicyUnloadScope.changed.equals(policyData.getUnloadScope())) { + if (oldPolicy != null && NamespaceIsolationPolicyUnloadScope.changed.equals(policyData.getUnloadScope())) { // We also compare that the previous primary broker list is same as current, in case all namespaces need // to be placed again anyway. List oldBrokers = new ArrayList<>(policyData.getPrimary()); @@ -816,7 +817,8 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus if (oldBrokers.isEmpty()) { // list is same, so we continue finding the changed namespaces. - List oldNamespacePatterns = oldPolicy.getNamespaces().stream().map(Pattern::compile).toList(); + List oldNamespacePatterns = oldPolicy.getNamespaces().stream() + .map(Pattern::compile).toList(); // We create a union regex list contains old + new regexes Set combinedPatterns = new HashSet<>(oldNamespacePatterns); combinedPatterns.addAll(namespacePatterns); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java index 1bdacb42e7825..52480d91eefa4 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/NamespaceIsolationPolicy.java @@ -45,7 +45,7 @@ public interface NamespaceIsolationPolicy { List getSecondaryBrokers(); /** - * Get the unload scope for the policy set call + * Get the unload scope for the policy set call. */ NamespaceIsolationPolicyUnloadScope getUnloadScope(); From e467025578208e56f7a4692d9e4b89d758d3c650 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Mon, 12 Aug 2024 17:44:57 +0530 Subject: [PATCH 04/10] minor comment name fix --- .../common/policies/data/NamespaceIsolationDataImpl.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java index a46438b164032..1e72f0e50ee05 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/NamespaceIsolationDataImpl.java @@ -76,10 +76,10 @@ public class NamespaceIsolationDataImpl implements NamespaceIsolationData { private AutoFailoverPolicyData autoFailoverPolicy; @ApiModelProperty( - name = "unload", + name = "unload_scope", value = "The type of unload to perform while applying the new isolation policy.", - example = "'all' (default) for unloading all matching namespaces. 'none' for not unloading any namespace." - + " 'changed' for unloading only the namespaces whose placement is actually changing" + example = "'all_matching' (default) for unloading all matching namespaces. 'none' for not unloading " + + "any namespace. 'changed' for unloading only the namespaces whose placement is actually changing" ) @JsonProperty("unload_scope") private NamespaceIsolationPolicyUnloadScope unloadScope; From 55c5e759bbbd085f637af600233b366f8e9c15d7 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Fri, 16 Aug 2024 16:56:28 +0530 Subject: [PATCH 05/10] Replace Pattern with String for sets --- .../broker/admin/impl/ClustersBase.java | 25 +++++++++++-------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 21ae837ba1175..7b634ac8cb578 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -780,7 +780,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus } // compile regex patterns once List namespacePatterns = policyData.getNamespaces().stream().map(Pattern::compile).toList(); - // TODO for 4.x, we should include both old and new namespace regex pattern for unload `all` option + // TODO for 4.x, we should include both old and new namespace regex pattern for unload `all_matching` option return adminClient.tenants().getTenantsAsync().thenCompose(tenants -> { List>> filteredNamespacesForEachTenant = tenants.stream() .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant).thenCompose(namespaces -> { @@ -809,6 +809,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus // If unload type is 'changed', we need to figure out a further subset of namespaces whose placement might // actually have been changed. + log.debug("Old policy: {} ; new policy: {}", oldPolicy, policyData); if (oldPolicy != null && NamespaceIsolationPolicyUnloadScope.changed.equals(policyData.getUnloadScope())) { // We also compare that the previous primary broker list is same as current, in case all namespaces need // to be placed again anyway. @@ -817,21 +818,25 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus if (oldBrokers.isEmpty()) { // list is same, so we continue finding the changed namespaces. - List oldNamespacePatterns = oldPolicy.getNamespaces().stream() - .map(Pattern::compile).toList(); + // We create a union regex list contains old + new regexes - Set combinedPatterns = new HashSet<>(oldNamespacePatterns); - combinedPatterns.addAll(namespacePatterns); + Set combinedNamespaces = new HashSet<>(oldPolicy.getNamespaces()); + combinedNamespaces.addAll(policyData.getNamespaces()); // We create a intersection of the old and new regexes. These won't need to be unloaded - Set commonPatterns = new HashSet<>(oldNamespacePatterns); - commonPatterns.retainAll(namespacePatterns); + Set commonNamespaces = new HashSet<>(oldPolicy.getNamespaces()); + commonNamespaces.retainAll(policyData.getNamespaces()); + + log.debug("combined: regexes{}; common regexes:{}", combinedNamespaces, combinedNamespaces); // Find the changed regexes (new - new ∩ old). TODO for 4.x, make this (new U old - new ∩ old) - combinedPatterns.removeAll(commonPatterns); + combinedNamespaces.removeAll(commonNamespaces); + + log.debug("changed regexes: {}", commonNamespaces); - // Now we further filter the filtered namespaces based on this combinedPatterns set + // Now we further filter the filtered namespaces based on this combinedNamespaces set shouldUnloadNamespaces = shouldUnloadNamespaces.stream() - .filter(name -> combinedPatterns.stream() + .filter(name -> combinedNamespaces.stream() + .map(Pattern::compile) .anyMatch(pattern -> pattern.matcher(name).matches()) ).toList(); From 7c8ab79aee17c3f0ee5059c11d0f86b2fd9602fe Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Fri, 23 Aug 2024 20:50:08 +0530 Subject: [PATCH 06/10] Test correct namespaces are unloaded based on the scope --- .../pulsar/broker/admin/AdminApi2Test.java | 188 ++++++++++++++++-- ...ApiNamespaceIsolationMultiBrokersTest.java | 54 +---- 2 files changed, 176 insertions(+), 66 deletions(-) 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 40e2ca8cce905..ad593ad8d3993 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 @@ -22,6 +22,7 @@ import static org.apache.commons.lang3.StringUtils.isBlank; import static org.apache.pulsar.broker.BrokerTestUtil.newUniqueName; import static org.apache.pulsar.broker.resources.LoadBalanceResources.BUNDLE_DATA_BASE_PATH; +import static org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope.*; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.spy; @@ -53,6 +54,7 @@ import java.util.TreeSet; import java.util.UUID; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import javax.ws.rs.NotAcceptableException; @@ -109,27 +111,7 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicDomain; import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.policies.data.AutoFailoverPolicyData; -import org.apache.pulsar.common.policies.data.AutoFailoverPolicyType; -import org.apache.pulsar.common.policies.data.AutoTopicCreationOverride; -import org.apache.pulsar.common.policies.data.BacklogQuota; -import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationData; -import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationDataImpl; -import org.apache.pulsar.common.policies.data.BundlesData; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.ConsumerStats; -import org.apache.pulsar.common.policies.data.EntryFilters; -import org.apache.pulsar.common.policies.data.FailureDomain; -import org.apache.pulsar.common.policies.data.NamespaceIsolationData; -import org.apache.pulsar.common.policies.data.NonPersistentTopicStats; -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.RetentionPolicies; -import org.apache.pulsar.common.policies.data.SubscriptionStats; -import org.apache.pulsar.common.policies.data.TenantInfoImpl; -import org.apache.pulsar.common.policies.data.TopicStats; -import org.apache.pulsar.common.policies.data.TopicType; +import org.apache.pulsar.common.policies.data.*; import org.apache.pulsar.common.policies.data.impl.BacklogQuotaImpl; import org.apache.pulsar.common.protocol.schema.SchemaData; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; @@ -3496,4 +3478,168 @@ public void testGetStatsIfPartitionNotExists() throws Exception { // cleanup. admin.topics().deletePartitionedTopic(partitionedTp); } + + private NamespaceIsolationData createPolicyData(NamespaceIsolationPolicyUnloadScope scope, List namespaces) { + // setup ns-isolation-policy in both the clusters. + Map parameters1 = new HashMap<>(); + parameters1.put("min_limit", "1"); + parameters1.put("usage_threshold", "100"); + List nsRegexList = new ArrayList<>(namespaces); + + return NamespaceIsolationData.builder() + // "prop-ig/ns1" is present in test cluster, policy set on test2 should work + .namespaces(nsRegexList) + .primary(Collections.singletonList(".*")) + .secondary(Collections.singletonList("")) + .autoFailoverPolicy(AutoFailoverPolicyData.builder() + .policyType(AutoFailoverPolicyType.min_available) + .parameters(parameters1) + .build()) + .unloadScope(scope) + .build(); + } + + private boolean allTopicsUnloaded(List topics) { + for (String topic : topics) { + if (pulsar.getBrokerService().getTopicReference(topic).isPresent()) { + return false; + } + } + return true; + } + + private void loadTopics(List topics) throws PulsarClientException, ExecutionException, InterruptedException { + // create a topic by creating a producer so that the topic is present on the broker + for (String topic : topics) { + Producer producer = pulsarClient.newProducer().topic(topic).create(); + producer.close(); + pulsar.getBrokerService().getTopicIfExists(topic).get(); + } + + // All namespaces are loaded onto broker. Assert that + for (String topic : topics) { + assertTrue(pulsar.getBrokerService().getTopicReference(topic).isPresent()); + } + } + + /** + * Validates that the namespace isolation policy set and update is unloading only the relevant namespaces based on + * the unload scope provided. + * + * @param topicType persistent or non persistent. + * @param policyName policy name. + * @param nsPrefix unique namespace prefix. + * @param totalNamespaces total namespaces to create. Only the end part. Each namespace also gets a topic t1. + * @param initialScope unload scope while creating the policy. + * @param initialNamespaceRegex namespace regex while creating the policy. + * @param initialLoadedNS expected namespaces to be still loaded after the policy create call. Remaining namespaces + * will be asserted to be unloaded within 20 seconds. + * @param updatedScope unload scope while updating the policy. + * @param updatedNamespaceRegex namespace regex while updating the policy. + * @param updatedLoadedNS expected namespaces to be loaded after policy update call. Remaining namespaces will be + * asserted to be unloaded within 20 seconds. + * @throws PulsarAdminException + * @throws PulsarClientException + * @throws ExecutionException + * @throws InterruptedException + */ + private void testIsolationPolicyUnloadsNSWithScope(String topicType, String policyName, String nsPrefix, + List totalNamespaces, + NamespaceIsolationPolicyUnloadScope initialScope, + List initialNamespaceRegex, List initialLoadedNS, + NamespaceIsolationPolicyUnloadScope updatedScope, + List updatedNamespaceRegex, List updatedLoadedNS) + throws PulsarAdminException, PulsarClientException, ExecutionException, InterruptedException { + + // Create all namespaces + List allTopics = new ArrayList<>(); + for (String namespacePart: totalNamespaces) { + admin.namespaces().createNamespace(nsPrefix + namespacePart, Set.of("test")); + allTopics.add(topicType + "://" + nsPrefix + namespacePart + "/t1"); + } + // Load all topics so that they are present. Assume topic t1 under each namespace + loadTopics(allTopics); + + // Create the policy + NamespaceIsolationData nsPolicyData1 = createPolicyData(initialScope, initialNamespaceRegex); + admin.clusters().createNamespaceIsolationPolicy("test", policyName, nsPolicyData1); + + List initialLoadedTopics = new ArrayList<>(); + for (String namespacePart: initialLoadedNS) { + initialLoadedTopics.add(topicType + "://" + nsPrefix + namespacePart + "/t1"); + } + + List initialUnloadedTopics = new ArrayList<>(allTopics); + initialUnloadedTopics.removeAll(initialLoadedTopics); + + // Assert that all topics (and thus ns) not under initialLoadedNS namespaces are unloaded + if (initialUnloadedTopics.isEmpty()) { + // Just wait a bit to ensure we don't miss lazy unloading of topics we expect not to unload + TimeUnit.SECONDS.sleep(5); + } else { + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .until(() -> allTopicsUnloaded(initialUnloadedTopics)); + } + // Assert that all topics under initialLoadedNS are still present + initialLoadedTopics.forEach(t -> assertTrue(pulsar.getBrokerService().getTopicReference(t).isPresent())); + + // Load the topics again + loadTopics(allTopics); + + // Update policy using updatedScope with updated namespace regex + nsPolicyData1 = createPolicyData(updatedScope, updatedNamespaceRegex); + admin.clusters().updateNamespaceIsolationPolicy("test", policyName, nsPolicyData1); + + List updatedLoadedTopics = new ArrayList<>(); + for (String namespacePart : updatedLoadedNS) { + updatedLoadedTopics.add(topicType + "://" + nsPrefix + namespacePart + "/t1"); + } + + List updatedUnloadedTopics = new ArrayList<>(allTopics); + updatedUnloadedTopics.removeAll(updatedLoadedTopics); + + // Assert that all topics (and thus ns) not under updatedLoadedNS namespaces are unloaded + if (updatedUnloadedTopics.isEmpty()) { + // Just wait a bit to ensure we don't miss lazy unloading of topics we expect not to unload + TimeUnit.SECONDS.sleep(5); + } else { + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .until(() -> allTopicsUnloaded(updatedUnloadedTopics)); + } + // Assert that all topics under updatedLoadedNS are still present + updatedLoadedTopics.forEach(t -> assertTrue(pulsar.getBrokerService().getTopicReference(t).isPresent())); + + } + + @Test(dataProvider = "topicType") + public void testIsolationPolicyUnloadsNSWithAllScope(final String topicType) throws Exception { + String nsPrefix = newUniqueName(defaultTenant + "/") + "-unload-test-"; + testIsolationPolicyUnloadsNSWithScope( + topicType, "policy-all", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), + all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), + all_matching, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("b1", "b2") + ); + } + + @Test(dataProvider = "topicType") + public void testIsolationPolicyUnloadsNSWithChangedScope(final String topicType) throws Exception { + String nsPrefix = newUniqueName(defaultTenant + "/") + "-unload-test-"; + testIsolationPolicyUnloadsNSWithScope( + topicType, "policy-changed", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), + all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), + changed, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("a1", "a2", "b1", "b2") + ); + } + + @Test(dataProvider = "topicType") + public void testIsolationPolicyUnloadsNSWithNoneScope(final String topicType) throws Exception { + String nsPrefix = newUniqueName(defaultTenant + "/") + "-unload-test-"; + testIsolationPolicyUnloadsNSWithScope( + topicType, "policy-none", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), + all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), + none, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("a1", "a2", "b1", "b2", "c1") + ); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java index 326f8cce34279..5da93e8104a35 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.broker.admin; import static org.testng.Assert.assertEquals; + import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -26,16 +27,15 @@ import java.util.List; import java.util.Map; import java.util.Set; + import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.MultiBrokerBaseTest; import org.apache.pulsar.client.admin.PulsarAdmin; -import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.policies.data.AutoFailoverPolicyData; import org.apache.pulsar.common.policies.data.AutoFailoverPolicyType; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.NamespaceIsolationData; import org.apache.pulsar.common.policies.data.TenantInfoImpl; -import org.apache.pulsar.common.policies.data.NamespaceIsolationPolicyUnloadScope; import org.testng.Assert; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; @@ -71,43 +71,23 @@ public void setupClusters() throws Exception { .createCluster("cluster-1", ClusterData.builder().serviceUrl(localBrokerWebService).build()); remoteAdmin.clusters() .createCluster("cluster-2", ClusterData.builder().serviceUrl(remoteBrokerWebService).build()); - setupForTenant("A"); - setupForTenant("B"); - } - - private void setupForTenant(String prefix) throws PulsarAdminException { - TenantInfoImpl tenantInfo = new TenantInfoImpl(Set.of(""), Set.of("test", "cluster-1", "cluster-2")); - localAdmin.tenants().createTenant(prefix + "prop-ig", tenantInfo); - localAdmin.namespaces().createNamespace(prefix + "prop-ig/ns1", Set.of("test", "cluster-1")); - localAdmin.namespaces().createNamespace(prefix + "prop-ig/n1", Set.of("test", "cluster-1")); - localAdmin.topics().createNonPartitionedTopic(prefix + "prop-ig/ns1/t1"); - } - - public void testNamespaceIsolationPolicyForReplNSWithAllUnload() throws Exception { - testNamespaceIsolationPolicyForReplNS("A", "policy-1", NamespaceIsolationPolicyUnloadScope.all_matching); - } - - public void testNamespaceIsolationPolicyForReplNSWithChangedUnload() throws Exception { - testNamespaceIsolationPolicyForReplNS("A", "policy-2", NamespaceIsolationPolicyUnloadScope.changed); + localAdmin.tenants().createTenant("prop-ig", tenantInfo); + localAdmin.namespaces().createNamespace("prop-ig/ns1", Set.of("test", "cluster-1")); } - public void testNamespaceIsolationPolicyForReplNSWithoutUnload() throws Exception { - testNamespaceIsolationPolicyForReplNS("B", "policy-3", NamespaceIsolationPolicyUnloadScope.none); - } + public void testNamespaceIsolationPolicyForReplNS() throws Exception { - private void testNamespaceIsolationPolicyForReplNS(String prefix, String policyName1, NamespaceIsolationPolicyUnloadScope unload) throws Exception { - // Verify that namespaces are not present in cluster-2. - Set replicationClusters = localAdmin.namespaces().getPolicies(prefix + "prop-ig/ns1").replication_clusters; - Assert.assertFalse(replicationClusters.contains("cluster-2")); - replicationClusters = localAdmin.namespaces().getPolicies(prefix + "prop-ig/n1").replication_clusters; + // Verify that namespace is not present in cluster-2. + Set replicationClusters = localAdmin.namespaces().getPolicies("prop-ig/ns1").replication_clusters; Assert.assertFalse(replicationClusters.contains("cluster-2")); // setup ns-isolation-policy in both the clusters. + String policyName1 = "policy-1"; Map parameters1 = new HashMap<>(); parameters1.put("min_limit", "1"); parameters1.put("usage_threshold", "100"); - List nsRegexList = new ArrayList<>(Arrays.asList(prefix + "prop-ig/ns1.*")); + List nsRegexList = new ArrayList<>(Arrays.asList("prop-ig/.*")); NamespaceIsolationData nsPolicyData1 = NamespaceIsolationData.builder() // "prop-ig/ns1" is present in test cluster, policy set on test2 should work @@ -118,35 +98,19 @@ private void testNamespaceIsolationPolicyForReplNS(String prefix, String policyN .policyType(AutoFailoverPolicyType.min_available) .parameters(parameters1) .build()) - .unloadScope(unload) .build(); - // 1. Create policy should work in local cluster localAdmin.clusters().createNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); // verify policy is present in local cluster Map policiesMap = localAdmin.clusters().getNamespaceIsolationPolicies("test"); assertEquals(policiesMap.get(policyName1), nsPolicyData1); - // 2. Create policy should work in remote cluster remoteAdmin.clusters().createNamespaceIsolationPolicy("cluster-2", policyName1, nsPolicyData1); // verify policy is present in remote cluster policiesMap = remoteAdmin.clusters().getNamespaceIsolationPolicies("cluster-2"); assertEquals(policiesMap.get(policyName1), nsPolicyData1); - // 3. Update (add) policy should work in local cluster - nsPolicyData1.getNamespaces().add(prefix + "prop-ig/n1.*"); // this will add public/.* namespaces - localAdmin.clusters().updateNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); - // verify policy is present in local cluster - policiesMap = localAdmin.clusters().getNamespaceIsolationPolicies("test"); - assertEquals(policiesMap.get(policyName1), nsPolicyData1); - - // 4. Update (remove) policy should work in local cluster - nsPolicyData1.getNamespaces().remove(prefix + "prop-ig/n1.*"); // this will add public/.* namespaces - localAdmin.clusters().updateNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); - // verify policy is present in local cluster - policiesMap = localAdmin.clusters().getNamespaceIsolationPolicies("test"); - assertEquals(policiesMap.get(policyName1), nsPolicyData1); } } From 61e02008a44b69c2ada5a60d801dd7b634608bef Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Fri, 23 Aug 2024 20:52:21 +0530 Subject: [PATCH 07/10] remove whitespace --- .../admin/AdminApiNamespaceIsolationMultiBrokersTest.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java index 5da93e8104a35..da7d95d677af8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java @@ -19,7 +19,6 @@ package org.apache.pulsar.broker.admin; import static org.testng.Assert.assertEquals; - import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -27,7 +26,6 @@ import java.util.List; import java.util.Map; import java.util.Set; - import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.MultiBrokerBaseTest; import org.apache.pulsar.client.admin.PulsarAdmin; From a799493554c8276d5ae187fd3d3db412d3f73878 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Thu, 29 Aug 2024 14:34:23 +0530 Subject: [PATCH 08/10] typo fix in log line --- .../java/org/apache/pulsar/broker/admin/impl/ClustersBase.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 7b634ac8cb578..1f2d54b43e983 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -826,7 +826,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus Set commonNamespaces = new HashSet<>(oldPolicy.getNamespaces()); commonNamespaces.retainAll(policyData.getNamespaces()); - log.debug("combined: regexes{}; common regexes:{}", combinedNamespaces, combinedNamespaces); + log.debug("combined regexes: {}; common regexes:{}", combinedNamespaces, combinedNamespaces); // Find the changed regexes (new - new ∩ old). TODO for 4.x, make this (new U old - new ∩ old) combinedNamespaces.removeAll(commonNamespaces); From e1372da72f7efc60267148440004dac53ca2f263 Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Fri, 30 Aug 2024 10:31:04 +0530 Subject: [PATCH 09/10] Update pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java Co-authored-by: Zixuan Liu --- .../java/org/apache/pulsar/broker/admin/impl/ClustersBase.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 1f2d54b43e983..8d460f1cdc6bb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -813,7 +813,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus if (oldPolicy != null && NamespaceIsolationPolicyUnloadScope.changed.equals(policyData.getUnloadScope())) { // We also compare that the previous primary broker list is same as current, in case all namespaces need // to be placed again anyway. - List oldBrokers = new ArrayList<>(policyData.getPrimary()); + List oldBrokers = new ArrayList<>(oldPolicy.getPrimary()); oldBrokers.removeAll(policyData.getPrimary()); if (oldBrokers.isEmpty()) { From 746f07233d7fbdebb390f8cd56221c10b81d582a Mon Sep 17 00:00:00 2001 From: Girish Sharma Date: Fri, 30 Aug 2024 16:40:35 +0530 Subject: [PATCH 10/10] Add test for primary broker changed case --- .../broker/admin/impl/ClustersBase.java | 5 +-- .../pulsar/broker/admin/AdminApi2Test.java | 36 ++++++++++++++----- 2 files changed, 29 insertions(+), 12 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java index 8d460f1cdc6bb..132c99ce16bec 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/ClustersBase.java @@ -813,10 +813,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus if (oldPolicy != null && NamespaceIsolationPolicyUnloadScope.changed.equals(policyData.getUnloadScope())) { // We also compare that the previous primary broker list is same as current, in case all namespaces need // to be placed again anyway. - List oldBrokers = new ArrayList<>(oldPolicy.getPrimary()); - oldBrokers.removeAll(policyData.getPrimary()); - - if (oldBrokers.isEmpty()) { + if (CollectionUtils.isEqualCollection(oldPolicy.getPrimary(), policyData.getPrimary())) { // list is same, so we continue finding the changed namespaces. // We create a union regex list contains old + new regexes 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 ad593ad8d3993..155994c814c11 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 @@ -3479,7 +3479,9 @@ public void testGetStatsIfPartitionNotExists() throws Exception { admin.topics().deletePartitionedTopic(partitionedTp); } - private NamespaceIsolationData createPolicyData(NamespaceIsolationPolicyUnloadScope scope, List namespaces) { + private NamespaceIsolationData createPolicyData(NamespaceIsolationPolicyUnloadScope scope, List namespaces, + List primaryBrokers + ) { // setup ns-isolation-policy in both the clusters. Map parameters1 = new HashMap<>(); parameters1.put("min_limit", "1"); @@ -3489,7 +3491,7 @@ private NamespaceIsolationData createPolicyData(NamespaceIsolationPolicyUnloadSc return NamespaceIsolationData.builder() // "prop-ig/ns1" is present in test cluster, policy set on test2 should work .namespaces(nsRegexList) - .primary(Collections.singletonList(".*")) + .primary(primaryBrokers) .secondary(Collections.singletonList("")) .autoFailoverPolicy(AutoFailoverPolicyData.builder() .policyType(AutoFailoverPolicyType.min_available) @@ -3548,7 +3550,8 @@ private void testIsolationPolicyUnloadsNSWithScope(String topicType, String poli NamespaceIsolationPolicyUnloadScope initialScope, List initialNamespaceRegex, List initialLoadedNS, NamespaceIsolationPolicyUnloadScope updatedScope, - List updatedNamespaceRegex, List updatedLoadedNS) + List updatedNamespaceRegex, List updatedLoadedNS, + List updatedBrokerRegex) throws PulsarAdminException, PulsarClientException, ExecutionException, InterruptedException { // Create all namespaces @@ -3561,7 +3564,9 @@ private void testIsolationPolicyUnloadsNSWithScope(String topicType, String poli loadTopics(allTopics); // Create the policy - NamespaceIsolationData nsPolicyData1 = createPolicyData(initialScope, initialNamespaceRegex); + NamespaceIsolationData nsPolicyData1 = createPolicyData( + initialScope, initialNamespaceRegex, Collections.singletonList(".*") + ); admin.clusters().createNamespaceIsolationPolicy("test", policyName, nsPolicyData1); List initialLoadedTopics = new ArrayList<>(); @@ -3588,7 +3593,7 @@ private void testIsolationPolicyUnloadsNSWithScope(String topicType, String poli loadTopics(allTopics); // Update policy using updatedScope with updated namespace regex - nsPolicyData1 = createPolicyData(updatedScope, updatedNamespaceRegex); + nsPolicyData1 = createPolicyData(updatedScope, updatedNamespaceRegex, updatedBrokerRegex); admin.clusters().updateNamespaceIsolationPolicy("test", policyName, nsPolicyData1); List updatedLoadedTopics = new ArrayList<>(); @@ -3619,7 +3624,8 @@ public void testIsolationPolicyUnloadsNSWithAllScope(final String topicType) thr testIsolationPolicyUnloadsNSWithScope( topicType, "policy-all", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), - all_matching, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("b1", "b2") + all_matching, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("b1", "b2"), + Collections.singletonList(".*") ); } @@ -3629,7 +3635,8 @@ public void testIsolationPolicyUnloadsNSWithChangedScope(final String topicType) testIsolationPolicyUnloadsNSWithScope( topicType, "policy-changed", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), - changed, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("a1", "a2", "b1", "b2") + changed, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("a1", "a2", "b1", "b2"), + Collections.singletonList(".*") ); } @@ -3639,7 +3646,20 @@ public void testIsolationPolicyUnloadsNSWithNoneScope(final String topicType) th testIsolationPolicyUnloadsNSWithScope( topicType, "policy-none", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), - none, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("a1", "a2", "b1", "b2", "c1") + none, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("a1", "a2", "b1", "b2", "c1"), + Collections.singletonList(".*") + ); + } + + @Test(dataProvider = "topicType") + public void testIsolationPolicyUnloadsNSWithPrimaryChanged(final String topicType) throws Exception { + String nsPrefix = newUniqueName(defaultTenant + "/") + "-unload-test-"; + // As per changed flag, only c1 should unload, but due to primary change, both a* and c* will. + testIsolationPolicyUnloadsNSWithScope( + topicType, "policy-primary-changed", nsPrefix, List.of("a1", "a2", "b1", "b2", "c1"), + all_matching, List.of(".*-unload-test-a.*"), List.of("b1", "b2", "c1"), + changed, List.of(".*-unload-test-a.*", ".*-unload-test-c.*"), List.of("b1", "b2"), + List.of(".*", "broker.*") ); } }