diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authorization/PulsarAuthorizationProvider.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authorization/PulsarAuthorizationProvider.java index de586345fd0da..0e60c8bacafca 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authorization/PulsarAuthorizationProvider.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/authorization/PulsarAuthorizationProvider.java @@ -247,7 +247,8 @@ public CompletableFuture grantPermissionAsync(TopicName topicName, Set { policies.auth_policies.getTopicAuthentication() diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 8a0afd93213ea..e7fc3348af5d4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -273,6 +273,24 @@ protected CompletableFuture>> internalGetPermissions } } } + + // If topic is partitioned, add based topic permission + if (topicName.isPartitioned() && auth.getTopicAuthentication().containsKey( + topicName.getPartitionedTopicName())) { + for (Map.Entry> entry : + auth.getTopicAuthentication().get(topicName.getPartitionedTopicName()).entrySet()) { + String role = entry.getKey(); + Set topicPermissions = entry.getValue(); + + if (!permissions.containsKey(role)) { + permissions.put(role, topicPermissions); + } else { + // Do the union between namespace and topic level + Set union = Sets.union(permissions.get(role), topicPermissions); + permissions.put(role, union); + } + } + } return permissions; })); } @@ -326,20 +344,9 @@ protected void internalGrantPermissionsOnTopic(final AsyncResponse asyncResponse // This operation should be reading from zookeeper and it should be allowed without having admin privileges validateAdminAccessForTenantAsync(namespaceName.getTenant()) .thenCompose(__ -> validatePoliciesReadOnlyAccessAsync().thenCompose(unused1 -> - getPartitionedTopicMetadataAsync(topicName, true, false) - .thenCompose(metadata -> { - int numPartitions = metadata.partitions; - CompletableFuture future = CompletableFuture.completedFuture(null); - if (numPartitions > 0) { - for (int i = 0; i < numPartitions; i++) { - TopicName topicNamePartition = topicName.getPartition(i); - future = future.thenCompose(unused -> grantPermissionsAsync(topicNamePartition, role, - actions)); - } - } - return future.thenCompose(unused -> grantPermissionsAsync(topicName, role, actions)) - .thenAccept(unused -> asyncResponse.resume(Response.noContent().build())); - }))).exceptionally(ex -> { + grantPermissionsAsync(topicName, role, actions) + .thenAccept(unused -> asyncResponse.resume(Response.noContent().build())))) + .exceptionally(ex -> { Throwable realCause = FutureUtil.unwrapCompletionException(ex); log.error("[{}] Failed to get permissions for topic {}", clientAppId(), topicName, realCause); resumeAsyncResponseExceptionally(asyncResponse, realCause); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthenticatedProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthenticatedProducerConsumerTest.java index 0279116821094..8751fe6a779a2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthenticatedProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthenticatedProducerConsumerTest.java @@ -401,11 +401,6 @@ public void testDeleteAuthenticationPoliciesOfTopic() throws Exception { Awaitility.await().untilAsserted(() -> { assertTrue(pulsar.getPulsarResources().getNamespaceResources().getPolicies(NamespaceName.get("p1/ns1")) .get().auth_policies.getTopicAuthentication().containsKey(partitionedTopic)); - for (int i = 0; i < numPartitions; i++) { - assertTrue(pulsar.getPulsarResources().getNamespaceResources().getPolicies(NamespaceName.get("p1/ns1")) - .get().auth_policies.getTopicAuthentication() - .containsKey(TopicName.get(partitionedTopic).getPartition(i).toString())); - } }); admin.topics().deletePartitionedTopic("persistent://p1/ns1/partitioned-topic");