From fbbaab694872c1702af11f272a478e3b1db44142 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Thu, 28 Mar 2024 13:27:41 -0700 Subject: [PATCH] MINOR: AbstractConfig cleanup (#15597) Signed-off-by: Greg Harris Reviewers: Chris Egerton , Mickael Maison , Omnia G H Ibrahim , Matthias J. Sax --- .../kafka/common/config/AbstractConfig.java | 50 +++++-- .../org/apache/kafka/common/utils/Utils.java | 18 ++- .../common/config/AbstractConfigTest.java | 122 ++++++++++++++---- .../connect/mirror/MirrorClientConfig.java | 2 +- .../kafka/connect/runtime/WorkerConfig.java | 9 +- .../kafka/connect/runtime/WorkerTest.java | 1 + .../admin/BrokerApiVersionsCommand.scala | 2 +- .../kafka/server/DynamicBrokerConfig.scala | 2 +- .../main/scala/kafka/server/KafkaConfig.scala | 2 +- 9 files changed, 159 insertions(+), 49 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 1363716331196..14b68ca305f4e 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.apache.kafka.common.config.provider.ConfigProvider; 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,10 @@ 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") 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()) { @@ -528,6 +529,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); @@ -536,11 +538,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); @@ -554,6 +558,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()) { @@ -574,9 +589,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()) { @@ -587,9 +607,15 @@ private Map instantiateConfigProviders(Map configProviderInstances = new HashMap<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index a9c510bac3f32..c8b47cf739511 100755 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1387,13 +1387,23 @@ public static 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 c0c6f8cee37c7..2b830ef2afca5 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; @@ -26,7 +27,10 @@ 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; import org.junit.jupiter.api.Test; import java.util.Arrays; @@ -46,6 +50,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(""); @@ -252,12 +273,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 @@ -347,16 +363,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(); @@ -387,6 +393,43 @@ public void testOriginalsWithConfigProvidersProps() { MockFileConfigProvider.assertClosed(id); } + @Test + public void testOriginalsWithConfigProvidersPropsExcluded() { + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockVaultConfigProvider.class.getName() + " , " + FileConfigProvider.class.getName()); + Properties props = new Properties(); + + // 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(); + 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 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 @@ -398,7 +441,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); @@ -415,7 +458,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); @@ -435,7 +478,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")); @@ -451,12 +494,33 @@ 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 + public void testAutoConfigResolutionWithInvalidConfigProviderClassExcluded() { + String invalidConfigProvider = "org.apache.kafka.common.config.provider.InvalidConfigProvider"; + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, ""); + // 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", invalidConfigProvider); + props.put("testKey", "${test:/foo/bar/testpath:testKey}"); + KafkaException e = assertThrows(KafkaException.class, () -> new TestIndirectConfigResolution(props, Collections.emptyMap())); + assertTrue(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); + } + + @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}"); + KafkaException e = assertThrows(KafkaException.class, () -> new TestIndirectConfigResolution(props, Collections.emptyMap())); + assertFalse(e.getMessage().contains(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY)); } @Test @@ -494,13 +558,15 @@ 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")); } @Test public void testConfigProviderConfigurationWithConfigParams() { - // Test Case: Valid Test Case With Multiple ConfigProviders as a separate variable + // should have no effect + System.setProperty(AbstractConfig.AUTOMATIC_CONFIG_PROVIDERS_PROPERTY, MockFileConfigProvider.class.getName()); + // 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()); @@ -510,7 +576,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/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/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 388803da946ab..bb0514613c542 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 @@ -472,7 +472,7 @@ public static List pluginLocations(Map props) { } public WorkerConfig(ConfigDef definition, Map props) { - super(definition, props); + super(definition, props, Utils.castToStringObjectMap(props), true); logInternalConverterRemovalWarnings(props); logPluginPathConfigProviderWarning(props); } @@ -598,4 +598,11 @@ public String toString() { + "if any part of a header rule contains a comma"; } } + + @Override + public Map originals() { + Map map = super.originals(); + map.remove(AbstractConfig.CONFIG_PROVIDERS_CONFIG); + return map; + } } 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 d981c3265e1a6..b5c0c565d2c17 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 @@ -1089,6 +1089,7 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { Map connConfig = Collections.singletonMap("metadata.max.age.ms", "10000"); Map expectedConfigs = new HashMap<>(workerProps); + expectedConfigs.remove(AbstractConfig.CONFIG_PROVIDERS_CONFIG); 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 6d9793bd84000..185ec09bb2388 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/server/DynamicBrokerConfig.scala b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala index e66b075b5c47d..04d55299681da 100755 --- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala +++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala @@ -189,7 +189,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) diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index 83a1f488f41a7..27666608fc698 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -1528,7 +1528,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)