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 4ff236af5ac27..d5952f8627e50 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 @@ -328,7 +328,13 @@ private void revokePermissions(String topicUri, String role) { try { // Write the new policies to metadata store namespaceResources().setPolicies(namespaceName, p -> { - p.auth_policies.getTopicAuthentication().get(topicUri).remove(role); + p.auth_policies.getTopicAuthentication().computeIfPresent(topicUri, (k, roles) -> { + roles.remove(role); + if (roles.isEmpty()) { + return null; + } + return roles; + }); return p; }); log.info("[{}] Successfully revoke access for role {} - topic {}", clientAppId(), role, topicUri); 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 046b26846e2d3..6d9135af1acd4 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 @@ -403,4 +403,45 @@ public void testDeleteAuthenticationPoliciesOfTopic() throws Exception { admin.tenants().deleteTenant("p1"); admin.clusters().deleteCluster("test"); } + + @Test + public void testCleanupEmptyTopicAuthenticationMap() throws Exception { + Map authParams = new HashMap<>(); + authParams.put("tlsCertFile", TLS_CLIENT_CERT_FILE_PATH); + authParams.put("tlsKeyFile", TLS_CLIENT_KEY_FILE_PATH); + Authentication authTls = new AuthenticationTls(); + authTls.configure(authParams); + internalSetup(authTls); + + admin.clusters().createCluster("test", ClusterData.builder().build()); + admin.tenants().createTenant("p1", + new TenantInfoImpl(Collections.emptySet(), new HashSet<>(admin.clusters().getClusters()))); + admin.namespaces().createNamespace("p1/ns1"); + + String topic = "persistent://p1/ns1/topic"; + admin.topics().createNonPartitionedTopic(topic); + assertFalse(pulsar.getPulsarResources().getNamespaceResources().getPolicies(NamespaceName.get("p1/ns1")) + .get().auth_policies.getTopicAuthentication().containsKey(topic)); + + // grant permission + admin.topics().grantPermission(topic, "test-user-1", EnumSet.of(AuthAction.consume)); + Awaitility.await().untilAsserted(() -> { + assertTrue(pulsar.getPulsarResources().getNamespaceResources().getPolicies(NamespaceName.get("p1/ns1")) + .get().auth_policies.getTopicAuthentication().containsKey(topic)); + }); + + // revoke permission + admin.topics().revokePermissions(topic, "test-user-1"); + Awaitility.await().untilAsserted(() -> { + assertFalse(pulsar.getPulsarResources().getNamespaceResources().getPolicies(NamespaceName.get("p1/ns1")) + .get().auth_policies.getTopicAuthentication().containsKey(topic)); + }); + + // grant permission again + admin.topics().grantPermission(topic, "test-user-1", EnumSet.of(AuthAction.consume)); + Awaitility.await().untilAsserted(() -> { + assertTrue(pulsar.getPulsarResources().getNamespaceResources().getPolicies(NamespaceName.get("p1/ns1")) + .get().auth_policies.getTopicAuthentication().containsKey(topic)); + }); + } }