- 
                Notifications
    You must be signed in to change notification settings 
- Fork 0
Clone kafka 18894 #31
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
563f347
              08dff59
              9e221e3
              1dd8f41
              204c8a9
              96fd5a6
              6ef4001
              44ecfdb
              e26a06f
              27d9a55
              b92a2bb
              5989a26
              File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | 
|---|---|---|
| @@ -0,0 +1,58 @@ | ||
| /* | ||
| * Licensed to the Apache Software Foundation (ASF) under one or more | ||
| * contributor license agreements. See the NOTICE file distributed with | ||
| * this work for additional information regarding copyright ownership. | ||
| * The ASF licenses this file to You under the Apache License, Version 2.0 | ||
| * (the "License"); you may not use this file except in compliance with | ||
| * the License. You may obtain a copy of the License at | ||
| * | ||
| * http://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software | ||
| * distributed under the License is distributed on an "AS IS" BASIS, | ||
| * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
| * See the License for the specific language governing permissions and | ||
| * limitations under the License. | ||
| */ | ||
| package org.apache.kafka.common.config.provider; | ||
|  | ||
| import org.apache.kafka.common.MetricName; | ||
| import org.apache.kafka.common.config.ConfigData; | ||
| import org.apache.kafka.common.metrics.Measurable; | ||
| import org.apache.kafka.common.metrics.Monitorable; | ||
| import org.apache.kafka.common.metrics.PluginMetrics; | ||
|  | ||
| import java.io.IOException; | ||
| import java.util.Map; | ||
| import java.util.Set; | ||
|  | ||
| public class MonitorableConfigProvider implements ConfigProvider, Monitorable { | ||
| public static final String NAME = "name"; | ||
| public static final String DESCRIPTION = "description"; | ||
| protected boolean configured = false; | ||
|  | ||
| @Override | ||
| public void withPluginMetrics(PluginMetrics metrics) { | ||
| MetricName metricName = metrics.metricName(NAME, DESCRIPTION, Map.of()); | ||
| metrics.addMetric(metricName, (Measurable) (config, now) -> 123); | ||
| } | ||
|  | ||
| @Override | ||
| public ConfigData get(String path) { | ||
| return null; | ||
| } | ||
|  | ||
| @Override | ||
| public ConfigData get(String path, Set<String> keys) { | ||
| return null; | ||
| } | ||
|  | ||
| @Override | ||
| public void close() throws IOException { | ||
| } | ||
|  | ||
| @Override | ||
| public void configure(Map<String, ?> configs) { | ||
| configured = true; | ||
| } | ||
| } | 
| Original file line number | Diff line number | Diff line change | 
|---|---|---|
|  | @@ -23,6 +23,7 @@ | |
| import org.apache.kafka.common.config.ConfigDef.Type; | ||
| import org.apache.kafka.common.config.ConfigTransformer; | ||
| import org.apache.kafka.common.config.provider.ConfigProvider; | ||
| import org.apache.kafka.common.internals.Plugin; | ||
| import org.apache.kafka.common.security.auth.SecurityProtocol; | ||
| import org.apache.kafka.common.utils.Utils; | ||
| import org.apache.kafka.connect.runtime.WorkerConfig; | ||
|  | @@ -269,18 +270,19 @@ List<String> configProviders() { | |
| Map<String, String> transform(Map<String, String> props) { | ||
| // transform worker config according to config.providers | ||
| List<String> providerNames = configProviders(); | ||
| Map<String, ConfigProvider> providers = new HashMap<>(); | ||
| Map<String, Plugin<ConfigProvider>> providerPlugins = new HashMap<>(); | ||
| for (String name : providerNames) { | ||
| ConfigProvider configProvider = plugins.newConfigProvider( | ||
| Plugin<ConfigProvider> configProviderPlugin = plugins.newConfigProvider( | ||
| this, | ||
| CONFIG_PROVIDERS_CONFIG + "." + name, | ||
| Plugins.ClassLoaderUsage.PLUGINS | ||
| name, | ||
| Plugins.ClassLoaderUsage.PLUGINS, | ||
| null | ||
| ); | ||
| providers.put(name, configProvider); | ||
| providerPlugins.put(name, configProviderPlugin); | ||
| } | ||
| ConfigTransformer transformer = new ConfigTransformer(providers); | ||
| ConfigTransformer transformer = new ConfigTransformer(providerPlugins); | ||
| Map<String, String> transformed = transformer.transform(props).data(); | ||
| providers.values().forEach(x -> Utils.closeQuietly(x, "config provider")); | ||
| providerPlugins.values().forEach(x -> Utils.closeQuietly(x, "config provider plugin")); | ||
| return transformed; | ||
| } | ||
| 
      Comment on lines
    
      270
     to 
      287
    
   There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🧩 Analysis chainEnsure plugins are always closed; skip null plugins. If transform throws, provider plugins aren’t closed, leaking resources/classloaders. Also, newConfigProvider may return null; avoid putting nulls and guard close. Apply: -        Map<String, Plugin<ConfigProvider>> providerPlugins = new HashMap<>();
-        for (String name : providerNames) {
-            Plugin<ConfigProvider> configProviderPlugin = plugins.newConfigProvider(
-                    this,
-                    name,
-                    Plugins.ClassLoaderUsage.PLUGINS,
-                    null
-            );
-            providerPlugins.put(name, configProviderPlugin);
-        }
-        ConfigTransformer transformer = new ConfigTransformer(providerPlugins);
-        Map<String, String> transformed = transformer.transform(props).data();
-        providerPlugins.values().forEach(x -> Utils.closeQuietly(x, "config provider plugin"));
-        return transformed;
+        Map<String, Plugin<ConfigProvider>> providerPlugins = new HashMap<>();
+        for (String name : providerNames) {
+            Plugin<ConfigProvider> plugin = plugins.newConfigProvider(
+                    this,
+                    name,
+                    Plugins.ClassLoaderUsage.PLUGINS,
+                    null
+            );
+            if (plugin != null) {
+                providerPlugins.put(name, plugin);
+            }
+        }
+        ConfigTransformer transformer = new ConfigTransformer(providerPlugins);
+        try {
+            return transformer.transform(props).data();
+        } finally {
+            providerPlugins.values().forEach(p -> Utils.closeQuietly(p, "config provider plugin"));
+        }If you decide not to filter nulls, please confirm Utils.closeQuietly tolerates nulls. 🏁 Script executed: #!/bin/bash
set -e
echo "=== newConfigProvider definitions and usage ==="
rg -nR "newConfigProvider" -n -C5 --type java
echo ""
echo "=== Utils.closeQuietly implementations ==="
rg -nR "closeQuietly" -n -C5 --type javaLength of output: 452 🏁 Script executed: #!/bin/bash
echo "=== Searching for newConfigProvider definitions ==="
rg -n "newConfigProvider" -g '*.java' -C3
echo "=== Searching for Utils.closeQuietly ==="
rg -n "closeQuietly" -g '*.java' -C3Length of output: 168790 🏁 Script executed: #!/bin/bash
sed -n '600,670p' connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.javaLength of output: 3899 Always null-check and close config provider plugins in a finally block 
 🤖 Prompt for AI Agents | ||
|  | ||
|  | ||
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|  | @@ -209,16 +209,17 @@ public Worker( | |||||||||||||||||||||||||||||||||||||||||||
|  | ||||||||||||||||||||||||||||||||||||||||||||
| private WorkerConfigTransformer initConfigTransformer() { | ||||||||||||||||||||||||||||||||||||||||||||
| final List<String> providerNames = config.getList(WorkerConfig.CONFIG_PROVIDERS_CONFIG); | ||||||||||||||||||||||||||||||||||||||||||||
| Map<String, ConfigProvider> providerMap = new HashMap<>(); | ||||||||||||||||||||||||||||||||||||||||||||
| Map<String, Plugin<ConfigProvider>> providerPluginMap = new HashMap<>(); | ||||||||||||||||||||||||||||||||||||||||||||
| for (String providerName : providerNames) { | ||||||||||||||||||||||||||||||||||||||||||||
| ConfigProvider configProvider = plugins.newConfigProvider( | ||||||||||||||||||||||||||||||||||||||||||||
| config, | ||||||||||||||||||||||||||||||||||||||||||||
| WorkerConfig.CONFIG_PROVIDERS_CONFIG + "." + providerName, | ||||||||||||||||||||||||||||||||||||||||||||
| ClassLoaderUsage.PLUGINS | ||||||||||||||||||||||||||||||||||||||||||||
| Plugin<ConfigProvider> configProviderPlugin = plugins.newConfigProvider( | ||||||||||||||||||||||||||||||||||||||||||||
| config, | ||||||||||||||||||||||||||||||||||||||||||||
| providerName, | ||||||||||||||||||||||||||||||||||||||||||||
| ClassLoaderUsage.PLUGINS, | ||||||||||||||||||||||||||||||||||||||||||||
| metrics.metrics() | ||||||||||||||||||||||||||||||||||||||||||||
| ); | ||||||||||||||||||||||||||||||||||||||||||||
| providerMap.put(providerName, configProvider); | ||||||||||||||||||||||||||||||||||||||||||||
| providerPluginMap.put(providerName, configProviderPlugin); | ||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||
| 
      Comment on lines
    
      +214
     to 
      221
    
   There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Guard against null providers to avoid NPEs in transformer 
 -            Plugin<ConfigProvider> configProviderPlugin = plugins.newConfigProvider(
+            Plugin<ConfigProvider> configProviderPlugin = plugins.newConfigProvider(
                 config,
                 providerName,
                 ClassLoaderUsage.PLUGINS,
                 metrics.metrics()
             );
-            providerPluginMap.put(providerName, configProviderPlugin);
+            if (configProviderPlugin != null) {
+                providerPluginMap.put(providerName, configProviderPlugin);
+            } else {
+                log.warn("Skipping config provider '{}' as no class configured under {}.{}.class",
+                        providerName, WorkerConfig.CONFIG_PROVIDERS_CONFIG, providerName);
+            }📝 Committable suggestion
 
        Suggested change
       
 🤖 Prompt for AI Agents | ||||||||||||||||||||||||||||||||||||||||||||
| return new WorkerConfigTransformer(this, providerMap); | ||||||||||||||||||||||||||||||||||||||||||||
| return new WorkerConfigTransformer(this, providerPluginMap); | ||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||
|  | ||||||||||||||||||||||||||||||||||||||||||||
| public WorkerConfigTransformer configTransformer() { | ||||||||||||||||||||||||||||||||||||||||||||
|  | ||||||||||||||||||||||||||||||||||||||||||||
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|  | @@ -19,6 +19,8 @@ | |||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.common.Configurable; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.common.config.AbstractConfig; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.common.config.provider.ConfigProvider; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.common.internals.Plugin; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.common.metrics.Metrics; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.common.utils.Utils; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.connect.components.Versioned; | ||||||||||||||||||||||||||||||||||||||||
| import org.apache.kafka.connect.connector.Connector; | ||||||||||||||||||||||||||||||||||||||||
|  | @@ -627,7 +629,8 @@ private <U> U newVersionedPlugin( | |||||||||||||||||||||||||||||||||||||||
| return plugin; | ||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||
|  | ||||||||||||||||||||||||||||||||||||||||
| public ConfigProvider newConfigProvider(AbstractConfig config, String providerPrefix, ClassLoaderUsage classLoaderUsage) { | ||||||||||||||||||||||||||||||||||||||||
| public Plugin<ConfigProvider> newConfigProvider(AbstractConfig config, String providerName, ClassLoaderUsage classLoaderUsage, Metrics metrics) { | ||||||||||||||||||||||||||||||||||||||||
| String providerPrefix = WorkerConfig.CONFIG_PROVIDERS_CONFIG + "." + providerName; | ||||||||||||||||||||||||||||||||||||||||
| String classPropertyName = providerPrefix + ".class"; | ||||||||||||||||||||||||||||||||||||||||
| Map<String, String> originalConfig = config.originalsStrings(); | ||||||||||||||||||||||||||||||||||||||||
| if (!originalConfig.containsKey(classPropertyName)) { | ||||||||||||||||||||||||||||||||||||||||
|  | @@ -643,7 +646,7 @@ public ConfigProvider newConfigProvider(AbstractConfig config, String providerPr | |||||||||||||||||||||||||||||||||||||||
| try (LoaderSwap loaderSwap = safeLoaderSwapper().apply(plugin.getClass().getClassLoader())) { | ||||||||||||||||||||||||||||||||||||||||
| plugin.configure(configProviderConfig); | ||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||
| return plugin; | ||||||||||||||||||||||||||||||||||||||||
| return Plugin.wrapInstance(plugin, metrics, WorkerConfig.CONFIG_PROVIDERS_CONFIG, Map.of("provider", providerName)); | ||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||
| 
      Comment on lines
    
      646
     to 
      650
    
   There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🛠️ Refactor suggestion | 🟠 Major Close provider if configure fails to prevent leaks If  -        try (LoaderSwap loaderSwap = safeLoaderSwapper().apply(plugin.getClass().getClassLoader())) {
-            plugin.configure(configProviderConfig);
-        }
-        return Plugin.wrapInstance(plugin, metrics, WorkerConfig.CONFIG_PROVIDERS_CONFIG, Map.of("provider", providerName));
+        try (LoaderSwap loaderSwap = safeLoaderSwapper().apply(plugin.getClass().getClassLoader())) {
+            plugin.configure(configProviderConfig);
+        } catch (Throwable t) {
+            // Ensure we don't leak a partially configured provider
+            Utils.maybeCloseQuietly(plugin, "config provider");
+            throw t;
+        }
+        return Plugin.wrapInstance(
+            plugin,
+            metrics,
+            WorkerConfig.CONFIG_PROVIDERS_CONFIG,
+            Map.of("provider", providerName)
+        );📝 Committable suggestion
 
        Suggested change
       
 🤖 Prompt for AI Agents | ||||||||||||||||||||||||||||||||||||||||
|  | ||||||||||||||||||||||||||||||||||||||||
| /** | ||||||||||||||||||||||||||||||||||||||||
|  | ||||||||||||||||||||||||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🛠️ Refactor suggestion | 🟠 Major
Always close provider plugins even if transform throws
Wrap
transform(...)in try/finally so provider plugins are closed on all paths. Prevents leaking provider resources and plugin metrics.📝 Committable suggestion
🤖 Prompt for AI Agents