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 555b634ec251e..c591e72081535 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 @@ -55,6 +55,8 @@ public class AbstractConfig { private static final String CONFIG_PROVIDERS_CONFIG = "config.providers"; + private static final String CONFIG_PROVIDERS_PARAM = ".param."; + /** * Construct a configuration with a ConfigDef and the configuration properties, which can include properties * for zero or more {@link ConfigProvider} that will be used to resolve variables in configuration property @@ -492,6 +494,16 @@ private Map configProviderProperties(String configProviderPrefix return result; } + /** + * Instantiates and configures the ConfigProviders. The config providers configs are defined as follows: + * config.providers : A comma-separated list of names for providers. + * config.providers.{name}.class : The Java class name for a provider. + * config.providers.{name}.param.{param-name} : A parameter to be passed to the above Java class on initialization. + * returns a map of config provider name and its instance. + * @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. + */ private Map instantiateConfigProviders(Map indirectConfigs, Map providerConfigProperties) { final String configProviders = indirectConfigs.get(CONFIG_PROVIDERS_CONFIG); @@ -511,7 +523,7 @@ private Map instantiateConfigProviders(Map configProviderInstances = new HashMap<>(); for (Map.Entry entry : providerMap.entrySet()) { try { - String prefix = CONFIG_PROVIDERS_CONFIG + "." + entry.getKey() + "."; + String prefix = CONFIG_PROVIDERS_CONFIG + "." + entry.getKey() + CONFIG_PROVIDERS_PARAM; Map configProperties = configProviderProperties(prefix, providerConfigProperties); ConfigProvider provider = Utils.newInstance(entry.getValue(), ConfigProvider.class); provider.configure(configProperties); 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 a007fd3707eb4..0f2d27084f61a 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 @@ -24,6 +24,8 @@ import org.apache.kafka.common.metrics.JmxReporter; import org.apache.kafka.common.metrics.MetricsReporter; 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.junit.Test; import java.util.Arrays; @@ -337,7 +339,7 @@ public void testOriginalsWithConfigProvidersProps() { // Test Case: Valid Test Case for ConfigProviders as part of config.properties props.put("config.providers", "file"); - props.put("config.providers.file.class", "org.apache.kafka.common.config.provider.MockFileConfigProvider"); + props.put("config.providers.file.class", MockFileConfigProvider.class.getName()); 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}"); @@ -354,7 +356,7 @@ public void testConfigProvidersPropsAsParam() { // Test Case: Valid Test Case for ConfigProviders as a separate variable Properties providers = new Properties(); providers.put("config.providers", "file"); - providers.put("config.providers.file.class", "org.apache.kafka.common.config.provider.MockFileConfigProvider"); + providers.put("config.providers.file.class", MockFileConfigProvider.class.getName()); Properties props = new Properties(); props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); props.put("sasl.kerberos.password", "${file:/usr/kerberos:password}"); @@ -368,8 +370,8 @@ public void testAutoConfigResolutionWithMultipleConfigProviders() { // Test Case: Valid Test Case With Multiple ConfigProviders as a separate variable Properties providers = new Properties(); providers.put("config.providers", "file,vault"); - providers.put("config.providers.file.class", "org.apache.kafka.common.config.provider.MockFileConfigProvider"); - providers.put("config.providers.vault.class", "org.apache.kafka.common.config.provider.MockVaultConfigProvider"); + providers.put("config.providers.file.class", MockFileConfigProvider.class.getName()); + providers.put("config.providers.vault.class", MockVaultConfigProvider.class.getName()); Properties props = new Properties(); props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); props.put("sasl.kerberos.password", "${file:/usr/kerberos:password}"); @@ -412,7 +414,7 @@ public void testAutoConfigResolutionWithMissingConfigKey() { // Test Case: Config Provider fails to resolve the config (key not present) Properties props = new Properties(); props.put("config.providers", "test"); - props.put("config.providers.test.class", "org.apache.kafka.common.config.provider.MockFileConfigProvider"); + props.put("config.providers.test.class", MockFileConfigProvider.class.getName()); props.put("random", "${test:/foo/bar/testpath:random}"); TestIndirectConfigResolution config = new TestIndirectConfigResolution(props); assertEquals(config.originals().get("random"), "${test:/foo/bar/testpath:random}"); @@ -423,17 +425,33 @@ public void testAutoConfigResolutionWithDuplicateConfigProvider() { // Test Case: If ConfigProvider is provided in both originals and provider. Only the ones in provider should be used. Properties providers = new Properties(); providers.put("config.providers", "test"); - providers.put("config.providers.test.class", "org.apache.kafka.common.config.provider.MockFileConfigProvider"); + providers.put("config.providers.test.class", MockVaultConfigProvider.class.getName()); Properties props = new Properties(); props.put("sasl.kerberos.key", "${file:/usr/kerberos:key}"); props.put("config.providers", "file"); - props.put("config.providers.file.class", "org.apache.kafka.common.config.provider.MockFileConfigProvider"); + props.put("config.providers.file.class", MockVaultConfigProvider.class.getName()); TestIndirectConfigResolution config = new TestIndirectConfigResolution(props, convertPropertiesToMap(providers)); assertEquals(config.originals().get("sasl.kerberos.key"), "${file:/usr/kerberos:key}"); } + @Test + public void testConfigProviderConfigurationWithConfigParams() { + // Test Case: Valid Test Case With Multiple ConfigProviders as a separate variable + Properties providers = new Properties(); + providers.put("config.providers", "vault"); + providers.put("config.providers.vault.class", MockVaultConfigProvider.class.getName()); + providers.put("config.providers.vault.param.key", "randomKey"); + providers.put("config.providers.vault.param.location", "/usr/vault"); + Properties props = new Properties(); + props.put("sasl.truststore.location", "${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)); + assertEquals(config.originals().get("sasl.truststore.location"), "/usr/vault"); + } + private static class TestIndirectConfigResolution extends AbstractConfig { private static final ConfigDef CONFIG; diff --git a/clients/src/test/java/org/apache/kafka/common/config/provider/MockVaultConfigProvider.java b/clients/src/test/java/org/apache/kafka/common/config/provider/MockVaultConfigProvider.java index dbfe15813941c..c741798a72972 100644 --- a/clients/src/test/java/org/apache/kafka/common/config/provider/MockVaultConfigProvider.java +++ b/clients/src/test/java/org/apache/kafka/common/config/provider/MockVaultConfigProvider.java @@ -19,11 +19,27 @@ import java.io.IOException; import java.io.Reader; import java.io.StringReader; +import java.util.Map; public class MockVaultConfigProvider extends FileConfigProvider { + Map vaultConfigs; + private boolean configured = false; + private static final String LOCATION = "location"; + @Override protected Reader reader(String path) throws IOException { - return new StringReader("truststoreKey=testTruststoreKey\ntruststorePassword=randomtruststorePassword"); + String vaultLocation = (String) vaultConfigs.get(LOCATION); + return new StringReader("truststoreKey=testTruststoreKey\ntruststorePassword=randomtruststorePassword\n" + "truststoreLocation=" + vaultLocation + "\n"); + } + + @Override + public void configure(Map configs) { + this.vaultConfigs = configs; + configured = true; + } + + public boolean configured() { + return configured; } }