From 91e3efa5f57e78cddb3c86943517b5ea37a8eda0 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Tue, 10 May 2022 17:19:47 +0800 Subject: [PATCH 01/12] Step commit --- .../broker/resources/NamespaceResources.java | 15 + .../broker/admin/impl/ClustersBase.java | 342 ++++++++---------- 2 files changed, 175 insertions(+), 182 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/NamespaceResources.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/NamespaceResources.java index c24df6c586fb4..2223e951f6628 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/NamespaceResources.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/NamespaceResources.java @@ -212,6 +212,21 @@ public void setIsolationData(String cluster, set(joinPath(BASE_CLUSTERS_PATH, cluster, NAMESPACE_ISOLATION_POLICIES), modifyFunction); } + public CompletableFuture setIsolationDataAsync(String cluster, + Function, + Map> modifyFunction) { + return setAsync(joinPath(BASE_CLUSTERS_PATH, cluster, NAMESPACE_ISOLATION_POLICIES), modifyFunction); + } + + public CompletableFuture setIsolationDataWithCreateAsync(String cluster, + Function>, + Map> + createFunction) { + return setWithCreateAsync(joinPath(BASE_CLUSTERS_PATH, cluster, NAMESPACE_ISOLATION_POLICIES), + createFunction); + } + public void setIsolationDataWithCreate(String cluster, Function>, Map> createFunction) 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 80309b2a026e6..92d85f6dab5bb 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 @@ -18,7 +18,6 @@ */ package org.apache.pulsar.broker.admin.impl; -import com.google.common.collect.Lists; import com.google.common.collect.Maps; import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiParam; @@ -33,7 +32,6 @@ import java.util.Map; import java.util.Optional; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; import javax.ws.rs.DELETE; import javax.ws.rs.GET; @@ -46,11 +44,14 @@ import javax.ws.rs.core.MediaType; import javax.ws.rs.core.Response; import javax.ws.rs.core.Response.Status; + +import org.apache.bookkeeper.common.util.JsonUtil; import org.apache.commons.collections4.CollectionUtils; +import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.resources.ClusterResources.FailureDomainResources; import org.apache.pulsar.broker.web.RestException; -import org.apache.pulsar.client.admin.Namespaces; +import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamedEntity; import org.apache.pulsar.common.policies.data.BrokerNamespaceIsolationData; @@ -63,12 +64,13 @@ import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicies; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicyImpl; import org.apache.pulsar.common.util.FutureUtil; -import org.apache.pulsar.common.util.ObjectMapperFactory; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.MetadataStoreException.NotFoundException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import static javax.ws.rs.core.Response.Status.PRECONDITION_FAILED; + public class ClustersBase extends AdminResource { @GET @@ -170,7 +172,7 @@ public void createCluster( log.error("[{}] Failed to create cluster {}", clientAppId(), cluster, ex); Throwable realCause = FutureUtil.unwrapCompletionException(ex); if (realCause instanceof IllegalArgumentException) { - asyncResponse.resume(new RestException(Status.PRECONDITION_FAILED, + asyncResponse.resume(new RestException(PRECONDITION_FAILED, "Cluster name is not valid")); return null; } @@ -278,13 +280,13 @@ private CompletableFuture innerSetPeerClusterNamesAsync(String cluster, if (CollectionUtils.isNotEmpty(peerClusterNames)) { future = FutureUtil.waitForAll(peerClusterNames.stream().map(peerCluster -> { if (cluster.equalsIgnoreCase(peerCluster)) { - return FutureUtil.failedFuture(new RestException(Status.PRECONDITION_FAILED, + return FutureUtil.failedFuture(new RestException(PRECONDITION_FAILED, cluster + " itself can't be part of peer-list")); } return clusterResources().getClusterAsync(peerCluster) .thenAccept(peerClusterOpt -> { if (!peerClusterOpt.isPresent()) { - throw new RestException(Status.PRECONDITION_FAILED, + throw new RestException(PRECONDITION_FAILED, "Peer cluster " + peerCluster + " does not exist"); } }); @@ -365,14 +367,14 @@ private CompletableFuture internalDeleteClusterAsync(String cluster) { return pulsar().getPulsarResources().getClusterResources().isClusterUsedAsync(cluster) .thenCompose(isClusterUsed -> { if (isClusterUsed) { - throw new RestException(Status.PRECONDITION_FAILED, "Cluster not empty"); + throw new RestException(PRECONDITION_FAILED, "Cluster not empty"); } // check the namespaceIsolationPolicies associated with the cluster return namespaceIsolationPolicies().getIsolationDataPoliciesAsync(cluster); }).thenCompose(nsIsolationPoliciesOpt -> { if (nsIsolationPoliciesOpt.isPresent()) { if (!nsIsolationPoliciesOpt.get().getPolicies().isEmpty()) { - throw new RestException(Status.PRECONDITION_FAILED, "Cluster not empty"); + throw new RestException(PRECONDITION_FAILED, "Cluster not empty"); } // Need to delete the isolation policies if present return namespaceIsolationPolicies().deleteIsolationDataAsync(cluster); @@ -543,51 +545,58 @@ private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( @ApiResponse(code = 412, message = "Cluster doesn't exist."), @ApiResponse(code = 500, message = "Internal server error.") }) - public BrokerNamespaceIsolationData getBrokerWithNamespaceIsolationPolicy( - @ApiParam( - value = "The cluster name", - required = true - ) + public void getBrokerWithNamespaceIsolationPolicy( + @Suspended AsyncResponse asyncResponse, + @ApiParam(value = "The cluster name", required = true) @PathParam("cluster") String cluster, - @ApiParam( - value = "The broker name (:)", - required = true, - example = "broker1:8080" - ) + @ApiParam(value = "The broker name (:)", required = true, + example = "broker1:8080") @PathParam("broker") String broker) { - validateSuperUserAccess(); - validateClusterExists(cluster); + validateSuperUserAccessAsync() + .thenCompose(__ -> validateClusterExistAsync(cluster, PRECONDITION_FAILED)) + .thenCompose(__ -> internalGetNamespaceIsolationPolicies(cluster)) + .thenApply(policies -> internalGetBrokerNsIsolationData(broker, policies)) + .thenAccept(asyncResponse::resume) + .exceptionally(ex -> { + log.error("[{}] Failed to get namespace isolation-policies {}", clientAppId(), cluster, ex); + resumeAsyncResponseExceptionally(asyncResponse, ex); + return null; + }); + } - Map nsPolicies; - try { - Optional nsPoliciesResult = namespaceIsolationPolicies() - .getIsolationDataPolicies(cluster); - if (!nsPoliciesResult.isPresent()) { - throw new RestException(Status.NOT_FOUND, "namespace-isolation policies not found for " + cluster); - } - nsPolicies = nsPoliciesResult.get().getPolicies(); - } catch (Exception e) { - log.error("[{}] Failed to get namespace isolation-policies {}", clientAppId(), cluster, e); - throw new RestException(e); + private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( + String broker, + Map policies) { + BrokerNamespaceIsolationData.Builder brokerIsolationData = + BrokerNamespaceIsolationData.builder().brokerName(broker); + if (policies == null) { + return brokerIsolationData.build(); } - BrokerNamespaceIsolationData.Builder brokerIsolationData = BrokerNamespaceIsolationData.builder() - .brokerName(broker); - if (nsPolicies != null) { - List namespaceRegexes = new ArrayList<>(); - nsPolicies.forEach((name, policyData) -> { - NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); - boolean isPrimary = nsPolicyImpl.isPrimaryBroker(broker); - if (isPrimary || nsPolicyImpl.isSecondaryBroker(broker)) { - namespaceRegexes.addAll(policyData.getNamespaces()); - brokerIsolationData.primary(isPrimary); - brokerIsolationData.policyName(name); + List namespaceRegexes = new ArrayList<>(); + policies.forEach((name, policyData) -> { + NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); + if (nsPolicyImpl.isPrimaryBroker(broker) || nsPolicyImpl.isSecondaryBroker(broker)) { + namespaceRegexes.addAll(policyData.getNamespaces()); + if (nsPolicyImpl.isPrimaryBroker(broker)) { + brokerIsolationData.primary(true); } - }); - brokerIsolationData.namespaceRegex(namespaceRegexes); - } + } + }); + brokerIsolationData.namespaceRegex(namespaceRegexes); return brokerIsolationData.build(); } + private CompletableFuture> internalGetNamespaceIsolationPolicies( + String cluster) { + return namespaceIsolationPolicies().getIsolationDataPoliciesAsync(cluster) + .thenApply(namespaceIsolationPolicies -> { + if (!namespaceIsolationPolicies.isPresent()) { + throw new RestException(Status.NOT_FOUND, + "NamespaceIsolationPolicies for cluster " + cluster + " does not exist"); + } + return namespaceIsolationPolicies.get().getPolicies(); + }); + } @POST @Path("/{cluster}/namespaceIsolationPolicies/{policyName}") @ApiOperation( @@ -603,150 +612,105 @@ public BrokerNamespaceIsolationData getBrokerWithNamespaceIsolationPolicy( }) public void setNamespaceIsolationPolicy( @Suspended final AsyncResponse asyncResponse, - @ApiParam( - value = "The cluster name", - required = true - ) + @ApiParam(value = "The cluster name", required = true) @PathParam("cluster") String cluster, - @ApiParam( - value = "The namespace isolation policy name", - required = true - ) + @ApiParam(value = "The namespace isolation policy name", required = true) @PathParam("policyName") String policyName, - @ApiParam( - value = "The namespace isolation policy data", - required = true - ) - NamespaceIsolationDataImpl policyData + @ApiParam(value = "The namespace isolation policy data", required = true) + NamespaceIsolationDataImpl policyData ) { - validateSuperUserAccess(); - validateClusterExists(cluster); - validatePoliciesReadOnlyAccess(); - - String jsonInput = null; - try { - // validate the policy data before creating the node - policyData.validate(); - jsonInput = ObjectMapperFactory.create().writeValueAsString(policyData); - - NamespaceIsolationPolicies nsIsolationPolicies = namespaceIsolationPolicies() - .getIsolationDataPolicies(cluster).orElseGet(() -> { + validateSuperUserAccessAsync() + .thenCompose(__ -> validatePoliciesReadOnlyAccessAsync()) + .thenCompose(__ -> validateClusterExistAsync(cluster, PRECONDITION_FAILED)) + .thenCompose(__ -> { + // validate the policy data before creating the node + policyData.validate(); + return namespaceIsolationPolicies().getIsolationDataPoliciesAsync(cluster); + }).thenCompose(nsIsolationPoliciesOpt -> + nsIsolationPoliciesOpt.map(CompletableFuture::completedFuture) + .orElseGet(() -> namespaceIsolationPolicies() + .setIsolationDataWithCreateAsync(cluster, (p) -> Collections.emptyMap()) + .thenApply(__ -> new NamespaceIsolationPolicies())) + ).thenCompose(nsIsolationPolicies -> namespaceIsolationPolicies() + .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()) + ).thenAccept(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) + .thenAccept(__ -> { + log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", + clientAppId(), cluster, policyName); + asyncResponse.resume(Response.noContent().build()); + }).exceptionally(ex -> { + Throwable realCause = FutureUtil.unwrapCompletionException(ex); + if (realCause instanceof IllegalArgumentException) { + String jsonData; try { - namespaceIsolationPolicies().setIsolationDataWithCreate(cluster, - (p) -> Collections.emptyMap()); - return new NamespaceIsolationPolicies(); - } catch (Exception e) { - throw new RestException(e); + jsonData = JsonUtil.toJson(policyData); + } catch (JsonUtil.ParseJsonException e) { + jsonData = "[Failed to serialize]"; } - }); - - nsIsolationPolicies.setPolicy(policyName, policyData); - namespaceIsolationPolicies().setIsolationData(cluster, old -> nsIsolationPolicies.getPolicies()); - - // whether or not make the isolation update on time. - if (pulsar().getConfiguration().isEnableNamespaceIsolationUpdateOnTime()) { - filterAndUnloadMatchedNameSpaces(asyncResponse, policyData); - } else { - asyncResponse.resume(Response.noContent().build()); - return; - } - } catch (IllegalArgumentException iae) { - log.info("[{}] Failed to update clusters/{}/namespaceIsolationPolicies/{}. Input data is invalid", - clientAppId(), cluster, policyName, iae); - asyncResponse.resume(new RestException(Status.BAD_REQUEST, - "Invalid format of input policy data. policy: " + policyName + "; data: " + jsonInput)); - } catch (NotFoundException nne) { - log.warn("[{}] Failed to update clusters/{}/namespaceIsolationPolicies: Does not exist", clientAppId(), - cluster); - asyncResponse.resume(new RestException(Status.NOT_FOUND, - "NamespaceIsolationPolicies for cluster " + cluster + " does not exist")); - } catch (Exception e) { - log.error("[{}] Failed to update clusters/{}/namespaceIsolationPolicies/{}", clientAppId(), cluster, - policyName, e); - asyncResponse.resume(new RestException(e)); - } - } - - // get matched namespaces; call unload for each namespaces; - private void filterAndUnloadMatchedNameSpaces(AsyncResponse asyncResponse, - NamespaceIsolationDataImpl policyData) throws Exception { - Namespaces namespaces = pulsar().getAdminClient().namespaces(); - - List nssToUnload = Lists.newArrayList(); - - pulsar().getAdminClient().tenants().getTenantsAsync() - .whenComplete((tenants, ex) -> { - if (ex != null) { - log.error("[{}] Failed to get tenants when setNamespaceIsolationPolicy.", clientAppId(), ex); - return; - } - AtomicInteger tenantsNumber = new AtomicInteger(tenants.size()); - // get all tenants now, for each tenants, get its namespaces - tenants.forEach(tenant -> namespaces.getNamespacesAsync(tenant) - .whenComplete((nss, e) -> { - int leftTenantsToHandle = tenantsNumber.decrementAndGet(); - if (e != null) { - log.error("[{}] Failed to get namespaces for tenant {} when setNamespaceIsolationPolicy.", - clientAppId(), tenant, e); - - if (leftTenantsToHandle == 0) { - unloadMatchedNamespacesList(asyncResponse, nssToUnload, namespaces); - } - - return; - } - - AtomicInteger nssNumber = new AtomicInteger(nss.size()); - - // get all namespaces for this tenant now. - nss.forEach(namespaceName -> { - int leftNssToHandle = nssNumber.decrementAndGet(); - - // if namespace match any policy regex, add it to ns list to be unload. - if (policyData.getNamespaces().stream() - .anyMatch(nsnameRegex -> namespaceName.matches(nsnameRegex))) { - nssToUnload.add(namespaceName); - } - - // all the tenants & namespaces get filtered. - if (leftNssToHandle == 0 && leftTenantsToHandle == 0) { - unloadMatchedNamespacesList(asyncResponse, nssToUnload, namespaces); - } - }); - })); - }); + asyncResponse.resume(new RestException(Status.BAD_REQUEST, + "Invalid format of input policy data. policy: " + policyName + "; data: " + jsonData)); + return null; + } else if (realCause instanceof NotFoundException) { + log.warn("[{}] Failed to update clusters/{}/namespaceIsolationPolicies: Does not exist", + clientAppId(), cluster); + asyncResponse.resume(new RestException(Status.NOT_FOUND, + "NamespaceIsolationPolicies for cluster " + cluster + " does not exist")); + return null; + } + log.info("[{}] Failed to update clusters/{}/namespaceIsolationPolicies/{}. Input data is invalid", + clientAppId(), cluster, policyName, realCause); + resumeAsyncResponseExceptionally(asyncResponse, ex); + return null; + }); } - private void unloadMatchedNamespacesList(AsyncResponse asyncResponse, - List nssToUnload, - Namespaces namespaces) { - if (nssToUnload.size() == 0) { - asyncResponse.resume(Response.noContent().build()); - return; + /** + * Get matched namespaces; call unload for each namespaces; + */ + private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIsolationDataImpl policyData) { + PulsarAdmin adminClient; + try { + adminClient = pulsar().getAdminClient(); + } catch (PulsarServerException e) { + return FutureUtil.failedFuture(e); } - - List> futures = nssToUnload.stream() - .map(namespaceName -> namespaces.unloadAsync(namespaceName)) - .collect(Collectors.toList()); - - FutureUtil.waitForAll(futures).whenComplete((result, exception) -> { - if (exception != null) { - log.error("[{}] Failed to unload namespace while setNamespaceIsolationPolicy.", - clientAppId(), exception); - asyncResponse.resume(new RestException(exception)); - return; - } - - 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); - } - - asyncResponse.resume(Response.noContent().build()); - return; - }); + return adminClient.tenants().getTenantsAsync() + .thenCompose(tenants -> { + List>> namespaceFutures = tenants.stream() + .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant)) + .collect(Collectors.toList()); + return FutureUtil.waitForAll(namespaceFutures) + .thenApply(__ -> { + List namespaces = namespaceFutures.stream() + .map(CompletableFuture::join) + .flatMap(List::stream) + .collect(Collectors.toList()); + // if namespace match any policy regex, add it to ns list to be unload. + return namespaces.stream() + .filter(namespaceName -> + policyData.getNamespaces().stream().anyMatch(namespaceName::matches)) + .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) + // Because triggering a write load report is a synchronous method, + // We should use the common pool to handle it to avoid impacting the metadata thread until + // load management is refactored into an asynchronous method. + .thenAcceptAsync(__ -> { + 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 @@ -1013,4 +977,18 @@ private void validateBrokerExistsInOtherDomain(final String cluster, final Strin private static final Logger log = LoggerFactory.getLogger(ClustersBase.class); + /** + * Verify that the cluster exists. + * For compatibility to avoid breaking changes, we can specify a REST status code when it doesn't exist. + * @param cluster Cluster name + * @param notExistStatus REST status code + */ + private CompletableFuture validateClusterExistAsync(String cluster, Status notExistStatus) { + return clusterResources().clusterExistsAsync(cluster) + .thenAccept(clusterExist -> { + if (!clusterExist) { + throw new RestException(notExistStatus, "Cluster " + cluster + " does not exist."); + } + }); + } } From 9bed6ed3ef91796eac2f8105f31e0624cd095fdf Mon Sep 17 00:00:00 2001 From: mattison chao Date: Tue, 10 May 2022 17:29:32 +0800 Subject: [PATCH 02/12] Fix test --- .../org/apache/pulsar/broker/admin/impl/ClustersBase.java | 2 +- .../test/java/org/apache/pulsar/broker/admin/AdminTest.java | 4 ++-- 2 files changed, 3 insertions(+), 3 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 92d85f6dab5bb..f98785a309282 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 @@ -633,7 +633,7 @@ public void setNamespaceIsolationPolicy( .thenApply(__ -> new NamespaceIsolationPolicies())) ).thenCompose(nsIsolationPolicies -> namespaceIsolationPolicies() .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()) - ).thenAccept(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) + ).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); 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 583a87275bda4..e3840a7ca338d 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 @@ -253,8 +253,8 @@ public void clusters() throws Exception { .parameters(parameters1) .build()) .build(); - AsyncResponse response = mock(AsyncResponse.class); - clusters.setNamespaceIsolationPolicy(response,"use", "policy1", policyData); + asyncRequests(ctx -> clusters.setNamespaceIsolationPolicy(ctx, + "use", "policy1", policyData)); asyncRequests(ctx -> clusters.getNamespaceIsolationPolicies(ctx, "use")); try { From 39c7c8dcaafe5eaac71644a61d63853d739a7117 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Tue, 10 May 2022 17:33:58 +0800 Subject: [PATCH 03/12] Fix test --- .../org/apache/pulsar/broker/admin/impl/ClustersBase.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 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 f98785a309282..3e2049204dd0a 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 @@ -631,9 +631,11 @@ public void setNamespaceIsolationPolicy( .orElseGet(() -> namespaceIsolationPolicies() .setIsolationDataWithCreateAsync(cluster, (p) -> Collections.emptyMap()) .thenApply(__ -> new NamespaceIsolationPolicies())) - ).thenCompose(nsIsolationPolicies -> namespaceIsolationPolicies() - .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()) - ).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) + ).thenCompose(nsIsolationPolicies -> { + nsIsolationPolicies.setPolicy(policyName, policyData); + return namespaceIsolationPolicies() + .setIsolationDataAsync(cluster, old -> nsIsolationPolicies.getPolicies()); + }).thenCompose(__ -> filterAndUnloadMatchedNamespaceAsync(policyData)) .thenAccept(__ -> { log.info("[{}] Successful to update clusters/{}/namespaceIsolationPolicies/{}.", clientAppId(), cluster, policyName); From 4bfae06e4b8fda94d51c24a3350f355be9d3dcc1 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 07:07:58 +0800 Subject: [PATCH 04/12] Rebase to master --- .../broker/admin/impl/ClustersBase.java | 47 ------------------- 1 file changed, 47 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 3e2049204dd0a..d923738d3bae7 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 @@ -511,27 +511,6 @@ public void getBrokersWithNamespaceIsolationPolicy( }); } - private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( - String broker, - Map policies) { - BrokerNamespaceIsolationData.Builder brokerIsolationData = - BrokerNamespaceIsolationData.builder().brokerName(broker); - if (policies == null) { - return brokerIsolationData.build(); - } - List namespaceRegexes = new ArrayList<>(); - policies.forEach((name, policyData) -> { - NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); - if (nsPolicyImpl.isPrimaryBroker(broker) || nsPolicyImpl.isSecondaryBroker(broker)) { - namespaceRegexes.addAll(policyData.getNamespaces()); - brokerIsolationData.primary(nsPolicyImpl.isPrimaryBroker(broker)); - brokerIsolationData.policyName(name); - } - }); - brokerIsolationData.namespaceRegex(namespaceRegexes); - return brokerIsolationData.build(); - } - @GET @Path("/{cluster}/namespaceIsolationPolicies/brokers/{broker}") @ApiOperation( @@ -586,17 +565,6 @@ private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( return brokerIsolationData.build(); } - private CompletableFuture> internalGetNamespaceIsolationPolicies( - String cluster) { - return namespaceIsolationPolicies().getIsolationDataPoliciesAsync(cluster) - .thenApply(namespaceIsolationPolicies -> { - if (!namespaceIsolationPolicies.isPresent()) { - throw new RestException(Status.NOT_FOUND, - "NamespaceIsolationPolicies for cluster " + cluster + " does not exist"); - } - return namespaceIsolationPolicies.get().getPolicies(); - }); - } @POST @Path("/{cluster}/namespaceIsolationPolicies/{policyName}") @ApiOperation( @@ -978,19 +946,4 @@ private void validateBrokerExistsInOtherDomain(final String cluster, final Strin } private static final Logger log = LoggerFactory.getLogger(ClustersBase.class); - - /** - * Verify that the cluster exists. - * For compatibility to avoid breaking changes, we can specify a REST status code when it doesn't exist. - * @param cluster Cluster name - * @param notExistStatus REST status code - */ - private CompletableFuture validateClusterExistAsync(String cluster, Status notExistStatus) { - return clusterResources().clusterExistsAsync(cluster) - .thenAccept(clusterExist -> { - if (!clusterExist) { - throw new RestException(notExistStatus, "Cluster " + cluster + " does not exist."); - } - }); - } } From f1ab5afa4693e640098c6d3a6846edbb30115909 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 07:10:56 +0800 Subject: [PATCH 05/12] Fix checkstyle --- .../org/apache/pulsar/broker/admin/impl/ClustersBase.java | 5 +---- 1 file changed, 1 insertion(+), 4 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 d923738d3bae7..0ce47431ff1a0 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 @@ -18,6 +18,7 @@ */ package org.apache.pulsar.broker.admin.impl; +import static javax.ws.rs.core.Response.Status.PRECONDITION_FAILED; import com.google.common.collect.Maps; import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiParam; @@ -44,7 +45,6 @@ import javax.ws.rs.core.MediaType; import javax.ws.rs.core.Response; import javax.ws.rs.core.Response.Status; - import org.apache.bookkeeper.common.util.JsonUtil; import org.apache.commons.collections4.CollectionUtils; import org.apache.pulsar.broker.PulsarServerException; @@ -59,7 +59,6 @@ import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.ClusterDataImpl; import org.apache.pulsar.common.policies.data.FailureDomainImpl; -import org.apache.pulsar.common.policies.data.NamespaceIsolationData; import org.apache.pulsar.common.policies.data.NamespaceIsolationDataImpl; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicies; import org.apache.pulsar.common.policies.impl.NamespaceIsolationPolicyImpl; @@ -69,8 +68,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import static javax.ws.rs.core.Response.Status.PRECONDITION_FAILED; - public class ClustersBase extends AdminResource { @GET From 692f72fa43a79d77a4f6b4a2019aa33001bb30cf Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 07:14:24 +0800 Subject: [PATCH 06/12] Rollback and fix tests --- .../broker/admin/impl/ClustersBase.java | 45 ++++++++++--------- .../apache/pulsar/broker/admin/AdminTest.java | 4 +- 2 files changed, 25 insertions(+), 24 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 0ce47431ff1a0..1be73a26acfd5 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 @@ -508,6 +508,29 @@ public void getBrokersWithNamespaceIsolationPolicy( }); } + + private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( + String broker, + Map policies) { + BrokerNamespaceIsolationData.Builder brokerIsolationData = + BrokerNamespaceIsolationData.builder().brokerName(broker); + if (policies == null) { + return brokerIsolationData.build(); + } + List namespaceRegexes = new ArrayList<>(); + policies.forEach((name, policyData) -> { + NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); + if (nsPolicyImpl.isPrimaryBroker(broker) || nsPolicyImpl.isSecondaryBroker(broker)) { + namespaceRegexes.addAll(policyData.getNamespaces()); + if (nsPolicyImpl.isPrimaryBroker(broker)) { + brokerIsolationData.primary(true); + } + } + }); + brokerIsolationData.namespaceRegex(namespaceRegexes); + return brokerIsolationData.build(); + } + @GET @Path("/{cluster}/namespaceIsolationPolicies/brokers/{broker}") @ApiOperation( @@ -540,28 +563,6 @@ public void getBrokerWithNamespaceIsolationPolicy( }); } - private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( - String broker, - Map policies) { - BrokerNamespaceIsolationData.Builder brokerIsolationData = - BrokerNamespaceIsolationData.builder().brokerName(broker); - if (policies == null) { - return brokerIsolationData.build(); - } - List namespaceRegexes = new ArrayList<>(); - policies.forEach((name, policyData) -> { - NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); - if (nsPolicyImpl.isPrimaryBroker(broker) || nsPolicyImpl.isSecondaryBroker(broker)) { - namespaceRegexes.addAll(policyData.getNamespaces()); - if (nsPolicyImpl.isPrimaryBroker(broker)) { - brokerIsolationData.primary(true); - } - } - }); - brokerIsolationData.namespaceRegex(namespaceRegexes); - return brokerIsolationData.build(); - } - @POST @Path("/{cluster}/namespaceIsolationPolicies/{policyName}") @ApiOperation( 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 e3840a7ca338d..ec026af80459a 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 @@ -403,8 +403,8 @@ public void clusters() throws Exception { } catch (RestException e) { assertEquals(e.getResponse().getStatus(), Status.PRECONDITION_FAILED.getStatusCode()); } - verify(clusters, times(22)).validateSuperUserAccessAsync(); - verify(clusters, times(2)).validateSuperUserAccess(); + verify(clusters, times(23)).validateSuperUserAccessAsync(); + verify(clusters, times(1)).validateSuperUserAccess(); } @Test From c56b6728a07032f4d324276db219a4f6f05dbff1 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 07:41:05 +0800 Subject: [PATCH 07/12] Refactor some methods --- .../pulsar/broker/admin/impl/ClustersBase.java | 14 +++++--------- .../org/apache/pulsar/common/util/FutureUtil.java | 15 +++++++++++++++ 2 files changed, 20 insertions(+), 9 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 1be73a26acfd5..242de5d25588d 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 @@ -34,6 +34,7 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; 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; @@ -644,15 +645,10 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIs } return adminClient.tenants().getTenantsAsync() .thenCompose(tenants -> { - List>> namespaceFutures = tenants.stream() - .map(tenant -> adminClient.namespaces().getNamespacesAsync(tenant)) - .collect(Collectors.toList()); - return FutureUtil.waitForAll(namespaceFutures) - .thenApply(__ -> { - List namespaces = namespaceFutures.stream() - .map(CompletableFuture::join) - .flatMap(List::stream) - .collect(Collectors.toList()); + Stream>> completableFutureStream = tenants.stream() + .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. return namespaces.stream() .filter(namespaceName -> diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java index e5c2caeb7d04b..e1176cc8b9648 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java @@ -19,7 +19,9 @@ package org.apache.pulsar.common.util; import java.time.Duration; +import java.util.ArrayList; import java.util.Collection; +import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; @@ -31,6 +33,7 @@ import java.util.function.Predicate; import java.util.function.Supplier; import java.util.stream.Collectors; +import java.util.stream.Stream; /** * This class is aimed at simplifying work with {@code CompletableFuture}. @@ -47,6 +50,18 @@ public static CompletableFuture waitForAll(Collection CompletableFuture> waitForAll(Stream>> futures) { + return futures.reduce(CompletableFuture.completedFuture(new ArrayList<>()), + (pre, curr) -> pre.thenCompose(preV -> curr.thenApply(currV -> { + preV.addAll(currV); + return preV; + }))); + } + + public static CompletableFuture> waitForAll(List>> futures) { + return waitForAll(futures.stream()); + } + /** * Return a future that represents the completion of any future in the provided Collection. * From f59e97c219d5c240e52e169d8cbbbc77a8871c67 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 07:53:01 +0800 Subject: [PATCH 08/12] Fix checkstyle --- .../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 242de5d25588d..6fa4f5733f92e 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 @@ -634,7 +634,7 @@ public void setNamespaceIsolationPolicy( } /** - * Get matched namespaces; call unload for each namespaces; + * Get matched namespaces; call unload for each namespaces. */ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIsolationDataImpl policyData) { PulsarAdmin adminClient; From 9992d7a0e4f33cade5800caeca6a63057f88f8e7 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 13:19:02 +0800 Subject: [PATCH 09/12] Remove useless method. --- .../main/java/org/apache/pulsar/common/util/FutureUtil.java | 4 ---- 1 file changed, 4 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java index e1176cc8b9648..afad61eb66918 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/FutureUtil.java @@ -58,10 +58,6 @@ public static CompletableFuture> waitForAll(Stream CompletableFuture> waitForAll(List>> futures) { - return waitForAll(futures.stream()); - } - /** * Return a future that represents the completion of any future in the provided Collection. * From b5f2fc0430f09ee318ddb95f4c41885dd0c53d72 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 14:24:49 +0800 Subject: [PATCH 10/12] Fix the issue caused by rebase operation. --- .../org/apache/pulsar/broker/admin/impl/ClustersBase.java | 6 ++---- 1 file changed, 2 insertions(+), 4 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 6fa4f5733f92e..3f5816053c8b8 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 @@ -522,10 +522,8 @@ private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( policies.forEach((name, policyData) -> { NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); if (nsPolicyImpl.isPrimaryBroker(broker) || nsPolicyImpl.isSecondaryBroker(broker)) { - namespaceRegexes.addAll(policyData.getNamespaces()); - if (nsPolicyImpl.isPrimaryBroker(broker)) { - brokerIsolationData.primary(true); - } + brokerIsolationData.primary(nsPolicyImpl.isPrimaryBroker(broker)); + brokerIsolationData.policyName(name); } }); brokerIsolationData.namespaceRegex(namespaceRegexes); From d7d06ecf19dca73976c32eae445a7c06769ac671 Mon Sep 17 00:00:00 2001 From: mattison chao Date: Mon, 16 May 2022 15:00:42 +0800 Subject: [PATCH 11/12] Remove use common-pool run async method. --- .../org/apache/pulsar/broker/admin/impl/ClustersBase.java | 5 +---- 1 file changed, 1 insertion(+), 4 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 3f5816053c8b8..7e203b61e890a 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 @@ -661,10 +661,7 @@ private CompletableFuture filterAndUnloadMatchedNamespaceAsync(NamespaceIs .map(namespaceName -> adminClient.namespaces().unloadAsync(namespaceName)) .collect(Collectors.toList()); return FutureUtil.waitForAll(futures) - // Because triggering a write load report is a synchronous method, - // We should use the common pool to handle it to avoid impacting the metadata thread until - // load management is refactored into an asynchronous method. - .thenAcceptAsync(__ -> { + .thenAccept(__ -> { try { // write load info to load manager to make the load happens fast pulsar().getLoadManager().get().writeLoadReportOnZookeeper(true); From 1c15efb6e1653968de1e0f9c9b615df0d4afa4de Mon Sep 17 00:00:00 2001 From: mattison chao Date: Wed, 18 May 2022 09:13:41 +0800 Subject: [PATCH 12/12] Fix test --- .../java/org/apache/pulsar/broker/admin/impl/ClustersBase.java | 1 + 1 file changed, 1 insertion(+) 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 7e203b61e890a..a1fd8a247334a 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 @@ -522,6 +522,7 @@ private BrokerNamespaceIsolationData internalGetBrokerNsIsolationData( policies.forEach((name, policyData) -> { NamespaceIsolationPolicyImpl nsPolicyImpl = new NamespaceIsolationPolicyImpl(policyData); if (nsPolicyImpl.isPrimaryBroker(broker) || nsPolicyImpl.isSecondaryBroker(broker)) { + namespaceRegexes.addAll(policyData.getNamespaces()); brokerIsolationData.primary(nsPolicyImpl.isPrimaryBroker(broker)); brokerIsolationData.policyName(name); }