From 6906ec7e1cfe11c63ec1cac77fad2bf8b075cccb Mon Sep 17 00:00:00 2001 From: Lei Zhiyuan Date: Mon, 19 Sep 2022 20:33:05 +0800 Subject: [PATCH 1/2] [fix][broker] Fix if dynamicConfig item in ZK do not exist in broker cause NPE (#17705) --- .../org/apache/pulsar/broker/service/BrokerService.java | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 77cc8f11ff553..de4b70c7046c9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2152,7 +2152,12 @@ private void handleDynamicConfigurationUpdates() { } Map data = optMap.get(); data.forEach((configKey, value) -> { - Field configField = dynamicConfigurationMap.get(configKey).field; + ConfigField configFieldWrapper = dynamicConfigurationMap.get(configKey); + if (configFieldWrapper == null) { + log.warn("{} does not exist in dynamicConfigurationMap, skip this config.", configKey); + return; + } + Field configField = configFieldWrapper.field; Object newValue = FieldParser.value(data.get(configKey), configField); if (configField != null) { Consumer listener = configRegisteredListeners.get(configKey); From c6e9f990d2ba33ed9d6584070d151e14814c650a Mon Sep 17 00:00:00 2001 From: Ruguo Yu Date: Sat, 29 Oct 2022 01:27:15 +0800 Subject: [PATCH 2/2] [fix][Authorization] Modify authorization to update topic's properties --- .../PulsarAuthorizationProvider.java | 1 + .../admin/impl/PersistentTopicsBase.java | 2 +- .../AuthorizationProducerConsumerTest.java | 58 +++++++++++++++++++ .../common/policies/data/TopicOperation.java | 1 + 4 files changed, 61 insertions(+), 1 deletion(-) 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 ab92881f062ed..a43591dd1d94f 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 @@ -554,6 +554,7 @@ public CompletableFuture allowTopicOperationAsync(TopicName topicName, case OFFLOAD: case UNLOAD: case DELETE_METADATA: + case UPDATE_METADATA: case ADD_BUNDLE_RANGE: case GET_BUNDLE_RANGE: case DELETE_BUNDLE_RANGE: 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 d28a4113cc27b..ce3fdd79edb8e 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 @@ -619,7 +619,7 @@ protected CompletableFuture internalUpdatePropertiesAsync(boolean authorit return CompletableFuture.completedFuture(null); } return validateTopicOwnershipAsync(topicName, authoritative) - .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.PRODUCE)) + .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.UPDATE_METADATA)) .thenCompose(__ -> { if (topicName.isPartitioned()) { return internalUpdateNonPartitionedTopicProperties(properties); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java index 29f20fce25402..0ce3b7df07d1f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/AuthorizationProducerConsumerTest.java @@ -30,6 +30,7 @@ import com.google.common.collect.Sets; import java.io.IOException; import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -548,6 +549,63 @@ public void testSchemaCompatibilityStrategyPermission() throws Exception { log.info("-- Exiting {} test --", methodName); } + @Test + public void testUpdateTopicPropertiesAuthorization() throws Exception { + log.info("-- Starting {} test --", methodName); + cleanup(); + conf.setAuthorizationProvider(PulsarAuthorizationProvider.class.getName()); + setup(); + + final String tenantRole = "tenant-role"; + final String generalRole = "general-role"; + final String namespace = "my-property/my-ns-sub-auth"; + final String topicName = "persistent://" + namespace + "/my-topic"; + clientAuthProviderSupportedRoles.add(generalRole); + + Authentication superUserAuthentication = new ClientAuthentication("superUser"); + @Cleanup + PulsarAdmin superAdmin = spy(PulsarAdmin.builder().serviceHttpUrl(brokerUrl.toString()) + .authentication(superUserAuthentication).build()); + + Authentication tenantAdminAuthentication = new ClientAuthentication(tenantRole); + @Cleanup + PulsarAdmin tenantAdmin = spy(PulsarAdmin.builder().serviceHttpUrl(brokerUrl.toString()) + .authentication(tenantAdminAuthentication).build()); + + Authentication generalAuthentication = new ClientAuthentication(generalRole); + @Cleanup + PulsarAdmin generalAdmin = spy(PulsarAdmin.builder().serviceHttpUrl(brokerUrl.toString()) + .authentication(generalAuthentication).build()); + + superAdmin.clusters().createCluster("test", + ClusterData.builder().serviceUrl(brokerUrl.toString()).build()); + superAdmin.tenants().createTenant("my-property", + new TenantInfoImpl(Sets.newHashSet(tenantRole), Sets.newHashSet("test"))); + superAdmin.namespaces().createNamespace(namespace, Sets.newHashSet("test")); + superAdmin.topics().createPartitionedTopic(topicName, 1); + + Map topicProperties = new HashMap(); + topicProperties.put("key1", "value1"); + // superUser and tenantAdminatrator have authorization to operate topic's properties + superAdmin.topics().updateProperties(topicName, topicProperties); + Awaitility.await().untilAsserted(() -> assertEquals( + superAdmin.topics().getProperties(topicName).get("key1"), "value1")); + topicProperties.put("key1", "value2"); + tenantAdmin.topics().updateProperties(topicName, topicProperties); + Awaitility.await().untilAsserted(() -> assertEquals( + tenantAdmin.topics().getProperties(topicName).get("key1"), "value2")); + + // generalRole doesn't have authorization to update topic's properties so it will fail + try { + generalAdmin.topics().updateProperties(topicName, topicProperties); + fail("should have failed with authorization exception"); + } catch (Exception e) { + assertTrue(e.getMessage().startsWith( + "Unauthorized to validateTopicOperation for operation [UPDATE_METADATA]")); + } + + log.info("-- Exiting {} test --", methodName); + } @Test public void testSubscriptionPrefixAuthorization() throws Exception { log.info("-- Starting {} test --", methodName); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/TopicOperation.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/TopicOperation.java index 340d7f006f1e5..ad16b58e0136f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/TopicOperation.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/TopicOperation.java @@ -50,6 +50,7 @@ public enum TopicOperation { GET_STATS, GET_METADATA, DELETE_METADATA, + UPDATE_METADATA, GET_BACKLOG_SIZE, SET_REPLICATED_SUBSCRIPTION_STATUS,