From debf0ff034c796b3892751793b9dae5d1be3af66 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Mon, 29 Jul 2024 21:53:41 +0530 Subject: [PATCH 01/13] fix ns isolation policy for repl namespaces --- .../broker/admin/impl/ClustersBase.java | 52 +++++++++++++++++-- .../apache/pulsar/broker/admin/AdminTest.java | 2 +- 2 files changed, 48 insertions(+), 6 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 2f064d7b37720..2fa26af409c30 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 @@ -36,6 +36,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import javax.ws.rs.DELETE; +import javax.ws.rs.DefaultValue; import javax.ws.rs.GET; import javax.ws.rs.POST; import javax.ws.rs.PUT; @@ -55,6 +56,7 @@ import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.PulsarAdmin; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamedEntity; import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationData; @@ -705,7 +707,9 @@ public void setNamespaceIsolationPolicy( @ApiParam(value = "The namespace isolation policy name", required = true) @PathParam("policyName") String policyName, @ApiParam(value = "The namespace isolation policy data", required = true) - NamespaceIsolationDataImpl policyData + NamespaceIsolationDataImpl policyData, + @DefaultValue("true") + @QueryParam("unloadBundles") boolean unload ) { validateSuperUserAccessAsync() .thenCompose(__ -> validatePoliciesReadOnlyAccessAsync()) @@ -723,7 +727,13 @@ public void setNamespaceIsolationPolicy( nsIsolationPolicies.setPolicy(policyName, policyData); return namespaceIsolationPolicies() .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); - }).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) + }).thenCompose(__ -> { + if (unload) { + return filterAndUnloadMatchedNamespaceAsync(cluster, policyData); + } else { + return CompletableFuture.completedFuture(null); + } + }) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); @@ -757,7 +767,8 @@ public void setNamespaceIsolationPolicy( /** * Get matched namespaces; call unload for each namespaces. */ - private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIsolationDataImpl policyData) { + private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String cluster, + NamespaceIsolationDataImpl policyData) { PulsarAdmin adminClient; try { adminClient = pulsar().getAdminClient(); @@ -770,8 +781,13 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIs .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant)); return FutureUtil.waitForAll(completableFutureStream) .thenApply(namespaces -> { - // if namespace match any policy regex, add it to ns list to be unload. + // Filter namespaces that have current cluster in their replication_clusters + // if namespace match any policy regex, add it to ns list to be unloaded. return namespaces.stream() + .filter(namespaceName -> adminClient.namespaces() + .getPoliciesAsync(namespaceName) + .thenApply(policies -> policies.replication_clusters.contains(cluster)) + .join()) .filter(namespaceName -> policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) .collect(Collectors.toList()); @@ -781,7 +797,33 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIs return CompletableFuture.completedFuture(null); } List> futures = shouldUnloadNamespaces.stream() - .map(namespaceName -> adminClient.namespaces().unloadAsync(namespaceName)) + .map(namespaceName -> { + try { + return adminClient.namespaces() + .getPolicies(namespaceName); + } catch (PulsarAdminException e) { + log.warn("[{}] Failed to get policy for {} namespace.", clientAppId(), + namespaceName, e); + throw new RuntimeException(e); + } + }) + .map(policies -> { + final List> unloadFutures = new ArrayList<>(); + List boundaries = policies.bundles.getBoundaries(); + for (int i = 0; i < boundaries.size() - 1; i++) { + String bundle = String.format("%s_%s", boundaries.get(i), boundaries.get(i + 1)); + try { + unloadFutures.add( + pulsar().getAdminClient().namespaces().unloadNamespaceBundleAsync( + namespaceName.toString(), bundle)); + } catch (PulsarServerException e) { + log.error("[{}] Failed to unload namespace {}", clientAppId(), namespaceName, + e); + throw new RestException(e); + } + } + return FutureUtil.waitForAll(unloadFutures); + }) .collect(Collectors.toList()); return FutureUtil.waitForAll(futures) .thenAccept(__ -> { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index 2894903c0d0c1..9d216014e7738 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -309,7 +309,7 @@ public void clusters() throws Exception { .build()) .build(); asyncRequests(ctx -> clusters.setNamespaceIsolationPolicy(ctx, - "use", "policy1", policyData)); + "use", "policy1", policyData, true)); asyncRequests(ctx -> clusters.getNamespaceIsolationPolicies(ctx, "use")); try { From 81ba58a40e11288592aa9ec95402906178e3bfe2 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Mon, 29 Jul 2024 22:12:20 +0530 Subject: [PATCH 02/13] Revert "fix ns isolation policy for repl namespaces" This reverts commit debf0ff034c796b3892751793b9dae5d1be3af66. --- .../broker/admin/impl/ClustersBase.java | 52 ++----------------- .../apache/pulsar/broker/admin/AdminTest.java | 2 +- 2 files changed, 6 insertions(+), 48 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 2fa26af409c30..2f064d7b37720 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 @@ -36,7 +36,6 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import javax.ws.rs.DELETE; -import javax.ws.rs.DefaultValue; import javax.ws.rs.GET; import javax.ws.rs.POST; import javax.ws.rs.PUT; @@ -56,7 +55,6 @@ import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.PulsarAdmin; -import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamedEntity; import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationData; @@ -707,9 +705,7 @@ public void setNamespaceIsolationPolicy( @ApiParam(value = "The namespace isolation policy name", required = true) @PathParam("policyName") String policyName, @ApiParam(value = "The namespace isolation policy data", required = true) - NamespaceIsolationDataImpl policyData, - @DefaultValue("true") - @QueryParam("unloadBundles") boolean unload + NamespaceIsolationDataImpl policyData ) { validateSuperUserAccessAsync() .thenCompose(__ -> validatePoliciesReadOnlyAccessAsync()) @@ -727,13 +723,7 @@ public void setNamespaceIsolationPolicy( nsIsolationPolicies.setPolicy(policyName, policyData); return namespaceIsolationPolicies() .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); - }).thenCompose(__ -> { - if (unload) { - return filterAndUnloadMatchedNamespaceAsync(cluster, policyData); - } else { - return CompletableFuture.completedFuture(null); - } - }) + }).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); @@ -767,8 +757,7 @@ public void setNamespaceIsolationPolicy( /** * Get matched namespaces; call unload for each namespaces. */ - private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String cluster, - NamespaceIsolationDataImpl policyData) { + private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIsolationDataImpl policyData) { PulsarAdmin adminClient; try { adminClient = pulsar().getAdminClient(); @@ -781,13 +770,8 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant)); return FutureUtil.waitForAll(completableFutureStream) .thenApply(namespaces -> { - // Filter namespaces that have current cluster in their replication_clusters - // if namespace match any policy regex, add it to ns list to be unloaded. + // if namespace match any policy regex, add it to ns list to be unload. return namespaces.stream() - .filter(namespaceName -> adminClient.namespaces() - .getPoliciesAsync(namespaceName) - .thenApply(policies -> policies.replication_clusters.contains(cluster)) - .join()) .filter(namespaceName -> policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) .collect(Collectors.toList()); @@ -797,33 +781,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus return CompletableFuture.completedFuture(null); } List> futures = shouldUnloadNamespaces.stream() - .map(namespaceName -> { - try { - return adminClient.namespaces() - .getPolicies(namespaceName); - } catch (PulsarAdminException e) { - log.warn("[{}] Failed to get policy for {} namespace.", clientAppId(), - namespaceName, e); - throw new RuntimeException(e); - } - }) - .map(policies -> { - final List> unloadFutures = new ArrayList<>(); - List boundaries = policies.bundles.getBoundaries(); - for (int i = 0; i < boundaries.size() - 1; i++) { - String bundle = String.format("%s_%s", boundaries.get(i), boundaries.get(i + 1)); - try { - unloadFutures.add( - pulsar().getAdminClient().namespaces().unloadNamespaceBundleAsync( - namespaceName.toString(), bundle)); - } catch (PulsarServerException e) { - log.error("[{}] Failed to unload namespace {}", clientAppId(), namespaceName, - e); - throw new RestException(e); - } - } - return FutureUtil.waitForAll(unloadFutures); - }) + .map(namespaceName -> adminClient.namespaces().unloadAsync(namespaceName)) .collect(Collectors.toList()); return FutureUtil.waitForAll(futures) .thenAccept(__ -> { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index 9d216014e7738..2894903c0d0c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -309,7 +309,7 @@ public void clusters() throws Exception { .build()) .build(); asyncRequests(ctx -> clusters.setNamespaceIsolationPolicy(ctx, - "use", "policy1", policyData, true)); + "use", "policy1", policyData)); asyncRequests(ctx -> clusters.getNamespaceIsolationPolicies(ctx, "use")); try { From 670f9a2e5505fa4f0a72e32575200ca737464d0d Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Mon, 29 Jul 2024 22:15:45 +0530 Subject: [PATCH 03/13] fix ns isolation policy for repl namespaces --- .../broker/admin/impl/ClustersBase.java | 26 +++++++++++++++---- 1 file changed, 21 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 2f064d7b37720..7472775abe2d7 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 @@ -36,6 +36,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import javax.ws.rs.DELETE; +import javax.ws.rs.DefaultValue; import javax.ws.rs.GET; import javax.ws.rs.POST; import javax.ws.rs.PUT; @@ -55,6 +56,7 @@ import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.PulsarAdmin; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamedEntity; import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationData; @@ -705,7 +707,9 @@ public void setNamespaceIsolationPolicy( @ApiParam(value = "The namespace isolation policy name", required = true) @PathParam("policyName") String policyName, @ApiParam(value = "The namespace isolation policy data", required = true) - NamespaceIsolationDataImpl policyData + NamespaceIsolationDataImpl policyData, + @DefaultValue("true") + @QueryParam("unloadBundles") boolean unload ) { validateSuperUserAccessAsync() .thenCompose(__ -> validatePoliciesReadOnlyAccessAsync()) @@ -722,8 +726,14 @@ public void setNamespaceIsolationPolicy( ).thenCompose(nsIsolationPolicies -> { nsIsolationPolicies.setPolicy(policyName, policyData); return namespaceIsolationPolicies() - .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); - }).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) + .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); + }).thenCompose(__ -> { + if (unload) { + return filterAndUnloadMatchedNamespaceAsync(cluster, policyData); + } else { + return CompletableFuture.completedFuture(null); + } + }) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); @@ -757,7 +767,8 @@ public void setNamespaceIsolationPolicy( /** * Get matched namespaces; call unload for each namespaces. */ - private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIsolationDataImpl policyData) { + private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String cluster, + NamespaceIsolationDataImpl policyData) { PulsarAdmin adminClient; try { adminClient = pulsar().getAdminClient(); @@ -770,8 +781,13 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIs .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant)); return FutureUtil.waitForAll(completableFutureStream) .thenApply(namespaces -> { - // if namespace match any policy regex, add it to ns list to be unload. + // Filter namespaces that have current cluster in their replication_clusters + // if namespace match any policy regex, add it to ns list to be unloaded. return namespaces.stream() + .filter(namespaceName -> adminClient.namespaces() + .getPoliciesAsync(namespaceName) + .thenApply(policies -> policies.replication_clusters.contains(cluster)) + .join()) .filter(namespaceName -> policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) .collect(Collectors.toList()); From 32da1002a88a8997adc0107b6b238c85240f5080 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Mon, 29 Jul 2024 22:30:12 +0530 Subject: [PATCH 04/13] fix admin tests --- .../src/test/java/org/apache/pulsar/broker/admin/AdminTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index 2894903c0d0c1..9d216014e7738 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -309,7 +309,7 @@ public void clusters() throws Exception { .build()) .build(); asyncRequests(ctx -> clusters.setNamespaceIsolationPolicy(ctx, - "use", "policy1", policyData)); + "use", "policy1", policyData, true)); asyncRequests(ctx -> clusters.getNamespaceIsolationPolicies(ctx, "use")); try { From 00819ecddc03caf623d3101905e920277aee0027 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Mon, 29 Jul 2024 22:58:19 +0530 Subject: [PATCH 05/13] remove unused import --- .../java/org/apache/pulsar/broker/admin/impl/ClustersBase.java | 1 - 1 file changed, 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 23fb1c7d445f4..36ea193a643b9 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 @@ -56,7 +56,6 @@ import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.PulsarAdmin; -import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamedEntity; import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationData; From a2e37fe304c4479ddf6c83cf6b62946b69ed90ba Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 01:38:47 +0530 Subject: [PATCH 06/13] add unload bundle flag in pulsar admin (+CLI) and change default --- .../broker/admin/impl/ClustersBase.java | 2 +- .../apache/pulsar/client/admin/Clusters.java | 37 +++++++++++++++++-- .../client/admin/internal/ClustersImpl.java | 25 +++++++------ .../cli/CmdNamespaceIsolationPolicy.java | 6 ++- 4 files changed, 52 insertions(+), 18 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 36ea193a643b9..f31fa10bac6fe 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 @@ -708,7 +708,7 @@ public void setNamespaceIsolationPolicy( @PathParam("policyName") String policyName, @ApiParam(value = "The namespace isolation policy data", required = true) NamespaceIsolationDataImpl policyData, - @DefaultValue("true") + @DefaultValue("false") @QueryParam("unloadBundles") boolean unload ) { validateSuperUserAccessAsync() diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java index 53e6680946566..3cf014862cc02 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java @@ -38,6 +38,12 @@ * Admin interface for clusters management. */ public interface Clusters { + + /** + * Defaults for all the flags. + */ + boolean UNLOAD_BUNDLE_DEFAULT = false; + /** * Get the list of clusters. *

@@ -418,9 +424,15 @@ Map getNamespaceIsolationPolicies(String cluster * Unexpected error */ void createNamespaceIsolationPolicy( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException; + default void createNamespaceIsolationPolicy( + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) + throws PulsarAdminException { + createNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); + } + /** * Create a namespace isolation policy for a cluster asynchronously. *

@@ -437,7 +449,13 @@ void createNamespaceIsolationPolicy( * @return */ CompletableFuture createNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData); + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles); + + default CompletableFuture createNamespaceIsolationPolicyAsync( + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { + return createNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); + } + /** * Returns list of active brokers with namespace-isolation policies attached to it. @@ -506,9 +524,15 @@ CompletableFuture getBrokerWithNamespaceIsolationP * Unexpected error */ void updateNamespaceIsolationPolicy( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException; + default void updateNamespaceIsolationPolicy( + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) + throws PulsarAdminException { + updateNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); + } + /** * Update a namespace isolation policy for a cluster asynchronously. *

@@ -526,7 +550,12 @@ void updateNamespaceIsolationPolicy( * */ CompletableFuture updateNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData); + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles); + + default CompletableFuture updateNamespaceIsolationPolicyAsync( + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { + return updateNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); + } /** * Delete a namespace isolation policy for a cluster. diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java index 231d4506d6173..02385ccc55c49 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java @@ -202,26 +202,26 @@ public CompletableFuture getBrokerWithNamespaceIso @Override public void createNamespaceIsolationPolicy(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData) throws PulsarAdminException { - setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData); + NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException { + setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, unloadBundles); } @Override public CompletableFuture createNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { - return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData); + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) { + return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, unloadBundles); } @Override public void updateNamespaceIsolationPolicy(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData) throws PulsarAdminException { - setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData); + NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException { + setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, unloadBundles); } @Override public CompletableFuture updateNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { - return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData); + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) { + return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, unloadBundles); } @Override @@ -236,13 +236,14 @@ public CompletableFuture deleteNamespaceIsolationPolicyAsync(String cluste } private void setNamespaceIsolationPolicy(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData) throws PulsarAdminException { - sync(() -> setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData)); + NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException { + sync(() -> setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, unloadBundles)); } private CompletableFuture setNamespaceIsolationPolicyAsync(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData) { - WebTarget path = adminClusters.path(cluster).path("namespaceIsolationPolicies").path(policyName); + NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) { + WebTarget path = adminClusters.path(cluster).path("namespaceIsolationPolicies").path(policyName) + .queryParam("unloadBundles", unloadBundles); return asyncPostRequest(path, Entity.entity(namespaceIsolationData, MediaType.APPLICATION_JSON)); } 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..8ede00294491d 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 @@ -73,12 +73,16 @@ private class SetPolicy extends CliCommand { required = true, split = ",") private Map autoFailoverPolicyParams; + @Option(names = "--unloadBundles", description = "Unload namespace bundles after applying policy") + private boolean unloadBundles; + void run() throws PulsarAdminException { // validate and create the POJO NamespaceIsolationData namespaceIsolationData = createNamespaceIsolationData(namespaces, primary, secondary, autoFailoverPolicyTypeName, autoFailoverPolicyParams); - getAdmin().clusters().createNamespaceIsolationPolicy(clusterName, policyName, namespaceIsolationData); + getAdmin().clusters() + .createNamespaceIsolationPolicy(clusterName, policyName, namespaceIsolationData, unloadBundles); } } From 70cd51d6c299071783f3400877982230867c2d26 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 13:04:52 +0530 Subject: [PATCH 07/13] Revert "add unload bundle flag in pulsar admin (+CLI) and change default" This reverts commit a2e37fe304c4479ddf6c83cf6b62946b69ed90ba. --- .../broker/admin/impl/ClustersBase.java | 2 +- .../apache/pulsar/client/admin/Clusters.java | 37 ++----------------- .../client/admin/internal/ClustersImpl.java | 25 ++++++------- .../cli/CmdNamespaceIsolationPolicy.java | 6 +-- 4 files changed, 18 insertions(+), 52 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 f31fa10bac6fe..36ea193a643b9 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 @@ -708,7 +708,7 @@ public void setNamespaceIsolationPolicy( @PathParam("policyName") String policyName, @ApiParam(value = "The namespace isolation policy data", required = true) NamespaceIsolationDataImpl policyData, - @DefaultValue("false") + @DefaultValue("true") @QueryParam("unloadBundles") boolean unload ) { validateSuperUserAccessAsync() diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java index 3cf014862cc02..53e6680946566 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java @@ -38,12 +38,6 @@ * Admin interface for clusters management. */ public interface Clusters { - - /** - * Defaults for all the flags. - */ - boolean UNLOAD_BUNDLE_DEFAULT = false; - /** * Get the list of clusters. *

@@ -424,14 +418,8 @@ Map getNamespaceIsolationPolicies(String cluster * Unexpected error */ void createNamespaceIsolationPolicy( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) - throws PulsarAdminException; - - default void createNamespaceIsolationPolicy( String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) - throws PulsarAdminException { - createNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); - } + throws PulsarAdminException; /** * Create a namespace isolation policy for a cluster asynchronously. @@ -449,13 +437,7 @@ default void createNamespaceIsolationPolicy( * @return */ CompletableFuture createNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles); - - default CompletableFuture createNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { - return createNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); - } - + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData); /** * Returns list of active brokers with namespace-isolation policies attached to it. @@ -524,14 +506,8 @@ CompletableFuture getBrokerWithNamespaceIsolationP * Unexpected error */ void updateNamespaceIsolationPolicy( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) - throws PulsarAdminException; - - default void updateNamespaceIsolationPolicy( String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) - throws PulsarAdminException { - updateNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); - } + throws PulsarAdminException; /** * Update a namespace isolation policy for a cluster asynchronously. @@ -550,12 +526,7 @@ default void updateNamespaceIsolationPolicy( * */ CompletableFuture updateNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles); - - default CompletableFuture updateNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { - return updateNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, UNLOAD_BUNDLE_DEFAULT); - } + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData); /** * Delete a namespace isolation policy for a cluster. diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java index 02385ccc55c49..231d4506d6173 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/ClustersImpl.java @@ -202,26 +202,26 @@ public CompletableFuture getBrokerWithNamespaceIso @Override public void createNamespaceIsolationPolicy(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException { - setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, unloadBundles); + NamespaceIsolationData namespaceIsolationData) throws PulsarAdminException { + setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData); } @Override public CompletableFuture createNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) { - return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, unloadBundles); + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { + return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData); } @Override public void updateNamespaceIsolationPolicy(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException { - setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData, unloadBundles); + NamespaceIsolationData namespaceIsolationData) throws PulsarAdminException { + setNamespaceIsolationPolicy(cluster, policyName, namespaceIsolationData); } @Override public CompletableFuture updateNamespaceIsolationPolicyAsync( - String cluster, String policyName, NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) { - return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, unloadBundles); + String cluster, String policyName, NamespaceIsolationData namespaceIsolationData) { + return setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData); } @Override @@ -236,14 +236,13 @@ public CompletableFuture deleteNamespaceIsolationPolicyAsync(String cluste } private void setNamespaceIsolationPolicy(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) throws PulsarAdminException { - sync(() -> setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData, unloadBundles)); + NamespaceIsolationData namespaceIsolationData) throws PulsarAdminException { + sync(() -> setNamespaceIsolationPolicyAsync(cluster, policyName, namespaceIsolationData)); } private CompletableFuture setNamespaceIsolationPolicyAsync(String cluster, String policyName, - NamespaceIsolationData namespaceIsolationData, boolean unloadBundles) { - WebTarget path = adminClusters.path(cluster).path("namespaceIsolationPolicies").path(policyName) - .queryParam("unloadBundles", unloadBundles); + NamespaceIsolationData namespaceIsolationData) { + WebTarget path = adminClusters.path(cluster).path("namespaceIsolationPolicies").path(policyName); return asyncPostRequest(path, Entity.entity(namespaceIsolationData, MediaType.APPLICATION_JSON)); } 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 8ede00294491d..e9896decd8c96 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 @@ -73,16 +73,12 @@ private class SetPolicy extends CliCommand { required = true, split = ",") private Map autoFailoverPolicyParams; - @Option(names = "--unloadBundles", description = "Unload namespace bundles after applying policy") - private boolean unloadBundles; - void run() throws PulsarAdminException { // validate and create the POJO NamespaceIsolationData namespaceIsolationData = createNamespaceIsolationData(namespaces, primary, secondary, autoFailoverPolicyTypeName, autoFailoverPolicyParams); - getAdmin().clusters() - .createNamespaceIsolationPolicy(clusterName, policyName, namespaceIsolationData, unloadBundles); + getAdmin().clusters().createNamespaceIsolationPolicy(clusterName, policyName, namespaceIsolationData); } } From 8518315f9aa8acbae07eb23151676a3233382b82 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 13:05:10 +0530 Subject: [PATCH 08/13] Revert "fix admin tests" This reverts commit 32da1002a88a8997adc0107b6b238c85240f5080. --- .../src/test/java/org/apache/pulsar/broker/admin/AdminTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index 9d216014e7738..2894903c0d0c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -309,7 +309,7 @@ public void clusters() throws Exception { .build()) .build(); asyncRequests(ctx -> clusters.setNamespaceIsolationPolicy(ctx, - "use", "policy1", policyData, true)); + "use", "policy1", policyData)); asyncRequests(ctx -> clusters.getNamespaceIsolationPolicies(ctx, "use")); try { From 7bd978f1be8ee2168c82c4a3d0f320c770bdabac Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 13:08:15 +0530 Subject: [PATCH 09/13] revert flag changes --- .../pulsar/broker/admin/impl/ClustersBase.java | 13 ++----------- 1 file changed, 2 insertions(+), 11 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 36ea193a643b9..7d3e3ec030429 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 @@ -36,7 +36,6 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import javax.ws.rs.DELETE; -import javax.ws.rs.DefaultValue; import javax.ws.rs.GET; import javax.ws.rs.POST; import javax.ws.rs.PUT; @@ -707,9 +706,7 @@ public void setNamespaceIsolationPolicy( @ApiParam(value = "The namespace isolation policy name", required = true) @PathParam("policyName") String policyName, @ApiParam(value = "The namespace isolation policy data", required = true) - NamespaceIsolationDataImpl policyData, - @DefaultValue("true") - @QueryParam("unloadBundles") boolean unload + NamespaceIsolationDataImpl policyData ) { validateSuperUserAccessAsync() .thenCompose(__ -> validatePoliciesReadOnlyAccessAsync()) @@ -727,13 +724,7 @@ public void setNamespaceIsolationPolicy( nsIsolationPolicies.setPolicy(policyName, policyData); return namespaceIsolationPolicies() .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); - }).thenCompose(__ -> { - if (unload) { - return filterAndUnloadMatchedNamespaceAsync(cluster, policyData); - } else { - return CompletableFuture.completedFuture(null); - } - }) + }).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(cluster, policyData)) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); From 22de1255ae5e5c6a6cc601e262f6b538dc3f9fdb Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 13:15:25 +0530 Subject: [PATCH 10/13] get policy post ns filter --- .../org/apache/pulsar/broker/admin/impl/ClustersBase.java | 4 ++-- 1 file changed, 2 insertions(+), 2 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 7d3e3ec030429..f94f949fb7691 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 @@ -775,12 +775,12 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus // Filter namespaces that have current cluster in their replication_clusters // if namespace match any policy regex, add it to ns list to be unloaded. return namespaces.stream() + .filter(namespaceName -> + policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) .filter(namespaceName -> adminClient.namespaces() .getPoliciesAsync(namespaceName) .thenApply(policies -> policies.replication_clusters.contains(cluster)) .join()) - .filter(namespaceName -> - policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) .collect(Collectors.toList()); }); }).thenCompose(shouldUnloadNamespaces -> { From 7a6e0d57d25677999ae70fa4e9af24dc7e7ca66b Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Tue, 30 Jul 2024 13:09:33 +0300 Subject: [PATCH 11/13] Make filterAndUnloadMatchedNamespaceAsync asynchronous --- .../broker/admin/impl/ClustersBase.java | 75 ++++++++++--------- 1 file changed, 40 insertions(+), 35 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 f94f949fb7691..4fe8a01e679da 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,8 +33,8 @@ import java.util.Objects; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.regex.Pattern; import java.util.stream.Collectors; -import java.util.stream.Stream; import javax.ws.rs.DELETE; import javax.ws.rs.GET; import javax.ws.rs.POST; @@ -766,40 +766,45 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(String clus } catch (PulsarServerException e) { return FutureUtil.failedFuture(e); } - return adminClient.tenants().getTenantsAsync() - .thenCompose(tenants -> { - Stream>> completableFutureStream = tenants.stream() - .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant)); - return FutureUtil.waitForAll(completableFutureStream) - .thenApply(namespaces -> { - // Filter namespaces that have current cluster in their replication_clusters - // if namespace match any policy regex, add it to ns list to be unloaded. - return namespaces.stream() - .filter(namespaceName -> - policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) - .filter(namespaceName -> adminClient.namespaces() - .getPoliciesAsync(namespaceName) - .thenApply(policies -> policies.replication_clusters.contains(cluster)) - .join()) - .collect(Collectors.toList()); - }); - }).thenCompose(shouldUnloadNamespaces -> { - if (CollectionUtils.isEmpty(shouldUnloadNamespaces)) { - return CompletableFuture.completedFuture(null); - } - List> futures = shouldUnloadNamespaces.stream() - .map(namespaceName -> adminClient.namespaces().unloadAsync(namespaceName)) - .collect(Collectors.toList()); - return FutureUtil.waitForAll(futures) - .thenAccept(__ -> { - try { - // write load info to load manager to make the load happens fast - pulsar().getLoadManager().get().writeLoadReportOnZookeeper(true); - } catch (Exception e) { - log.warn("[{}] Failed to writeLoadReportOnZookeeper.", clientAppId(), e); - } - }); - }); + // compile regex patterns once + List namespacePatterns = policyData.getNamespaces().stream().map(Pattern::compile).toList(); + return adminClient.tenants().getTenantsAsync().thenCompose(tenants -> { + List>> filteredNamespacesForEachTenant = tenants.stream() + .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant).thenCompose(namespaces -> { + List> namespaceNamesInCluster = namespaces.stream() + .filter(namespaceName -> namespacePatterns.stream() + .anyMatch(pattern -> pattern.matcher(namespaceName).matches())) + .map(namespaceName -> adminClient.namespaces().getPoliciesAsync(namespaceName) + .thenApply(policies -> policies.replication_clusters.contains(cluster) + ? namespaceName : null)) + .collect(Collectors.toList()); + return FutureUtil.waitForAll(namespaceNamesInCluster).thenApply( + __ -> namespaceNamesInCluster.stream() + .map(CompletableFuture::join) + .filter(Objects::nonNull) + .collect(Collectors.toList())); + })).toList(); + return FutureUtil.waitForAll(filteredNamespacesForEachTenant) + .thenApply(__ -> filteredNamespacesForEachTenant.stream() + .map(CompletableFuture::join) + .flatMap(List::stream) + .collect(Collectors.toList())); + }).thenCompose(shouldUnloadNamespaces -> { + if (CollectionUtils.isEmpty(shouldUnloadNamespaces)) { + return CompletableFuture.completedFuture(null); + } + List> futures = shouldUnloadNamespaces.stream() + .map(namespaceName -> adminClient.namespaces().unloadAsync(namespaceName)) + .collect(Collectors.toList()); + return FutureUtil.waitForAll(futures).thenAccept(__ -> { + try { + // write load info to load manager to make the load happens fast + pulsar().getLoadManager().get().writeLoadReportOnZookeeper(true); + } catch (Exception e) { + log.warn("[{}] Failed to writeLoadReportOnZookeeper.", clientAppId(), e); + } + }); + }); } @DELETE From b8746a4224c22a08069a7c11bd65d97551b62594 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 20:37:15 +0530 Subject: [PATCH 12/13] add multi broker ns isolation policy test for repl ns --- ...ApiNamespaceIsolationMultiBrokersTest.java | 115 ++++++++++++++++++ 1 file changed, 115 insertions(+) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java 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 new file mode 100644 index 0000000000000..d5b2db37daf9b --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiNamespaceIsolationMultiBrokersTest.java @@ -0,0 +1,115 @@ +/* + * 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.broker.admin; + +import static org.testng.Assert.assertEquals; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +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.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.testng.Assert; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.Test; + +/** + * Test multi-broker admin api. + */ +@Slf4j +@Test(groups = "broker-admin") +public class AdminApiNamespaceIsolationMultiBrokersTest extends MultiBrokerBaseTest { + + PulsarAdmin localAdmin; + PulsarAdmin remoteAdmin; + + @Override + protected void doInitConf() throws Exception { + super.doInitConf(); + this.conf.setManagedLedgerMaxEntriesPerLedger(10); + } + + @Override + protected void onCleanup() { + super.onCleanup(); + } + + @BeforeClass + public void setupClusters() throws Exception { + localAdmin = getAllAdmins().get(1); + remoteAdmin = getAllAdmins().get(2); + String localBrokerWebService = additionalPulsarTestContexts.get(0).getPulsarService().getWebServiceAddress(); + String remoteBrokerWebService = additionalPulsarTestContexts.get(1).getPulsarService().getWebServiceAddress(); + localAdmin.clusters() + .createCluster("cluster-1", ClusterData.builder().serviceUrl(localBrokerWebService).build()); + remoteAdmin.clusters() + .createCluster("cluster-2", ClusterData.builder().serviceUrl(remoteBrokerWebService).build()); + 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.namespaces().setNamespaceReplicationClusters("prop-ig/ns1", Set.of("cluster-1")); + } + + public void testNamespaceIsolationPolicyForReplNS() throws Exception { + + // Verify that namespace is only present in one cluster. + Set replicationClusters = localAdmin.namespaces().getPolicies("prop-ig/ns1").replication_clusters; + Assert.assertEquals(replicationClusters, Set.of("cluster-1")); + + // 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/.*")); + + NamespaceIsolationData nsPolicyData1 = 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()) + .build(); + + localAdmin.clusters().createNamespaceIsolationPolicy("test", policyName1, nsPolicyData1); + // verify policy is present in local cluster + Map policiesMap = + localAdmin.clusters().getNamespaceIsolationPolicies("test"); + assertEquals(policiesMap.get(policyName1), nsPolicyData1); + + 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); + + } + +} From 13bf97a9242120186494292f21c0774bd57fd340 Mon Sep 17 00:00:00 2001 From: iosdev747 Date: Tue, 30 Jul 2024 21:23:20 +0530 Subject: [PATCH 13/13] test cleanup --- .../admin/AdminApiNamespaceIsolationMultiBrokersTest.java | 5 ++--- 1 file changed, 2 insertions(+), 3 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 d5b2db37daf9b..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 @@ -72,14 +72,13 @@ public void setupClusters() throws Exception { 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.namespaces().setNamespaceReplicationClusters("prop-ig/ns1", Set.of("cluster-1")); } public void testNamespaceIsolationPolicyForReplNS() throws Exception { - // Verify that namespace is only present in one cluster. + // Verify that namespace is not present in cluster-2. Set replicationClusters = localAdmin.namespaces().getPolicies("prop-ig/ns1").replication_clusters; - Assert.assertEquals(replicationClusters, Set.of("cluster-1")); + Assert.assertFalse(replicationClusters.contains("cluster-2")); // setup ns-isolation-policy in both the clusters. String policyName1 = "policy-1";