From a2077b0f5e6154954795c6e39f5bb7935642d07d Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Mon, 25 Mar 2024 08:54:10 -0700 Subject: [PATCH 1/8] MINOR: AbstractConfig cleanup Signed-off-by: Greg Harris --- .../kafka/common/config/AbstractConfig.java | 50 ++++++++++++---- .../org/apache/kafka/common/utils/Utils.java | 18 ++++-- .../common/config/AbstractConfigTest.java | 57 +++++++++++++++++++ .../connect/mirror/MirrorClientConfig.java | 2 +- .../connect/mirror/MirrorMakerConfig.java | 2 +- .../apache/kafka/connect/runtime/Worker.java | 6 +- .../kafka/connect/runtime/WorkerConfig.java | 2 +- .../runtime/rest/RestServerConfig.java | 2 +- .../kafka/connect/runtime/WorkerTest.java | 3 + .../admin/BrokerApiVersionsCommand.scala | 2 +- .../controller/PartitionStateMachine.scala | 2 +- .../kafka/server/DynamicBrokerConfig.scala | 6 +- .../main/scala/kafka/server/KafkaConfig.scala | 2 +- 13 files changed, 128 insertions(+), 26 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java index 84bae97a03a8c..c598a28c0af06 100644 --- a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java +++ b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java @@ -25,6 +25,7 @@ import org.slf4j.LoggerFactory; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; @@ -33,6 +34,8 @@ import java.util.Set; import java.util.TreeMap; import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Predicate; +import java.util.stream.Collectors; /** * A convenient base class for configurations to extend. @@ -58,6 +61,8 @@ public class AbstractConfig { private final ConfigDef definition; + public static final String AUTOMATIC_CONFIG_PROVIDERS_PROPERTY = "org.apache.kafka.automatic.config.providers"; + public static final String CONFIG_PROVIDERS_CONFIG = "config.providers"; private static final String CONFIG_PROVIDERS_PARAM = ".param."; @@ -101,14 +106,11 @@ public class AbstractConfig { * the constructor to resolve any variables in {@code originals}; may be null or empty * @param doLog whether the configurations should be logged */ - @SuppressWarnings({"unchecked", "this-escape"}) + @SuppressWarnings({"this-escape"}) public AbstractConfig(ConfigDef definition, Map originals, Map configProviderProps, boolean doLog) { - /* check that all the keys are really strings */ - for (Map.Entry entry : originals.entrySet()) - if (!(entry.getKey() instanceof String)) - throw new ConfigException(entry.getKey().toString(), entry.getValue(), "Key must be a string."); + Map originalMap = Utils.castToStringObjectMap(originals); - this.originals = resolveConfigVariables(configProviderProps, (Map) originals); + this.originals = resolveConfigVariables(configProviderProps, originalMap); this.values = definition.parse(this.originals); Map configUpdates = postProcessParsedConfig(Collections.unmodifiableMap(this.values)); for (Map.Entry update : configUpdates.entrySet()) { @@ -521,6 +523,7 @@ private Map extractPotentialVariables(Map configMap) { private Map resolveConfigVariables(Map configProviderProps, Map originals) { Map providerConfigString; Map configProperties; + Predicate classNameFilter; Map resolvedOriginals = new HashMap<>(); // As variable configs are strings, parse the originals and obtain the potential variable configs. Map indirectVariables = extractPotentialVariables(originals); @@ -529,11 +532,13 @@ private Map extractPotentialVariables(Map configMap) { if (configProviderProps == null || configProviderProps.isEmpty()) { providerConfigString = indirectVariables; configProperties = originals; + classNameFilter = automaticConfigProvidersFilter(); } else { providerConfigString = extractPotentialVariables(configProviderProps); configProperties = configProviderProps; + classNameFilter = ignored -> true; } - Map providers = instantiateConfigProviders(providerConfigString, configProperties); + Map providers = instantiateConfigProviders(providerConfigString, configProperties, classNameFilter); if (!providers.isEmpty()) { ConfigTransformer configTransformer = new ConfigTransformer(providers); @@ -547,6 +552,17 @@ private Map extractPotentialVariables(Map configMap) { return new ResolvingMap<>(resolvedOriginals, originals); } + private Predicate automaticConfigProvidersFilter() { + String systemProperty = System.getProperty(AUTOMATIC_CONFIG_PROVIDERS_PROPERTY); + if (systemProperty == null) { + return ignored -> true; + } else { + return Arrays.stream(systemProperty.split(",")) + .map(String::trim) + .collect(Collectors.toSet())::contains; + } + } + private Map configProviderProperties(String configProviderPrefix, Map providerConfigProperties) { Map result = new HashMap<>(); for (Map.Entry entry : providerConfigProperties.entrySet()) { @@ -567,9 +583,14 @@ private Map configProviderProperties(String configProviderPrefix * * @param indirectConfigs The map of potential variable configs * @param providerConfigProperties The map of config provider configs - * @return map map of config provider name and its instance. + * @param classNameFilter Filter for config provider class names + * @return map of config provider name and its instance. */ - private Map instantiateConfigProviders(Map indirectConfigs, Map providerConfigProperties) { + private Map instantiateConfigProviders( + Map indirectConfigs, + Map providerConfigProperties, + Predicate classNameFilter + ) { final String configProviders = indirectConfigs.get(CONFIG_PROVIDERS_CONFIG); if (configProviders == null || configProviders.isEmpty()) { @@ -580,8 +601,15 @@ private Map instantiateConfigProviders(Map Map filterMap(final Map map, final Predicate propsToMap(Properties properties) { - Map map = new HashMap<>(properties.size()); - for (Map.Entry entry : properties.entrySet()) { + return castToStringObjectMap(properties); + } + + /** + * Cast a map with arbitrary type keys to be keyed on String. + * @param inputMap A map with unknown type keys + * @return A map with the same contents as the input map, but with String keys + * @throws ConfigException if any key is not a String + */ + public static Map castToStringObjectMap(Map inputMap) { + Map map = new HashMap<>(inputMap.size()); + for (Map.Entry entry : inputMap.entrySet()) { if (entry.getKey() instanceof String) { String k = (String) entry.getKey(); - map.put(k, properties.get(k)); + map.put(k, entry.getValue()); } else { - throw new ConfigException(entry.getKey().toString(), entry.getValue(), "Key must be a string."); + throw new ConfigException(String.valueOf(entry.getKey()), entry.getValue(), "Key must be a string."); } } return map; diff --git a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java index 5859dc1dc1278..65aee07483f14 100644 --- a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java +++ b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java @@ -27,6 +27,8 @@ import org.apache.kafka.common.config.provider.MockVaultConfigProvider; import org.apache.kafka.common.config.provider.MockFileConfigProvider; import org.apache.kafka.test.MockConsumerInterceptor; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import java.util.Arrays; @@ -46,6 +48,23 @@ public class AbstractConfigTest { + private String propertyValue; + + @BeforeEach + public void setup() { + propertyValue = System.getProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY); + System.clearProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY); + } + + @AfterEach + public void teardown() { + if (propertyValue != null) { + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, propertyValue); + } else { + System.clearProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY); + } + } + @Test public void testConfiguredInstances() { testValidInputs(" "); @@ -389,6 +408,23 @@ public void testOriginalsWithConfigProvidersProps() { MockFileConfigProvider.assertClosed(id); } + @Test + public void testOriginalsWithConfigProvidersPropsExcluded() { + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockVaultConfigProvider.class.getName()); + Properties props = new Properties(); + + // Test Case: VaultConfigProvider is automatic, but FileConfigProvider is specified + props.put("config.providers", "file"); + props.put("config.providers.file.class", MockFileConfigProvider.class.getName()); + String id = UUID.randomUUID().toString(); + props.put("config.providers.file.param.testId", id); + props.put("prefix.ssl.truststore.location.number", 5); + props.put("sasl.kerberos.service.name", "service name"); + props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); + props.put("sasl.kerberos.password", "${file:/usr/kerberos:password}"); + assertThrows(ConfigException.class, () -> new TestIndirectConfigResolution(props, Collections.emptyMap())); + } + @Test public void testConfigProvidersPropsAsParam() { // Test Case: Valid Test Case for ConfigProviders as a separate variable @@ -461,6 +497,25 @@ public void testAutoConfigResolutionWithInvalidConfigProviderClass() { } } + @Test + public void testAutoConfigResolutionWithInvalidConfigProviderClassExcluded() { + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, ""); + // Test Case: Invalid class for Config Provider + Properties props = new Properties(); + props.put("config.providers", "file"); + props.put("config.providers.file.class", + "org.apache.kafka.common.config.provider.InvalidConfigProvider"); + props.put("testKey", "${test:/foo/bar/testpath:testKey}"); + try { + new TestIndirectConfigResolution(props, Collections.emptyMap()); + fail("Expected a config exception due to invalid props :" + props); + } catch (KafkaException e) { + // deliver the disallowed message first to prevent probing the classloader + assertTrue(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); + // this is good + } + } + @Test public void testAutoConfigResolutionWithMissingConfigProvider() { // Test Case: Config Provider for a variable missing in config file. @@ -502,6 +557,8 @@ public void testAutoConfigResolutionWithDuplicateConfigProvider() { @Test public void testConfigProviderConfigurationWithConfigParams() { + // should have no effect + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockFileConfigProvider.class.getName()); // Test Case: Valid Test Case With Multiple ConfigProviders as a separate variable Properties providers = new Properties(); providers.put("config.providers", "vault"); diff --git a/connect/mirror-client/src/main/java/org/apache/kafka/connect/mirror/MirrorClientConfig.java b/connect/mirror-client/src/main/java/org/apache/kafka/connect/mirror/MirrorClientConfig.java index 477459895c56c..053e594fbeb1d 100644 --- a/connect/mirror-client/src/main/java/org/apache/kafka/connect/mirror/MirrorClientConfig.java +++ b/connect/mirror-client/src/main/java/org/apache/kafka/connect/mirror/MirrorClientConfig.java @@ -76,7 +76,7 @@ public class MirrorClientConfig extends AbstractConfig { public static final String PRODUCER_CLIENT_PREFIX = "producer."; MirrorClientConfig(Map props) { - super(CONFIG_DEF, props, true); + super(CONFIG_DEF, props, Utils.castToStringObjectMap(props), true); } public ReplicationPolicy replicationPolicy() { diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java index fd672f56a6ce8..69978da026f5a 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java @@ -92,7 +92,7 @@ public class MirrorMakerConfig extends AbstractConfig { @SuppressWarnings("this-escape") public MirrorMakerConfig(Map props) { - super(config(), props, true); + super(config(), props, Utils.castToStringObjectMap(props), true); plugins = new Plugins(originalsStrings()); rawProperties = new HashMap<>(props); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index 2ce09ee28b6df..35c202a22debe 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -37,6 +37,7 @@ import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.MetricNameTemplate; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.AbstractConfig; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigValue; import org.apache.kafka.common.config.provider.ConfigProvider; @@ -925,10 +926,13 @@ static Map adminConfigs(String connName, // Ignore configs that begin with "admin." since those will be added next (with the prefix stripped) // and those that begin with "producer." and "consumer.", since we know they aren't intended for // the admin client + // Also ignore the config.providers configurations because the worker-configured ConfigProviders should + // already have been evaluated via the trusted WorkerConfig constructor Map nonPrefixedWorkerConfigs = config.originals().entrySet().stream() .filter(e -> !e.getKey().startsWith("admin.") && !e.getKey().startsWith("producer.") - && !e.getKey().startsWith("consumer.")) + && !e.getKey().startsWith("consumer.") + && !e.getKey().startsWith(AbstractConfig.CONFIG_PROVIDERS_CONFIG)) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, config.bootstrapServers()); adminProps.put(AdminClientConfig.CLIENT_ID_CONFIG, defaultClientId); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index fe28918a29a2b..5cb9a44fb97b8 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -438,7 +438,7 @@ public static PluginDiscoveryMode pluginDiscovery(Map props) { @SuppressWarnings("this-escape") public WorkerConfig(ConfigDef definition, Map props) { - super(definition, props); + super(definition, props, Utils.castToStringObjectMap(props), true); logInternalConverterRemovalWarnings(props); logPluginPathConfigProviderWarning(props); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/rest/RestServerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/rest/RestServerConfig.java index 0d6d06a4a596c..4b8b5acf93519 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/rest/RestServerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/rest/RestServerConfig.java @@ -258,7 +258,7 @@ public Integer rebalanceTimeoutMs() { } protected RestServerConfig(ConfigDef configDef, Map props) { - super(configDef, props); + super(configDef, props, Utils.castToStringObjectMap(props), true); } // Visible for testing diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index c09fa093fb32e..84dc8e6b088bf 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -1167,6 +1167,9 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { Map connConfig = Collections.singletonMap("metadata.max.age.ms", "10000"); Map expectedConfigs = new HashMap<>(workerProps); + workerProps.keySet() + .stream().filter(key -> key.startsWith(AbstractConfig.CONFIG_PROVIDERS_CONFIG)) + .forEach(expectedConfigs::remove); expectedConfigs.put("bootstrap.servers", "localhost:9092"); expectedConfigs.put("client.id", "testid"); expectedConfigs.put("metadata.max.age.ms", "10000"); diff --git a/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala b/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala index 45df36f7a1f87..6cb273f066f21 100644 --- a/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala +++ b/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala @@ -261,7 +261,7 @@ object BrokerApiVersionsCommand { config } - class AdminConfig(originals: Map[_,_]) extends AbstractConfig(AdminConfigDef, originals.asJava, false) + class AdminConfig(originals: Map[_,_]) extends AbstractConfig(AdminConfigDef, originals.asJava, Utils.castToStringObjectMap(originals.asJava), false) def create(props: Properties): AdminClient = create(props.asScala.toMap) diff --git a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala index 5dedad426406b..51634291b6d56 100755 --- a/core/src/main/scala/kafka/controller/PartitionStateMachine.scala +++ b/core/src/main/scala/kafka/controller/PartitionStateMachine.scala @@ -486,7 +486,7 @@ class ZkPartitionStateMachine(config: KafkaConfig, } else { val (logConfigs, failed) = zkClient.getLogConfigs( partitionsWithNoLiveInSyncReplicas.iterator.map { case (partition, _) => partition.topic }.toSet, - config.originals() + config.extractLogConfigMap ) partitionsWithNoLiveInSyncReplicas.map { case (partition, leaderAndIsr) => diff --git a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala index 74880971bc382..50f2e3111d50b 100755 --- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala +++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala @@ -197,7 +197,7 @@ object DynamicBrokerConfig { private[server] def resolveVariableConfigs(propsOriginal: Properties): Properties = { val props = new Properties - val config = new AbstractConfig(new ConfigDef(), propsOriginal, false) + val config = new AbstractConfig(new ConfigDef(), propsOriginal, Utils.castToStringObjectMap(propsOriginal), false) config.originals.forEach { (key, value) => if (!key.startsWith(AbstractConfig.CONFIG_PROVIDERS_CONFIG)) { props.put(key, value) @@ -740,13 +740,13 @@ class DynamicLogConfig(logManager: LogManager, server: KafkaBroker) extends Brok val originalLogConfig = logManager.currentDefaultConfig val originalUncleanLeaderElectionEnable = originalLogConfig.uncleanLeaderElectionEnable val newBrokerDefaults = new util.HashMap[String, Object](originalLogConfig.originals) - newConfig.valuesFromThisConfig.forEach { (k, v) => + newConfig.extractLogConfigMap.forEach { (k, v) => if (DynamicLogConfig.ReconfigurableConfigs.contains(k)) { DynamicLogConfig.KafkaConfigToLogConfigName.get(k).foreach { configName => if (v == null) newBrokerDefaults.remove(configName) else - newBrokerDefaults.put(configName, v.asInstanceOf[AnyRef]) + newBrokerDefaults.put(configName, v) } } } diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index eaddb047da8d5..af3a00d460def 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -1374,7 +1374,7 @@ object KafkaConfig { } class KafkaConfig private(doLog: Boolean, val props: java.util.Map[_, _], dynamicConfigOverride: Option[DynamicBrokerConfig]) - extends AbstractConfig(KafkaConfig.configDef, props, doLog) with Logging { + extends AbstractConfig(KafkaConfig.configDef, props, Utils.castToStringObjectMap(props), doLog) with Logging { def this(props: java.util.Map[_, _]) = this(true, KafkaConfig.populateSynonyms(props), None) def this(props: java.util.Map[_, _], doLog: Boolean) = this(doLog, KafkaConfig.populateSynonyms(props), None) From a12c19e08dcd8fd3dba790de65dc8722c0ba0092 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Tue, 26 Mar 2024 13:59:40 -0700 Subject: [PATCH 2/8] fixup: add tests for including invalid classes, lists of classes with spaces Signed-off-by: Greg Harris --- .../common/config/AbstractConfigTest.java | 52 ++++++++++++++++--- 1 file changed, 45 insertions(+), 7 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java index 65aee07483f14..e3c66971e7615 100644 --- a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java +++ b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; +import org.apache.kafka.common.config.provider.FileConfigProvider; import org.apache.kafka.common.config.types.Password; import org.apache.kafka.common.metrics.FakeMetricsReporter; import org.apache.kafka.common.metrics.JmxReporter; @@ -410,10 +411,10 @@ public void testOriginalsWithConfigProvidersProps() { @Test public void testOriginalsWithConfigProvidersPropsExcluded() { - System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockVaultConfigProvider.class.getName()); + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockVaultConfigProvider.class.getName() + " , " + FileConfigProvider.class.getName()); Properties props = new Properties(); - // Test Case: VaultConfigProvider is automatic, but FileConfigProvider is specified + // Test Case: Config provider that is not an allowed class props.put("config.providers", "file"); props.put("config.providers.file.class", MockFileConfigProvider.class.getName()); String id = UUID.randomUUID().toString(); @@ -425,6 +426,26 @@ public void testOriginalsWithConfigProvidersPropsExcluded() { assertThrows(ConfigException.class, () -> new TestIndirectConfigResolution(props, Collections.emptyMap())); } + @Test + public void testOriginalsWithConfigProvidersPropsIncluded() { + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockFileConfigProvider.class.getName() + " , " + FileConfigProvider.class.getName()); + Properties props = new Properties(); + + // Test Case: Config provider that is an allowed class + props.put("config.providers", "file"); + props.put("config.providers.file.class", MockFileConfigProvider.class.getName()); + String id = UUID.randomUUID().toString(); + props.put("config.providers.file.param.testId", id); + props.put("prefix.ssl.truststore.location.number", 5); + props.put("sasl.kerberos.service.name", "service name"); + props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); + props.put("sasl.kerberos.password", "${file:/usr/kerberos:password}"); + TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, Collections.emptyMap()); + assertEquals("testKey", config.originals().get("sasl.kerberos.key")); + assertEquals("randomPassword", config.originals().get("sasl.kerberos.password")); + MockFileConfigProvider.assertClosed(id); + } + @Test public void testConfigProvidersPropsAsParam() { // Test Case: Valid Test Case for ConfigProviders as a separate variable @@ -499,12 +520,12 @@ public void testAutoConfigResolutionWithInvalidConfigProviderClass() { @Test public void testAutoConfigResolutionWithInvalidConfigProviderClassExcluded() { + String invalidConfigProvider = "org.apache.kafka.common.config.provider.InvalidConfigProvider"; System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, ""); - // Test Case: Invalid class for Config Provider + // Test Case: Any config provider specified while the system property is empty Properties props = new Properties(); props.put("config.providers", "file"); - props.put("config.providers.file.class", - "org.apache.kafka.common.config.provider.InvalidConfigProvider"); + props.put("config.providers.file.class", invalidConfigProvider); props.put("testKey", "${test:/foo/bar/testpath:testKey}"); try { new TestIndirectConfigResolution(props, Collections.emptyMap()); @@ -512,7 +533,24 @@ public void testAutoConfigResolutionWithInvalidConfigProviderClassExcluded() { } catch (KafkaException e) { // deliver the disallowed message first to prevent probing the classloader assertTrue(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); - // this is good + } + } + + @Test + public void testAutoConfigResolutionWithInvalidConfigProviderClassIncluded() { + String invalidConfigProvider = "org.apache.kafka.common.config.provider.InvalidConfigProvider"; + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, invalidConfigProvider); + // Test Case: Invalid config provider specified, but is also included in the system property + Properties props = new Properties(); + props.put("config.providers", "file"); + props.put("config.providers.file.class", invalidConfigProvider); + props.put("testKey", "${test:/foo/bar/testpath:testKey}"); + try { + new TestIndirectConfigResolution(props, Collections.emptyMap()); + fail("Expected a config exception due to invalid props :" + props); + } catch (KafkaException e) { + // deliver the disallowed message first to prevent probing the classloader + assertFalse(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); } } @@ -559,7 +597,7 @@ public void testAutoConfigResolutionWithDuplicateConfigProvider() { public void testConfigProviderConfigurationWithConfigParams() { // should have no effect System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockFileConfigProvider.class.getName()); - // Test Case: Valid Test Case With Multiple ConfigProviders as a separate variable + // Test Case: Specify a config provider not allowed, but passed via the trusted providers argument Properties providers = new Properties(); providers.put("config.providers", "vault"); providers.put("config.providers.vault.class", MockVaultConfigProvider.class.getName()); From d6da916a550a4189756001c1ab85f1ff66d4e0a2 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Tue, 26 Mar 2024 15:00:04 -0700 Subject: [PATCH 3/8] Use assertThrows in tests Signed-off-by: Greg Harris --- .../common/config/AbstractConfigTest.java | 32 ++++--------------- 1 file changed, 6 insertions(+), 26 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java index e3c66971e7615..d49eb4c718eb9 100644 --- a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java +++ b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java @@ -274,12 +274,7 @@ private void testInvalidInputs(String configValue) { Properties props = new Properties(); props.put(TestConfig.METRIC_REPORTER_CLASSES_CONFIG, configValue); TestConfig config = new TestConfig(props); - try { - config.getConfiguredInstances(TestConfig.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class); - fail("Expected a config exception due to invalid props :" + props); - } catch (KafkaException e) { - // this is good - } + assertThrows(KafkaException.class, () -> config.getConfiguredInstances(TestConfig.METRIC_REPORTER_CLASSES_CONFIG, MetricsReporter.class)); } @Test @@ -510,12 +505,7 @@ public void testAutoConfigResolutionWithInvalidConfigProviderClass() { props.put("config.providers.file.class", "org.apache.kafka.common.config.provider.InvalidConfigProvider"); props.put("testKey", "${test:/foo/bar/testpath:testKey}"); - try { - new TestIndirectConfigResolution(props); - fail("Expected a config exception due to invalid props :" + props); - } catch (KafkaException e) { - // this is good - } + assertThrows(KafkaException.class, () -> new TestIndirectConfigResolution(props)); } @Test @@ -527,13 +517,8 @@ public void testAutoConfigResolutionWithInvalidConfigProviderClassExcluded() { props.put("config.providers", "file"); props.put("config.providers.file.class", invalidConfigProvider); props.put("testKey", "${test:/foo/bar/testpath:testKey}"); - try { - new TestIndirectConfigResolution(props, Collections.emptyMap()); - fail("Expected a config exception due to invalid props :" + props); - } catch (KafkaException e) { - // deliver the disallowed message first to prevent probing the classloader - assertTrue(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); - } + KafkaException e = assertThrows(KafkaException.class, () -> new TestIndirectConfigResolution(props, Collections.emptyMap())); + assertTrue(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); } @Test @@ -545,13 +530,8 @@ public void testAutoConfigResolutionWithInvalidConfigProviderClassIncluded() { props.put("config.providers", "file"); props.put("config.providers.file.class", invalidConfigProvider); props.put("testKey", "${test:/foo/bar/testpath:testKey}"); - try { - new TestIndirectConfigResolution(props, Collections.emptyMap()); - fail("Expected a config exception due to invalid props :" + props); - } catch (KafkaException e) { - // deliver the disallowed message first to prevent probing the classloader - assertFalse(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); - } + KafkaException e = assertThrows(KafkaException.class, () -> new TestIndirectConfigResolution(props, Collections.emptyMap())); + assertFalse(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); } @Test From 568050b1578d240a85c34b48750915107fb29f02 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Wed, 27 Mar 2024 09:09:22 -0700 Subject: [PATCH 4/8] fixup: remove test util fn, whitespace, refactor Worker prefix test Signed-off-by: Greg Harris --- .../kafka/common/config/AbstractConfig.java | 1 - .../common/config/AbstractConfigTest.java | 21 ++++++------------- .../apache/kafka/connect/runtime/Worker.java | 7 +++---- 3 files changed, 9 insertions(+), 20 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java index c598a28c0af06..aeb7f07a29c7b 100644 --- a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java +++ b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java @@ -610,7 +610,6 @@ private Map instantiateConfigProviders( + AUTOMATIC_CONFIG_PROVIDERS_PROPERTY + "' to allow " + providerClassName); } } - } // Instantiate Config Providers Map configProviderInstances = new HashMap<>(); diff --git a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java index d49eb4c718eb9..bf018aebbfcb9 100644 --- a/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java +++ b/clients/src/test/java/org/apache/kafka/common/config/AbstractConfigTest.java @@ -27,6 +27,7 @@ import org.apache.kafka.common.security.TestSecurityConfig; import org.apache.kafka.common.config.provider.MockVaultConfigProvider; import org.apache.kafka.common.config.provider.MockFileConfigProvider; +import org.apache.kafka.common.utils.Utils; import org.apache.kafka.test.MockConsumerInterceptor; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -364,16 +365,6 @@ protected Class findClass(String name) throws ClassNotFoundException { } } - @SuppressWarnings("unchecked") - public Map convertPropertiesToMap(Map props) { - for (Map.Entry entry : props.entrySet()) { - if (!(entry.getKey() instanceof String)) - throw new ConfigException(entry.getKey().toString(), entry.getValue(), - "Key must be a string."); - } - return (Map) props; - } - @Test public void testOriginalWithOverrides() { Properties props = new Properties(); @@ -452,7 +443,7 @@ public void testConfigProvidersPropsAsParam() { Properties props = new Properties(); props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); props.put("sasl.kerberos.password", "${file:/usr/kerberos:password}"); - TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, convertPropertiesToMap(providers)); + TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, Utils.castToStringObjectMap(providers)); assertEquals("testKey", config.originals().get("sasl.kerberos.key")); assertEquals("randomPassword", config.originals().get("sasl.kerberos.password")); MockFileConfigProvider.assertClosed(id); @@ -469,7 +460,7 @@ public void testImmutableOriginalsWithConfigProvidersProps() { Properties props = new Properties(); props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); Map immutableMap = Collections.unmodifiableMap(props); - Map provMap = convertPropertiesToMap(providers); + Map provMap = Utils.castToStringObjectMap(providers); TestIndirectConfigResolution config = new TestIndirectConfigResolution(immutableMap, provMap); assertEquals("testKey", config.originals().get("sasl.kerberos.key")); MockFileConfigProvider.assertClosed(id); @@ -489,7 +480,7 @@ public void testAutoConfigResolutionWithMultipleConfigProviders() { props.put("sasl.kerberos.password", "${file:/usr/kerberos:password}"); props.put("sasl.truststore.key", "${vault:/usr/truststore:truststoreKey}"); props.put("sasl.truststore.password", "${vault:/usr/truststore:truststorePassword}"); - TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, convertPropertiesToMap(providers)); + TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, Utils.castToStringObjectMap(providers)); assertEquals("testKey", config.originals().get("sasl.kerberos.key")); assertEquals("randomPassword", config.originals().get("sasl.kerberos.password")); assertEquals("testTruststoreKey", config.originals().get("sasl.truststore.key")); @@ -569,7 +560,7 @@ public void testAutoConfigResolutionWithDuplicateConfigProvider() { props.put("config.providers", "file"); props.put("config.providers.file.class", MockVaultConfigProvider.class.getName()); - TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, convertPropertiesToMap(providers)); + TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, Utils.castToStringObjectMap(providers)); assertEquals("${file:/usr/kerberos:key}", config.originals().get("sasl.kerberos.key")); } @@ -587,7 +578,7 @@ public void testConfigProviderConfigurationWithConfigParams() { props.put("sasl.truststore.key", "${vault:/usr/truststore:truststoreKey}"); props.put("sasl.truststore.password", "${vault:/usr/truststore:truststorePassword}"); props.put("sasl.truststore.location", "${vault:/usr/truststore:truststoreLocation}"); - TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, convertPropertiesToMap(providers)); + TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, Utils.castToStringObjectMap(providers)); assertEquals("/usr/vault", config.originals().get("sasl.truststore.location")); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index 35c202a22debe..a9e5c3db78a44 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -103,6 +103,7 @@ import java.time.Duration; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.HashMap; @@ -928,11 +929,9 @@ static Map adminConfigs(String connName, // the admin client // Also ignore the config.providers configurations because the worker-configured ConfigProviders should // already have been evaluated via the trusted WorkerConfig constructor + List excludedPrefixes = Arrays.asList("admin.", "producer.", "consumer.", AbstractConfig.CONFIG_PROVIDERS_CONFIG); Map nonPrefixedWorkerConfigs = config.originals().entrySet().stream() - .filter(e -> !e.getKey().startsWith("admin.") - && !e.getKey().startsWith("producer.") - && !e.getKey().startsWith("consumer.") - && !e.getKey().startsWith(AbstractConfig.CONFIG_PROVIDERS_CONFIG)) + .filter(e -> excludedPrefixes.stream().noneMatch(p -> e.getKey().startsWith(p))) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, config.bootstrapServers()); adminProps.put(AdminClientConfig.CLIENT_ID_CONFIG, defaultClientId); From 5efcf02201956b0bfb2845022bd9d9ad8e02fb09 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Wed, 27 Mar 2024 11:57:21 -0700 Subject: [PATCH 5/8] fixup: Override WorkerConfig.originals to hide already-processed configs from downstream users Signed-off-by: Greg Harris --- .../org/apache/kafka/connect/cli/AbstractConnectCli.java | 2 +- .../java/org/apache/kafka/connect/runtime/Worker.java | 9 +++------ .../org/apache/kafka/connect/runtime/WorkerConfig.java | 6 ++++++ 3 files changed, 10 insertions(+), 7 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/AbstractConnectCli.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/AbstractConnectCli.java index de666c7bd60a2..c770a19624a1d 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/AbstractConnectCli.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/AbstractConnectCli.java @@ -125,7 +125,7 @@ public Connect startConnect(Map workerProps, String... extraArgs RestClient restClient = new RestClient(config); - ConnectRestServer restServer = new ConnectRestServer(config.rebalanceTimeout(), restClient, workerProps); + ConnectRestServer restServer = new ConnectRestServer(config.rebalanceTimeout(), restClient, config.originals()); restServer.initializeServer(); URI advertisedUrl = restServer.advertisedUrl(); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index a9e5c3db78a44..2ce09ee28b6df 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -37,7 +37,6 @@ import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.MetricNameTemplate; import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.config.AbstractConfig; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigValue; import org.apache.kafka.common.config.provider.ConfigProvider; @@ -103,7 +102,6 @@ import java.time.Duration; import java.util.ArrayList; -import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.HashMap; @@ -927,11 +925,10 @@ static Map adminConfigs(String connName, // Ignore configs that begin with "admin." since those will be added next (with the prefix stripped) // and those that begin with "producer." and "consumer.", since we know they aren't intended for // the admin client - // Also ignore the config.providers configurations because the worker-configured ConfigProviders should - // already have been evaluated via the trusted WorkerConfig constructor - List excludedPrefixes = Arrays.asList("admin.", "producer.", "consumer.", AbstractConfig.CONFIG_PROVIDERS_CONFIG); Map nonPrefixedWorkerConfigs = config.originals().entrySet().stream() - .filter(e -> excludedPrefixes.stream().noneMatch(p -> e.getKey().startsWith(p))) + .filter(e -> !e.getKey().startsWith("admin.") + && !e.getKey().startsWith("producer.") + && !e.getKey().startsWith("consumer.")) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, config.bootstrapServers()); adminProps.put(AdminClientConfig.CLIENT_ID_CONFIG, defaultClientId); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 5cb9a44fb97b8..6de49ebd55084 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -443,4 +443,10 @@ public WorkerConfig(ConfigDef definition, Map props) { logPluginPathConfigProviderWarning(props); } + @Override + public Map originals() { + Map map = super.originals(); + map.remove(AbstractConfig.CONFIG_PROVIDERS_CONFIG); + return map; + } } From fe0d42eb69c158831d18a6c8aff0026149cd77a2 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Wed, 27 Mar 2024 14:29:16 -0700 Subject: [PATCH 6/8] fixup: revert mirror changes Signed-off-by: Greg Harris --- .../java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java index 69978da026f5a..fd672f56a6ce8 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorMakerConfig.java @@ -92,7 +92,7 @@ public class MirrorMakerConfig extends AbstractConfig { @SuppressWarnings("this-escape") public MirrorMakerConfig(Map props) { - super(config(), props, Utils.castToStringObjectMap(props), true); + super(config(), props, true); plugins = new Plugins(originalsStrings()); rawProperties = new HashMap<>(props); From 997f0260e6ae606588d1680a8626d124a8821403 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Wed, 27 Mar 2024 18:06:29 -0700 Subject: [PATCH 7/8] fixup: failing test due to Worker change Signed-off-by: Greg Harris --- .../java/org/apache/kafka/connect/runtime/WorkerTest.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index 84dc8e6b088bf..4579794a2c4ff 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -1167,9 +1167,7 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { Map connConfig = Collections.singletonMap("metadata.max.age.ms", "10000"); Map expectedConfigs = new HashMap<>(workerProps); - workerProps.keySet() - .stream().filter(key -> key.startsWith(AbstractConfig.CONFIG_PROVIDERS_CONFIG)) - .forEach(expectedConfigs::remove); + expectedConfigs.remove(AbstractConfig.CONFIG_PROVIDERS_CONFIG); expectedConfigs.put("bootstrap.servers", "localhost:9092"); expectedConfigs.put("client.id", "testid"); expectedConfigs.put("metadata.max.age.ms", "10000"); From 352d7e74202967c0a9378639690323be258a223e Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Thu, 28 Mar 2024 09:48:24 -0700 Subject: [PATCH 8/8] fixup: broken PartitionStateMachineTest mocks Signed-off-by: Greg Harris --- .../unit/kafka/controller/PartitionStateMachineTest.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala b/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala index 10cbe58904564..183e8657e0d44 100644 --- a/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala +++ b/core/src/test/scala/unit/kafka/controller/PartitionStateMachineTest.scala @@ -258,7 +258,7 @@ class PartitionStateMachineTest { .thenReturn(Seq(GetDataResponse(Code.OK, null, Some(partition), TopicPartitionStateZNode.encode(leaderIsrAndControllerEpoch), stat, ResponseMetadata(0, 0)))) - when(mockZkClient.getLogConfigs(Set.empty, config.originals())) + when(mockZkClient.getLogConfigs(Set.empty, config.extractLogConfigMap)) .thenReturn((Map(partition.topic -> new LogConfig(new Properties)), Map.empty[String, Exception])) val leaderAndIsrAfterElection = leaderAndIsr.newLeader(brokerId) val updatedLeaderAndIsr = leaderAndIsrAfterElection.withPartitionEpoch(2) @@ -434,7 +434,7 @@ class PartitionStateMachineTest { } prepareMockToGetTopicPartitionsStatesRaw() def prepareMockToGetLogConfigs(): Unit = { - when(mockZkClient.getLogConfigs(Set.empty, config.originals())).thenReturn((Map.empty[String, LogConfig], Map.empty[String, Exception])) + when(mockZkClient.getLogConfigs(Set.empty, config.extractLogConfigMap)).thenReturn((Map.empty[String, LogConfig], Map.empty[String, Exception])) } prepareMockToGetLogConfigs()