From 27660e1246b0048dbe2c9615646db11e74cad0c6 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 12 Jan 2021 19:00:54 +0800 Subject: [PATCH 1/4] fix admin-api-brokers list failed --- .../pulsar/broker/web/PulsarWebResource.java | 22 +++++++++++++------ .../pulsar/broker/service/ReplicatorTest.java | 19 ++++++++++++++++ 2 files changed, 34 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java index 7a06c120331cc..c8242e99ea212 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java @@ -49,6 +49,8 @@ import org.apache.pulsar.broker.authorization.AuthorizationService; import org.apache.pulsar.broker.namespace.LookupOptions; import org.apache.pulsar.broker.namespace.NamespaceService; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.impl.PulsarServiceNameResolver; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceBundles; @@ -87,6 +89,8 @@ public abstract class PulsarWebResource { private PulsarService pulsar; + private final PulsarServiceNameResolver serviceNameResolver = new PulsarServiceNameResolver(); + protected PulsarService pulsar() { if (pulsar == null) { pulsar = (PulsarService) servletContext.getAttribute(WebService.ATTRIBUTE_PULSAR_NAME); @@ -356,14 +360,18 @@ protected void validateClusterOwnership(String cluster) throws WebApplicationExc } private URI getRedirectionUrl(ClusterData differentClusterData) throws MalformedURLException { - URL webUrl = null; - if (isRequestHttps() && pulsar.getConfiguration().getWebServicePortTls().isPresent() - && StringUtils.isNotBlank(differentClusterData.getServiceUrlTls())) { - webUrl = new URL(differentClusterData.getServiceUrlTls()); - } else { - webUrl = new URL(differentClusterData.getServiceUrl()); + try { + if (isRequestHttps() && pulsar.getConfiguration().getWebServicePortTls().isPresent() + && StringUtils.isNotBlank(differentClusterData.getServiceUrlTls())) { + serviceNameResolver.updateServiceUrl(differentClusterData.getServiceUrlTls()); + } else { + serviceNameResolver.updateServiceUrl(differentClusterData.getServiceUrl()); + } + URL webUrl = new URL(serviceNameResolver.resolveHostUri().toString()); + return UriBuilder.fromUri(uri.getRequestUri()).host(webUrl.getHost()).port(webUrl.getPort()).build(); + } catch (PulsarClientException.InvalidServiceURL exception) { + throw new MalformedURLException(exception.getMessage()); } - return UriBuilder.fromUri(uri.getRequestUri()).host(webUrl.getHost()).port(webUrl.getPort()).build(); } protected static CompletableFuture getClusterDataIfDifferentCluster(PulsarService pulsar, diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index b1580d1bda4dd..32061bf4dffad 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -69,9 +69,11 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.BacklogQuota.RetentionPolicy; +import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.ReplicatorStats; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; +import org.awaitility.Awaitility; import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -198,6 +200,23 @@ public Void call() throws Exception { // Case 3: TODO: Once automatic cleanup is implemented, add tests case to verify auto removal of clusters } + @Test + public void activeBrokerParse() throws Exception { + pulsar1.getConfiguration().setAuthorizationEnabled(true); + //init clusterData + ClusterData cluster2Data = new ClusterData(); + String cluster2ServiceUrls = String.format("%s,localhost:1234,localhost:5678", pulsar2.getWebServiceAddress()); + cluster2ServiceUrls = cluster2ServiceUrls.replace("pulsar", "http"); + cluster2Data.setServiceUrl(cluster2ServiceUrls); + String cluster2 = "activeCLuster2"; + admin2.clusters().createCluster(cluster2, cluster2Data); + Awaitility.await().atMost(3, TimeUnit.SECONDS).until(() + -> admin2.clusters().getCluster(cluster2) != null); + + List list = admin1.brokers().getActiveBrokers(cluster2); + assertEquals(list.get(0), url2.toString().replace("http://", "")); + } + @SuppressWarnings("unchecked") @Test(timeOut = 30000) public void testConcurrentReplicator() throws Exception { From cbbd8c03aec5bf3c418d58dcddb1859cdab43d39 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 12 Jan 2021 19:07:02 +0800 Subject: [PATCH 2/4] move code --- .../java/org/apache/pulsar/broker/web/PulsarWebResource.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java index c8242e99ea212..5b5fb5009002e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/PulsarWebResource.java @@ -89,8 +89,6 @@ public abstract class PulsarWebResource { private PulsarService pulsar; - private final PulsarServiceNameResolver serviceNameResolver = new PulsarServiceNameResolver(); - protected PulsarService pulsar() { if (pulsar == null) { pulsar = (PulsarService) servletContext.getAttribute(WebService.ATTRIBUTE_PULSAR_NAME); @@ -361,6 +359,7 @@ protected void validateClusterOwnership(String cluster) throws WebApplicationExc private URI getRedirectionUrl(ClusterData differentClusterData) throws MalformedURLException { try { + PulsarServiceNameResolver serviceNameResolver = new PulsarServiceNameResolver(); if (isRequestHttps() && pulsar.getConfiguration().getWebServicePortTls().isPresent() && StringUtils.isNotBlank(differentClusterData.getServiceUrlTls())) { serviceNameResolver.updateServiceUrl(differentClusterData.getServiceUrlTls()); From 8dd376bc35c91453a587aeee076be4963da408af Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Tue, 12 Jan 2021 19:09:31 +0800 Subject: [PATCH 3/4] change unit test --- .../java/org/apache/pulsar/broker/service/ReplicatorTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 32061bf4dffad..351bdcced2967 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -206,7 +206,6 @@ public void activeBrokerParse() throws Exception { //init clusterData ClusterData cluster2Data = new ClusterData(); String cluster2ServiceUrls = String.format("%s,localhost:1234,localhost:5678", pulsar2.getWebServiceAddress()); - cluster2ServiceUrls = cluster2ServiceUrls.replace("pulsar", "http"); cluster2Data.setServiceUrl(cluster2ServiceUrls); String cluster2 = "activeCLuster2"; admin2.clusters().createCluster(cluster2, cluster2Data); From 9eaa65b82454dcc656ccca9386e8782db5ff6222 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Wed, 13 Jan 2021 11:02:05 +0800 Subject: [PATCH 4/4] fix unit test --- .../java/org/apache/pulsar/broker/service/ReplicatorTest.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java index 351bdcced2967..41d085bb1e121 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTest.java @@ -200,7 +200,7 @@ public Void call() throws Exception { // Case 3: TODO: Once automatic cleanup is implemented, add tests case to verify auto removal of clusters } - @Test + @Test(timeOut = 10000) public void activeBrokerParse() throws Exception { pulsar1.getConfiguration().setAuthorizationEnabled(true); //init clusterData @@ -214,6 +214,8 @@ public void activeBrokerParse() throws Exception { List list = admin1.brokers().getActiveBrokers(cluster2); assertEquals(list.get(0), url2.toString().replace("http://", "")); + //restore configuration + pulsar1.getConfiguration().setAuthorizationEnabled(false); } @SuppressWarnings("unchecked")