Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.kafka.connect.errors.ConnectException;
import org.apache.kafka.connect.errors.NotFoundException;
import org.apache.kafka.connect.runtime.isolation.LoaderSwap;
import org.apache.kafka.connect.runtime.isolation.PluginDesc;
import org.apache.kafka.connect.runtime.isolation.Plugins;
import org.apache.kafka.connect.runtime.rest.entities.ActiveTopicsInfo;
import org.apache.kafka.connect.runtime.rest.entities.ConfigInfo;
Expand Down Expand Up @@ -661,10 +662,24 @@ ConfigInfos validateConnectorConfig(
}
}
String connType = connectorProps.get(ConnectorConfig.CONNECTOR_CLASS_CONFIG);
if (connType == null)
throw new BadRequestException("Connector config " + connectorProps + " contains no connector type");
if (connType == null) {
return createConnectorClassError("Config contains no connector type");
}

Connector connector;
try {
connector = getConnector(connType);
} catch (ConnectException e) {
return createConnectorClassError(e.getMessage());
Comment thread
urbandan marked this conversation as resolved.
Outdated
}

return validateConnectorConfig(connector, connType, connectorProps, reportStage, doLog);
}

Connector connector = getConnector(connType);
private ConfigInfos validateConnectorConfig(Connector connector, String connType,
Map<String, String> connectorProps,
Function<String, TemporaryStage> reportStage, boolean doLog) {
String stageDescription;
ClassLoader connectorLoader = plugins().connectorLoader(connType);
try (LoaderSwap loaderSwap = plugins().withClassLoader(connectorLoader)) {
org.apache.kafka.connect.health.ConnectorType connectorType;
Expand Down Expand Up @@ -950,6 +965,19 @@ private static ConfigValueInfo convertConfigValue(ConfigValue configValue, Type
return new ConfigValueInfo(configValue.name(), value, recommendedValues, configValue.errorMessages(), configValue.visible());
}

private ConfigInfos createConnectorClassError(String errorMessage) {
ConfigKey connectorClassConfigKey =
ConnectorConfig.configDef().configKeys().get(ConnectorConfig.CONNECTOR_CLASS_CONFIG);
ConfigValueInfo valueInfo = new ConfigValueInfo(ConnectorConfig.CONNECTOR_CLASS_CONFIG, null,
connectorClassNames(), Collections.singletonList(errorMessage), true);
ConfigInfo info = new ConfigInfo(convertConfigKey(connectorClassConfigKey), valueInfo);
return new ConfigInfos("", 1, Collections.emptyList(), Collections.singletonList(info));
}

private List<String> connectorClassNames() {
return plugins().connectors().stream().map(PluginDesc::className).collect(Collectors.toList());
}

protected Connector getConnector(String connType) {
return tempConnectors.computeIfAbsent(connType, k -> plugins().newConnector(k));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -317,16 +317,13 @@ public Class<? extends Connector> connectorClass(String connectorClassOrAlias) {
throw new ConnectException(
"Failed to find any class that implements Connector and which name matches "
+ connectorClassOrAlias
+ ", available connectors are: "
+ connectors.stream().map(PluginDesc::toString).collect(Collectors.joining(", "))
);
}
if (matches.size() > 1) {
throw new ConnectException(
"More than one connector matches alias "
+ connectorClassOrAlias
+ ". Please use full package and class name instead. Classes found: "
+ connectors.stream().map(PluginDesc::toString).collect(Collectors.joining(", "))
+ ". Please use full package and class name instead."
);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.junit.jupiter.MockitoExtension;
import org.mockito.junit.jupiter.MockitoSettings;
import org.mockito.quality.Strictness;
Expand Down Expand Up @@ -457,9 +458,44 @@ public void testBuildRestartPlanForNoRestart() {
public void testConfigValidationEmptyConfig() {
AbstractHerder herder = createConfigValidationHerder(SampleSourceConnector.class, noneConnectorClientConfigOverridePolicy, 0);

assertThrows(BadRequestException.class, () -> herder.validateConnectorConfig(Collections.emptyMap(), s -> null, false));
verify(transformer).transform(Collections.emptyMap());
assertEquals(worker.getPlugins(), plugins);
@SuppressWarnings("unchecked")
PluginDesc<Connector> mockPluginDesc = mock(PluginDesc.class);
when(mockPluginDesc.className()).thenReturn(SampleSourceConnector.class.getName());
when(plugins.connectors()).thenReturn(Collections.singleton(mockPluginDesc));

ConfigInfos result = herder.validateConnectorConfig(Collections.emptyMap(), s -> null, false);

assertEquals(1, result.errorCount());
assertEquals(1, result.values().size());
ConfigInfo configInfo = result.values().get(0);
assertEquals(ConnectorConfig.CONNECTOR_CLASS_CONFIG, configInfo.configKey().name());
assertEquals(Collections.singletonList(SampleSourceConnector.class.getName()),
configInfo.configValue().recommendedValues());
}

@Test
public void testConfigValidationInvalidClassConfig() {
AbstractHerder herder = createConfigValidationHerder(SampleSourceConnector.class, noneConnectorClientConfigOverridePolicy, 0);
@SuppressWarnings("unchecked")
PluginDesc<Connector> mockPluginDesc = mock(PluginDesc.class);
when(mockPluginDesc.className()).thenReturn(SampleSourceConnector.class.getName());
Mockito.reset(plugins);
when(plugins.newConnector("InvalidSourceConnector"))
.thenThrow(new ConnectException("Test: Invalid class"));
when(plugins.connectors()).thenReturn(Collections.singleton(mockPluginDesc));

Map<String, String> config = new HashMap<>();
config.put(ConnectorConfig.CONNECTOR_CLASS_CONFIG, "InvalidSourceConnector");

ConfigInfos result = herder.validateConnectorConfig(config, s -> null, false);

assertEquals(1, result.errorCount());
assertEquals(1, result.values().size());
ConfigInfo configInfo = result.values().get(0);
assertEquals(ConnectorConfig.CONNECTOR_CLASS_CONFIG, configInfo.configKey().name());
assertEquals(Collections.singletonList(SampleSourceConnector.class.getName()),
configInfo.configValue().recommendedValues());
verify(transformer).transform(Collections.singletonMap(ConnectorConfig.CONNECTOR_CLASS_CONFIG, "InvalidSourceConnector"));
}

@Test
Expand Down