From aecdb8e556b9a9f5f18d7c23aebd1143fd27abbf Mon Sep 17 00:00:00 2001 From: technoboy Date: Thu, 10 Mar 2022 12:41:52 +0800 Subject: [PATCH 1/3] Fix get wrong prompt exception when get non-persistent topic list with un-authorized permission. --- .../broker/admin/v2/NonPersistentTopics.java | 28 +++++++---------- .../pulsar/broker/auth/AuthorizationTest.java | 30 +++++++++++++++++-- 2 files changed, 38 insertions(+), 20 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index cbb06b657267f..0bbcfa1dea981 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java @@ -415,25 +415,19 @@ public void getList( } final List topics = Lists.newArrayList(); - FutureUtil.waitForAll(futures).handle((result, exception) -> { - for (int i = 0; i < futures.size(); i++) { - try { - if (futures.get(i).isDone() && futures.get(i).get() != null) { - topics.addAll(futures.get(i).get()); - } - } catch (InterruptedException | ExecutionException e) { - log.error("[{}] Failed to get list of topics under namespace {}", clientAppId(), namespaceName, e); - asyncResponse.resume(new RestException(e instanceof ExecutionException ? e.getCause() : e)); - return null; + FutureUtil.waitForAll(futures).whenComplete((result, ex) -> { + if (ex != null) { + resumeAsyncResponseExceptionally(asyncResponse, ex); + } else { + for (int i = 0; i < futures.size(); i++) { + topics.addAll(futures.get(i).join()); } + final List nonPersistentTopics = + topics.stream() + .filter(name -> !TopicName.get(name).isPersistent()) + .collect(Collectors.toList()); + asyncResponse.resume(nonPersistentTopics); } - - final List nonPersistentTopics = - topics.stream() - .filter(name -> !TopicName.get(name).isPersistent()) - .collect(Collectors.toList()); - asyncResponse.resume(nonPersistentTopics); - return null; }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java index 49d46cd9b8f98..2ebdb9e7457b4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java @@ -18,14 +18,17 @@ */ package org.apache.pulsar.broker.auth; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; - import java.util.EnumSet; - import org.apache.pulsar.broker.authorization.AuthorizationService; +import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminBuilder; +import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.common.naming.TopicDomain; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.ClusterData; @@ -35,7 +38,6 @@ import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; - import com.google.common.collect.Sets; @Test(groups = "flaky") @@ -229,6 +231,28 @@ public void simple() throws Exception { admin.clusters().deleteCluster("c1"); } + @Test + public void testGetListWithoutGetBundleOp() throws Exception { + String tenant = "p1"; + String namespace = "p1/ns1"; + admin.clusters().createCluster("c1", ClusterData.builder().build()); + admin.tenants().createTenant(tenant, new TenantInfoImpl(Sets.newHashSet("role1"), Sets.newHashSet("c1"))); + admin.namespaces().createNamespace(namespace); + admin.namespaces().grantPermissionOnNamespace(namespace, "pass.pass2", EnumSet.of(AuthAction.produce)); + PulsarAdmin admin2 = PulsarAdmin.builder().serviceHttpUrl(brokerUrl != null + ? brokerUrl.toString() + : brokerUrlTls.toString()) + .authentication(new MockAuthentication("pass.pass2")) + .build(); + when(pulsar.getAdminClient()).thenReturn(admin2); + try { + admin2.topics().getList(namespace, TopicDomain.non_persistent); + } catch (Exception ex) { + assertTrue(ex instanceof PulsarAdminException.NotAuthorizedException); + assertEquals(ex.getMessage(), "Unauthorized to validateNamespaceOperation for operation [GET_BUNDLE] on namespace [p1/ns1]"); + } + } + private static void waitForChange() { try { Thread.sleep(100); From c9885360ba592e94fb9b9b8928b854f95b1fd70d Mon Sep 17 00:00:00 2001 From: technoboy Date: Thu, 10 Mar 2022 14:26:01 +0800 Subject: [PATCH 2/3] update. --- .../apache/pulsar/broker/admin/v2/NonPersistentTopics.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index 0bbcfa1dea981..66867b68426a0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java @@ -414,13 +414,16 @@ public void getList( } } - final List topics = Lists.newArrayList(); FutureUtil.waitForAll(futures).whenComplete((result, ex) -> { if (ex != null) { resumeAsyncResponseExceptionally(asyncResponse, ex); } else { + final List topics = Lists.newArrayList(); for (int i = 0; i < futures.size(); i++) { - topics.addAll(futures.get(i).join()); + List topicList = futures.get(i).join(); + if (topicList != null) { + topics.addAll(topicList); + } } final List nonPersistentTopics = topics.stream() From 1c20868baaba352b2c580131b43ad61c9cfd6428 Mon Sep 17 00:00:00 2001 From: technoboy Date: Fri, 11 Mar 2022 21:46:19 +0800 Subject: [PATCH 3/3] add v1 support. --- .../broker/admin/v1/NonPersistentTopics.java | 25 +++++++++---------- .../pulsar/broker/auth/AuthorizationTest.java | 19 ++++++++++---- 2 files changed, 26 insertions(+), 18 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java index bb28ee7934944..dc466be86d67b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java @@ -200,6 +200,7 @@ public void getList(@Suspended final AsyncResponse asyncResponse, @PathParam("pr Policies policies = null; NamespaceName nsName = null; try { + validateNamespaceName(property, cluster, namespace); validateNamespaceOperation(namespaceName, NamespaceOperation.GET_TOPICS); policies = getNamespacePolicies(property, cluster, namespace); nsName = NamespaceName.get(property, cluster, namespace); @@ -237,22 +238,19 @@ public void getList(@Suspended final AsyncResponse asyncResponse, @PathParam("pr } } - final List topics = Lists.newArrayList(); - FutureUtil.waitForAll(futures).handle((result, exception) -> { - for (int i = 0; i < futures.size(); i++) { - try { - if (futures.get(i).isDone() && futures.get(i).get() != null) { - topics.addAll(futures.get(i).get()); + FutureUtil.waitForAll(futures).whenComplete((result, ex) -> { + if (ex != null) { + resumeAsyncResponseExceptionally(asyncResponse, ex); + } else { + final List topics = Lists.newArrayList(); + for (int i = 0; i < futures.size(); i++) { + List topicList = futures.get(i).join(); + if (topicList != null) { + topics.addAll(topicList); } - } catch (InterruptedException | ExecutionException e) { - log.error("[{}] Failed to get list of topics under namespace {}/{}/{}", clientAppId(), property, - cluster, namespace, e); - asyncResponse.resume(new RestException(e instanceof ExecutionException ? e.getCause() : e)); - return null; } + asyncResponse.resume(topics); } - asyncResponse.resume(topics); - return null; }); } @@ -269,6 +267,7 @@ public List getListFromBundle(@PathParam("property") String property, @P @PathParam("bundle") String bundleRange) { log.info("[{}] list of topics on namespace bundle {}/{}/{}/{}", clientAppId(), property, cluster, namespace, bundleRange); + validateNamespaceName(property, cluster, namespace); validateNamespaceOperation(namespaceName, NamespaceOperation.GET_BUNDLE); Policies policies = getNamespacePolicies(property, cluster, namespace); if (!cluster.equals(Constants.GLOBAL_CLUSTER)) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java index 2ebdb9e7457b4..574e7a14c43f4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/AuthorizationTest.java @@ -234,11 +234,14 @@ public void simple() throws Exception { @Test public void testGetListWithoutGetBundleOp() throws Exception { String tenant = "p1"; - String namespace = "p1/ns1"; + String namespaceV1 = "p1/global/ns1"; + String namespaceV2 = "p1/ns2"; admin.clusters().createCluster("c1", ClusterData.builder().build()); admin.tenants().createTenant(tenant, new TenantInfoImpl(Sets.newHashSet("role1"), Sets.newHashSet("c1"))); - admin.namespaces().createNamespace(namespace); - admin.namespaces().grantPermissionOnNamespace(namespace, "pass.pass2", EnumSet.of(AuthAction.produce)); + admin.namespaces().createNamespace(namespaceV1, Sets.newHashSet("c1")); + admin.namespaces().grantPermissionOnNamespace(namespaceV1, "pass.pass2", EnumSet.of(AuthAction.produce)); + admin.namespaces().createNamespace(namespaceV2, Sets.newHashSet("c1")); + admin.namespaces().grantPermissionOnNamespace(namespaceV2, "pass.pass2", EnumSet.of(AuthAction.produce)); PulsarAdmin admin2 = PulsarAdmin.builder().serviceHttpUrl(brokerUrl != null ? brokerUrl.toString() : brokerUrlTls.toString()) @@ -246,10 +249,16 @@ public void testGetListWithoutGetBundleOp() throws Exception { .build(); when(pulsar.getAdminClient()).thenReturn(admin2); try { - admin2.topics().getList(namespace, TopicDomain.non_persistent); + admin2.topics().getList(namespaceV1, TopicDomain.non_persistent); } catch (Exception ex) { assertTrue(ex instanceof PulsarAdminException.NotAuthorizedException); - assertEquals(ex.getMessage(), "Unauthorized to validateNamespaceOperation for operation [GET_BUNDLE] on namespace [p1/ns1]"); + assertEquals(ex.getMessage(), "Unauthorized to validateNamespaceOperation for operation [GET_BUNDLE] on namespace [p1/global/ns1]"); + } + try { + admin2.topics().getList(namespaceV2, TopicDomain.non_persistent); + } catch (Exception ex) { + assertTrue(ex instanceof PulsarAdminException.NotAuthorizedException); + assertEquals(ex.getMessage(), "Unauthorized to validateNamespaceOperation for operation [GET_BUNDLE] on namespace [p1/ns2]"); } }