diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java index 174b032ea2488..7fafd39ee7f44 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java @@ -33,7 +33,11 @@ import org.apache.kafka.connect.connector.policy.ConnectorClientConfigRequest; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.errors.NotFoundException; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; import org.apache.kafka.connect.runtime.isolation.LoaderSwap; +import org.apache.kafka.connect.runtime.isolation.PluginType; 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; @@ -137,7 +141,7 @@ public abstract class AbstractHerder implements Herder, TaskStatus.Listener, Con private final Time time; protected final Loggers loggers; - private final ConcurrentMap tempConnectors = new ConcurrentHashMap<>(); + private final ConcurrentMap> tempConnectors = new ConcurrentHashMap<>(); public AbstractHerder(Worker worker, String workerId, @@ -389,13 +393,13 @@ public ConnectorStateInfo.TaskState taskStatus(ConnectorTaskId id) { status.workerId(), status.trace()); } - protected Map validateSinkConnectorConfig(SinkConnector connector, ConfigDef configDef, Map config) { + protected Map validateSinkConnectorConfig(IsolatedSinkConnector connector, ConfigDef configDef, Map config) { Map result = configDef.validateAll(config); SinkConnectorConfig.validate(config, result); return result; } - protected Map validateSourceConnectorConfig(SourceConnector connector, ConfigDef configDef, Map config) { + protected Map validateSourceConnectorConfig(IsolatedSourceConnector connector, ConfigDef configDef, Map config) { return configDef.validateAll(config); } @@ -647,7 +651,7 @@ ConfigInfos validateConnectorConfig( Map connectorProps, Function reportStage, boolean doLog - ) { + ) throws Exception { String stageDescription; if (worker.configTransformer() != null) { stageDescription = "resolving transformed configuration properties for the connector"; @@ -659,25 +663,26 @@ ConfigInfos validateConnectorConfig( if (connType == null) throw new BadRequestException("Connector config " + connectorProps + " contains no connector type"); - Connector connector = getConnector(connType); + IsolatedConnector connector = getConnector(connType); ClassLoader connectorLoader = plugins().connectorLoader(connType); try (LoaderSwap loaderSwap = plugins().withClassLoader(connectorLoader)) { org.apache.kafka.connect.health.ConnectorType connectorType; ConfigDef enrichedConfigDef; Map validatedConnectorConfig; - if (connector instanceof SourceConnector) { + PluginType type = connector.type(); + if (type == PluginType.SOURCE) { connectorType = org.apache.kafka.connect.health.ConnectorType.SOURCE; enrichedConfigDef = ConnectorConfig.enrich(plugins(), SourceConnectorConfig.configDef(), connectorProps, false); stageDescription = "validating source connector-specific properties for the connector"; try (TemporaryStage stage = reportStage.apply(stageDescription)) { - validatedConnectorConfig = validateSourceConnectorConfig((SourceConnector) connector, enrichedConfigDef, connectorProps); + validatedConnectorConfig = validateSourceConnectorConfig((IsolatedSourceConnector) connector, enrichedConfigDef, connectorProps); } } else { connectorType = org.apache.kafka.connect.health.ConnectorType.SINK; enrichedConfigDef = ConnectorConfig.enrich(plugins(), SinkConnectorConfig.configDef(), connectorProps, false); stageDescription = "validating sink connector-specific properties for the connector"; try (TemporaryStage stage = reportStage.apply(stageDescription)) { - validatedConnectorConfig = validateSinkConnectorConfig((SinkConnector) connector, enrichedConfigDef, connectorProps); + validatedConnectorConfig = validateSinkConnectorConfig((IsolatedSinkConnector) connector, enrichedConfigDef, connectorProps); } } @@ -703,7 +708,7 @@ ConfigInfos validateConnectorConfig( throw new BadRequestException( String.format( "%s.config() must return a ConfigDef that is not null.", - connector.getClass().getName() + connector.pluginClass().getName() ) ); } @@ -717,7 +722,7 @@ ConfigInfos validateConnectorConfig( throw new BadRequestException( String.format( "%s.validate() must return a Config that is not null.", - connector.getClass().getName() + connector.pluginClass().getName() ) ); } @@ -757,7 +762,7 @@ ConfigInfos validateConnectorConfig( ConnectorConfig.CONNECTOR_CLIENT_PRODUCER_OVERRIDES_PREFIX, connectorConfig, ProducerConfig.configDef(), - connector.getClass(), + connector.pluginClass(), connectorType, ConnectorClientConfigRequest.ClientType.PRODUCER, connectorClientConfigOverridePolicy); @@ -771,7 +776,7 @@ ConfigInfos validateConnectorConfig( ConnectorConfig.CONNECTOR_CLIENT_ADMIN_OVERRIDES_PREFIX, connectorConfig, AdminClientConfig.configDef(), - connector.getClass(), + connector.pluginClass(), connectorType, ConnectorClientConfigRequest.ClientType.ADMIN, connectorClientConfigOverridePolicy); @@ -785,7 +790,7 @@ ConfigInfos validateConnectorConfig( ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX, connectorConfig, ConsumerConfig.configDef(), - connector.getClass(), + connector.pluginClass(), connectorType, ConnectorClientConfigRequest.ClientType.CONSUMER, connectorClientConfigOverridePolicy); @@ -945,7 +950,7 @@ private static ConfigValueInfo convertConfigValue(ConfigValue configValue, Type return new ConfigValueInfo(configValue.name(), value, recommendedValues, configValue.errorMessages(), configValue.visible()); } - protected Connector getConnector(String connType) { + protected IsolatedConnector getConnector(String connType) { return tempConnectors.computeIfAbsent(connType, k -> plugins().newConnector(k)); } @@ -964,7 +969,7 @@ public ConnectorType connectorType(Map connConfig) { return ConnectorType.UNKNOWN; } try { - return ConnectorType.from(getConnector(connClass).getClass()); + return ConnectorType.from(getConnector(connClass).pluginClass()); } catch (ConnectException e) { log.warn("Unable to retrieve connector type", e); return ConnectorType.UNKNOWN; 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..907693163540a 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 @@ -33,8 +33,8 @@ import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerConfig; -import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.MetricNameTemplate; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.config.ConfigDef; @@ -58,6 +58,13 @@ import org.apache.kafka.connect.json.JsonConverter; import org.apache.kafka.connect.json.JsonConverterConfig; import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; +import org.apache.kafka.connect.runtime.isolation.LoaderSwap; +import org.apache.kafka.connect.runtime.isolation.PluginType; +import org.apache.kafka.connect.source.SourceConnector; +import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.distributed.DistributedConfig; import org.apache.kafka.connect.runtime.errors.DeadLetterQueueReporter; import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; @@ -65,7 +72,6 @@ import org.apache.kafka.connect.runtime.errors.LogReporter; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.WorkerErrantRecordReporter; -import org.apache.kafka.connect.runtime.isolation.LoaderSwap; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.isolation.Plugins.ClassLoaderUsage; import org.apache.kafka.connect.runtime.rest.RestServer; @@ -75,11 +81,9 @@ import org.apache.kafka.connect.sink.SinkConnector; import org.apache.kafka.connect.sink.SinkRecord; import org.apache.kafka.connect.sink.SinkTask; -import org.apache.kafka.connect.source.SourceConnector; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.source.SourceTask; import org.apache.kafka.connect.storage.CloseableOffsetStorageReader; -import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.storage.ConnectorOffsetBackingStore; import org.apache.kafka.connect.storage.Converter; import org.apache.kafka.connect.storage.HeaderConverter; @@ -310,11 +314,11 @@ public void startConnector( ClassLoader connectorLoader = plugins.connectorLoader(connClass); try (LoaderSwap loaderSwap = plugins.withClassLoader(connectorLoader)) { log.info("Creating connector {} of type {}", connName, connClass); - final Connector connector = plugins.newConnector(connClass); + final IsolatedConnector connector = plugins.newConnector(connClass); final ConnectorConfig connConfig; final CloseableOffsetStorageReader offsetReader; final ConnectorOffsetBackingStore offsetStore; - if (ConnectUtils.isSinkConnector(connector)) { + if (connector.type() == PluginType.SINK) { connConfig = new SinkConnectorConfig(plugins, connProps); offsetReader = null; offsetStore = null; @@ -331,7 +335,7 @@ public void startConnector( } workerConnector = new WorkerConnector( connName, connector, connConfig, ctx, metrics, connectorStatusListener, offsetReader, offsetStore, connectorLoader); - log.info("Instantiated connector {} with version {} of type {}", connName, connector.version(), connector.getClass()); + log.info("Instantiated connector {} with version {} of type {}", connName, connector.version(), connector.pluginClass()); workerConnector.transitionTo(initialState, onConnectorStateChange); } catch (Throwable t) { log.error("Failed to start connector {}", connName, t); @@ -379,7 +383,7 @@ public boolean isSinkConnector(String connName) { * @param connName the connector name. * @return a list of updated tasks properties. */ - public List> connectorTaskConfigs(String connName, ConnectorConfig connConfig) { + public List> connectorTaskConfigs(String connName, ConnectorConfig connConfig) throws Exception { List> result = new ArrayList<>(); try (LoggingContext loggingContext = LoggingContext.forConnector(connName)) { log.trace("Reconfiguring connector tasks for {}", connName); @@ -391,7 +395,7 @@ public List> connectorTaskConfigs(String connName, Connector int maxTasks = connConfig.tasksMax(); Map connOriginals = connConfig.originalsStrings(); - Connector connector = workerConnector.connector(); + IsolatedConnector connector = workerConnector.connector(); try (LoaderSwap loaderSwap = plugins.withClassLoader(workerConnector.loader())) { String taskClassName = connector.taskClass().getName(); List> taskConfigs = connector.taskConfigs(maxTasks); @@ -1195,8 +1199,8 @@ public void connectorOffsets(String connName, Map connectorConfi ClassLoader connectorLoader = plugins.connectorLoader(connectorClassOrAlias); try (LoaderSwap loaderSwap = plugins.withClassLoader(connectorLoader)) { - Connector connector = plugins.newConnector(connectorClassOrAlias); - if (ConnectUtils.isSinkConnector(connector)) { + IsolatedConnector connector = plugins.newConnector(connectorClassOrAlias); + if (connector.type() == PluginType.SINK) { log.debug("Fetching offsets for sink connector: {}", connName); sinkConnectorOffsets(connName, connector, connectorConfig, cb); } else { @@ -1216,20 +1220,20 @@ public void connectorOffsets(String connName, Map connectorConfi * @param connectorConfig the sink connector's configurations * @param cb callback to invoke upon completion of the request */ - void sinkConnectorOffsets(String connName, Connector connector, Map connectorConfig, + void sinkConnectorOffsets(String connName, IsolatedConnector connector, Map connectorConfig, Callback cb) { Map adminConfig = adminConfigs( connName, "connector-worker-adminclient-" + connName, config, new SinkConnectorConfig(plugins, connectorConfig), - connector.getClass(), + connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SINK); String groupId = (String) baseConsumerConfigs( connName, "connector-consumer-", config, new SinkConnectorConfig(plugins, connectorConfig), - connector.getClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SINK).get(ConsumerConfig.GROUP_ID_CONFIG); + connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SINK).get(ConsumerConfig.GROUP_ID_CONFIG); Admin admin = adminFactory.apply(adminConfig); try { ListConsumerGroupOffsetsOptions listOffsetsOptions = new ListConsumerGroupOffsetsOptions() @@ -1260,7 +1264,7 @@ void sinkConnectorOffsets(String connName, Connector connector, Map connectorConfig, + private void sourceConnectorOffsets(String connName, IsolatedConnector connector, Map connectorConfig, Callback cb) { SourceConnectorConfig sourceConfig = new SourceConnectorConfig(plugins, connectorConfig, config.topicCreationEnable()); ConnectorOffsetBackingStore offsetStore = config.exactlyOnceSourceEnabled() @@ -1306,16 +1310,16 @@ public void modifyConnectorOffsets(String connName, Map connecto Map, Map> offsets, Callback cb) { String connectorClassOrAlias = connectorConfig.get(ConnectorConfig.CONNECTOR_CLASS_CONFIG); ClassLoader connectorLoader = plugins.connectorLoader(connectorClassOrAlias); - Connector connector; + IsolatedConnector connector; try (LoaderSwap loaderSwap = plugins.withClassLoader(connectorLoader)) { connector = plugins.newConnector(connectorClassOrAlias); - if (ConnectUtils.isSinkConnector(connector)) { + if (connector.type() == PluginType.SINK) { log.debug("Modifying offsets for sink connector: {}", connName); - modifySinkConnectorOffsets(connName, connector, connectorConfig, offsets, connectorLoader, cb); + modifySinkConnectorOffsets(connName, (IsolatedSinkConnector) connector, connectorConfig, offsets, connectorLoader, cb); } else { log.debug("Modifying offsets for source connector: {}", connName); - modifySourceConnectorOffsets(connName, connector, connectorConfig, offsets, connectorLoader, cb); + modifySourceConnectorOffsets(connName, (IsolatedSourceConnector) connector, connectorConfig, offsets, connectorLoader, cb); } } } @@ -1333,14 +1337,14 @@ public void modifyConnectorOffsets(String connName, Map connecto * @param connectorLoader the connector plugin's classloader to be used as the thread context classloader * @param cb callback to invoke upon completion */ - void modifySinkConnectorOffsets(String connName, Connector connector, Map connectorConfig, + void modifySinkConnectorOffsets(String connName, IsolatedSinkConnector connector, Map connectorConfig, Map, Map> offsets, ClassLoader connectorLoader, Callback cb) { executor.submit(plugins.withClassLoader(connectorLoader, () -> { try { Timer timer = time.timer(Duration.ofMillis(RestServer.DEFAULT_REST_REQUEST_TIMEOUT_MS)); boolean isReset = offsets == null; SinkConnectorConfig sinkConnectorConfig = new SinkConnectorConfig(plugins, connectorConfig); - Class sinkConnectorClass = connector.getClass(); + Class sinkConnectorClass = connector.pluginClass(); Map adminConfig = adminConfigs( connName, "connector-worker-adminclient-" + connName, @@ -1384,7 +1388,7 @@ void modifySinkConnectorOffsets(String connName, Connector connector, Map connectorConfig, + private void modifySourceConnectorOffsets(String connName, IsolatedSourceConnector connector, Map connectorConfig, Map, Map> offsets, ClassLoader connectorLoader, Callback cb) { SourceConnectorConfig sourceConfig = new SourceConnectorConfig(plugins, connectorConfig, config.topicCreationEnable()); Map producerProps = config.exactlyOnceSourceEnabled() ? exactlyOnceSourceTaskProducerConfigs(new ConnectorTaskId(connName, 0), config, sourceConfig, - connector.getClass(), connectorClientConfigOverridePolicy, kafkaClusterId) + connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId) : baseProducerConfigs(connName, "connector-offset-producer-" + connName, config, sourceConfig, - connector.getClass(), connectorClientConfigOverridePolicy, kafkaClusterId); + connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId); KafkaProducer producer = new KafkaProducer<>(producerProps); ConnectorOffsetBackingStore offsetStore = config.exactlyOnceSourceEnabled() @@ -1562,7 +1566,7 @@ private void modifySourceConnectorOffsets(String connName, Connector connector, } // Visible for testing - void modifySourceConnectorOffsets(String connName, Connector connector, Map connectorConfig, + void modifySourceConnectorOffsets(String connName, IsolatedSourceConnector connector, Map connectorConfig, Map, Map> offsets, ConnectorOffsetBackingStore offsetStore, KafkaProducer producer, OffsetStorageWriter offsetWriter, ClassLoader connectorLoader, Callback cb) { @@ -1592,7 +1596,7 @@ void modifySourceConnectorOffsets(String connName, Connector connector, Map doBuild( ConnectorOffsetBackingStore offsetStoreForRegularSourceConnector( SourceConnectorConfig sourceConfig, String connName, - Connector connector, + IsolatedConnector connector, Producer producer ) { String connectorSpecificOffsetsTopic = sourceConfig.offsetsTopic(); - Map producerProps = baseProducerConfigs(connName, "connector-producer-" + connName, config, sourceConfig, connector.getClass(), + Map producerProps = baseProducerConfigs(connName, "connector-producer-" + connName, config, sourceConfig, connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId); // We use a connector-specific store (i.e., a dedicated KafkaOffsetBackingStore for this connector) @@ -2017,12 +2021,12 @@ ConnectorOffsetBackingStore offsetStoreForRegularSourceConnector( if (usesConnectorSpecificStore) { Map consumerProps = regularSourceOffsetsConsumerConfigs( - connName, "connector-consumer-" + connName, config, sourceConfig, connector.getClass(), + connName, "connector-consumer-" + connName, config, sourceConfig, connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId); KafkaConsumer consumer = new KafkaConsumer<>(consumerProps); Map adminOverrides = adminConfigs(connName, "connector-adminclient-" + connName, config, - sourceConfig, connector.getClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SOURCE); + sourceConfig, connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SOURCE); TopicAdmin admin = new TopicAdmin(adminOverrides); @@ -2080,21 +2084,21 @@ ConnectorOffsetBackingStore offsetStoreForRegularSourceConnector( ConnectorOffsetBackingStore offsetStoreForExactlyOnceSourceConnector( SourceConnectorConfig sourceConfig, String connName, - Connector connector, + IsolatedConnector connector, Producer producer ) { String connectorSpecificOffsetsTopic = Optional.ofNullable(sourceConfig.offsetsTopic()).orElse(config.offsetsTopic()); - Map producerProps = baseProducerConfigs(connName, "connector-producer-" + connName, config, sourceConfig, connector.getClass(), + Map producerProps = baseProducerConfigs(connName, "connector-producer-" + connName, config, sourceConfig, connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId); Map consumerProps = exactlyOnceSourceOffsetsConsumerConfigs( - connName, "connector-consumer-" + connName, config, sourceConfig, connector.getClass(), + connName, "connector-consumer-" + connName, config, sourceConfig, connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId); KafkaConsumer consumer = new KafkaConsumer<>(consumerProps); Map adminOverrides = adminConfigs(connName, "connector-adminclient-" + connName, config, - sourceConfig, connector.getClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SOURCE); + sourceConfig, connector.pluginClass(), connectorClientConfigOverridePolicy, kafkaClusterId, ConnectorType.SOURCE); TopicAdmin admin = new TopicAdmin(adminOverrides); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConnector.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConnector.java index 84cec3640b9cd..b1ed15398fbcf 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConnector.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConnector.java @@ -17,17 +17,18 @@ package org.apache.kafka.connect.runtime; import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.connector.Connector; import org.apache.kafka.connect.connector.ConnectorContext; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.PluginDesc; +import org.apache.kafka.connect.runtime.isolation.PluginType; import org.apache.kafka.connect.sink.SinkConnectorContext; import org.apache.kafka.connect.source.SourceConnectorContext; import org.apache.kafka.connect.storage.CloseableOffsetStorageReader; import org.apache.kafka.connect.storage.ConnectorOffsetBackingStore; import org.apache.kafka.connect.storage.OffsetStorageReader; import org.apache.kafka.connect.util.Callback; -import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.LoggingContext; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -68,7 +69,7 @@ private enum State { private final ConnectorStatus.Listener statusListener; private final ClassLoader loader; private final CloseableConnectorContext ctx; - private final Connector connector; + private final IsolatedConnector connector; private final ConnectorMetricsGroup metrics; private final AtomicReference pendingTargetStateChange; private final AtomicReference> pendingStateChangeCallback; @@ -82,7 +83,7 @@ private enum State { private final ConnectorOffsetBackingStore offsetStore; public WorkerConnector(String connName, - Connector connector, + IsolatedConnector connector, ConnectorConfig connectorConfig, CloseableConnectorContext ctx, ConnectMetrics metrics, @@ -203,7 +204,7 @@ void initialize() { } } - private boolean doStart() { + private boolean doStart() throws Throwable { try { switch (state) { case STARTED: @@ -235,12 +236,12 @@ private synchronized void onFailure(Throwable t) { this.state = State.FAILED; } - private void resume() { + private void resume() throws Throwable { if (doStart()) statusListener.onResume(connName); } - private void start() { + private void start() throws Throwable { if (doStart()) statusListener.onStartup(connName); } @@ -392,7 +393,7 @@ void doTransitionTo(TargetState targetState, Callback stateChangeCa } } - private void doTransitionTo(TargetState targetState) { + private void doTransitionTo(TargetState targetState) throws Throwable { log.debug("{} Transition connector to {}", this, targetState); if (targetState == TargetState.PAUSED) { suspend(true); @@ -409,11 +410,11 @@ private void doTransitionTo(TargetState targetState) { } public final boolean isSinkConnector() { - return ConnectUtils.isSinkConnector(connector); + return connector.type() == PluginType.SINK; } public final boolean isSourceConnector() { - return ConnectUtils.isSourceConnector(connector); + return connector.type() == PluginType.SOURCE; } protected final String connectorType() { @@ -424,7 +425,7 @@ protected final String connectorType() { return "unknown"; } - public Connector connector() { + public IsolatedConnector connector() { return connector; } @@ -462,8 +463,14 @@ public ConnectorMetricsGroup(ConnectMetrics connectMetrics, AbstractStatus.State metricGroup.close(); metricGroup.addImmutableValueMetric(registry.connectorType, connectorType()); - metricGroup.addImmutableValueMetric(registry.connectorClass, connector.getClass().getName()); - metricGroup.addImmutableValueMetric(registry.connectorVersion, connector.version()); + metricGroup.addImmutableValueMetric(registry.connectorClass, connector.pluginClass().getName()); + String version; + try { + version = connector.version(); + } catch (Exception e) { + version = PluginDesc.UNDEFINED_VERSION; + } + metricGroup.addImmutableValueMetric(registry.connectorVersion, version); metricGroup.addValueMetric(registry.connectorStatus, now -> state.toString().toLowerCase(Locale.getDefault())); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java index c4e5bb1558420..c94ebb0f7197d 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java @@ -52,6 +52,8 @@ import org.apache.kafka.connect.runtime.TaskStatus; import org.apache.kafka.connect.runtime.TooManyTasksException; import org.apache.kafka.connect.runtime.Worker; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; import org.apache.kafka.connect.runtime.rest.InternalRequestSignature; import org.apache.kafka.connect.runtime.rest.RestClient; import org.apache.kafka.connect.runtime.rest.entities.ConnectorInfo; @@ -62,10 +64,8 @@ import org.apache.kafka.connect.runtime.rest.entities.TaskInfo; import org.apache.kafka.connect.runtime.rest.errors.BadRequestException; import org.apache.kafka.connect.runtime.rest.errors.ConnectRestException; -import org.apache.kafka.connect.sink.SinkConnector; import org.apache.kafka.connect.source.ConnectorTransactionBoundaries; import org.apache.kafka.connect.source.ExactlyOnceSupport; -import org.apache.kafka.connect.source.SourceConnector; import org.apache.kafka.connect.source.SourceTask; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.storage.ConfigBackingStore; @@ -942,14 +942,14 @@ public void deleteConnectorConfig(final String connName, final Callback validateSinkConnectorConfig(SinkConnector connector, ConfigDef configDef, Map config) { + protected Map validateSinkConnectorConfig(IsolatedSinkConnector connector, ConfigDef configDef, Map config) { Map result = super.validateSinkConnectorConfig(connector, configDef, config); validateSinkConnectorGroupId(config, result); return result; } @Override - protected Map validateSourceConnectorConfig(SourceConnector connector, ConfigDef configDef, Map config) { + protected Map validateSourceConnectorConfig(IsolatedSourceConnector connector, ConfigDef configDef, Map config) { Map result = super.validateSourceConnectorConfig(connector, configDef, config); validateSourceConnectorExactlyOnceSupport(config, result, connector); validateSourceConnectorTransactionBoundary(config, result, connector); @@ -982,7 +982,7 @@ private void validateSinkConnectorGroupId(Map config, Map rawConfig, Map validatedConfig, - SourceConnector connector) { + IsolatedSourceConnector connector) { ConfigValue validatedExactlyOnceSupport = validatedConfig.get(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG); if (validatedExactlyOnceSupport.errorMessages().isEmpty()) { // Should be safe to parse the enum from the user-provided value since it's passed validation so far @@ -1029,7 +1029,7 @@ private void validateSourceConnectorExactlyOnceSupport( private void validateSourceConnectorTransactionBoundary( Map rawConfig, Map validatedConfig, - SourceConnector connector) { + IsolatedSourceConnector connector) { ConfigValue validatedTransactionBoundary = validatedConfig.get(SourceConnectorConfig.TRANSACTION_BOUNDARY_CONFIG); if (validatedTransactionBoundary.errorMessages().isEmpty()) { // Should be safe to parse the enum from the user-provided value since it's passed validation so far diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedConnector.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedConnector.java new file mode 100644 index 0000000000000..b8c9ee71994a1 --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedConnector.java @@ -0,0 +1,73 @@ +/* + * 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.connect.runtime.isolation; + +import org.apache.kafka.common.config.Config; +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.connect.connector.Connector; +import org.apache.kafka.connect.connector.ConnectorContext; +import org.apache.kafka.connect.connector.Task; + +import java.util.List; +import java.util.Map; + +public abstract class IsolatedConnector

extends IsolatedPlugin

{ + + IsolatedConnector(Plugins plugins, P delegate, PluginType type) { + super(plugins, delegate, type); + } + + public String version() throws Exception { + return isolate(delegate::version); + } + + public void initialize(ConnectorContext ctx) throws Exception { + isolate(() -> delegate.initialize(ctx)); + } + + public void initialize(ConnectorContext ctx, List> taskConfigs) throws Exception { + isolate(() -> delegate.initialize(ctx, taskConfigs)); + } + + public void reconfigure(Map props) throws Exception { + isolate(() -> delegate.reconfigure(props)); + } + + public Config validate(Map connectorConfigs) throws Exception { + return isolate(() -> delegate.validate(connectorConfigs)); + } + + public void start(Map props) throws Exception { + isolate(() -> delegate.start(props)); + } + + public Class taskClass() throws Exception { + return isolate(delegate::taskClass); + } + + public List> taskConfigs(int maxTasks) throws Exception { + return isolate(() -> delegate.taskConfigs(maxTasks)); + } + + public void stop() throws Exception { + isolate(delegate::stop); + } + + public ConfigDef config() throws Exception { + return isolate(delegate::config); + } +} diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedPlugin.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedPlugin.java new file mode 100644 index 0000000000000..adaa447375a15 --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedPlugin.java @@ -0,0 +1,98 @@ +/* + * 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.connect.runtime.isolation; + +import java.util.Objects; +import java.util.concurrent.Callable; + +public abstract class IsolatedPlugin

{ + + private final Plugins plugins; + private final Class pluginClass; + protected final P delegate; + private final ClassLoader classLoader; + private final PluginType type; + + IsolatedPlugin(Plugins plugins, P delegate, PluginType type) { + this.plugins = Objects.requireNonNull(plugins, "plugins must be non-null"); + this.delegate = Objects.requireNonNull(delegate, "delegate plugin must be non-null"); + this.pluginClass = delegate.getClass(); + this.classLoader = plugins.pluginLoader(delegate); + this.type = Objects.requireNonNull(type, "plugin type must be non-null"); + } + + public PluginType type() { + return type; + } + + @SuppressWarnings("unchecked") + public Class pluginClass() { + return (Class) pluginClass; + } + + protected V isolate(Callable callable) throws Exception { + try (LoaderSwap loaderSwap = plugins.withClassLoader(classLoader)) { + return callable.call(); + } + } + + protected void isolate(ThrowingRunnable runnable) throws Exception { + isolate(() -> { + runnable.run(); + return null; + }); + } + + public interface ThrowingRunnable { + void run() throws Exception; + } + + @Override + public int hashCode() { + return Objects.hash( + plugins, + pluginClass, + delegate, + classLoader, + type + ); + } + + @Override + public boolean equals(Object obj) { + if (obj == null || this.getClass() != obj.getClass()) { + return false; + } + IsolatedPlugin other = (IsolatedPlugin) obj; + return + this.plugins == other.plugins + && this.pluginClass == other.pluginClass + // use reference equality, as plugin implementations may mis-implement equals + && this.delegate == other.delegate + && this.classLoader == other.classLoader + && this.type == other.type; + } + + @Override + public String toString() { + try { + return isolate(delegate::toString); + } catch (Throwable e) { + return "unable to evaluate plugin toString: " + e.getMessage(); + } + } +} diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedSinkConnector.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedSinkConnector.java new file mode 100644 index 0000000000000..53dad9cc21a84 --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedSinkConnector.java @@ -0,0 +1,34 @@ +/* + * 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.connect.runtime.isolation; + +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.connect.sink.SinkConnector; + +import java.util.Map; + + +public class IsolatedSinkConnector extends IsolatedConnector { + + IsolatedSinkConnector(Plugins plugins, SinkConnector delegate) { + super(plugins, delegate, PluginType.SINK); + } + + public boolean alterOffsets(Map connectorConfig, Map offsets) throws Exception { + return isolate(() -> delegate.alterOffsets(connectorConfig, offsets)); + } +} diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedSourceConnector.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedSourceConnector.java new file mode 100644 index 0000000000000..da6c36dce0c98 --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/IsolatedSourceConnector.java @@ -0,0 +1,43 @@ +/* + * 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.connect.runtime.isolation; + +import org.apache.kafka.connect.source.ConnectorTransactionBoundaries; +import org.apache.kafka.connect.source.ExactlyOnceSupport; +import org.apache.kafka.connect.source.SourceConnector; + +import java.util.Map; + +public class IsolatedSourceConnector extends IsolatedConnector { + + IsolatedSourceConnector(Plugins plugins, SourceConnector delegate) { + super(plugins, delegate, PluginType.SOURCE); + } + + public ExactlyOnceSupport exactlyOnceSupport(Map connectorConfig) throws Exception { + return isolate(() -> delegate.exactlyOnceSupport(connectorConfig)); + } + + public ConnectorTransactionBoundaries canDefineTransactionBoundaries(Map connectorConfig) throws Exception { + return isolate(() -> delegate.canDefineTransactionBoundaries(connectorConfig)); + } + + public boolean alterOffsets(Map connectorConfig, Map, Map> offsets) throws Exception { + return isolate(() -> delegate.alterOffsets(connectorConfig, offsets)); + } + +} diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java index 681394f21af5e..73cb524c8c9db 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java @@ -247,6 +247,15 @@ public ClassLoader connectorLoader(String connectorClassOrAlias) { return delegatingLoader.connectorLoader(connectorClassOrAlias); } + public ClassLoader pluginLoader(Object delegate) { + ClassLoader classLoader = delegate.getClass().getClassLoader(); + if (classLoader instanceof PluginClassLoader) { + return classLoader; + } else { + return delegatingLoader; + } + } + @SuppressWarnings({"unchecked", "rawtypes"}) public Set> connectors() { Set> connectors = new TreeSet<>((Set) sinkConnectors()); @@ -287,7 +296,19 @@ public Object newPlugin(String classOrAlias) throws ClassNotFoundException { return newPlugin(klass); } - public Connector newConnector(String connectorClassOrAlias) { + public IsolatedConnector newConnector(String connectorClassOrAlias) { + Connector connector = newRawConnector(connectorClassOrAlias); + if (connector instanceof SourceConnector) { + return new IsolatedSourceConnector(this, (SourceConnector) connector); + } else if (connector instanceof SinkConnector) { + return new IsolatedSinkConnector(this, (SinkConnector) connector); + } else { + throw new IllegalArgumentException( + "Unknown connector " + connector.getClass().getName() + " does not subclass any known connector type"); + } + } + + private Connector newRawConnector(String connectorClassOrAlias) { Class klass = connectorClass(connectorClassOrAlias); return newPlugin(klass); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java index e1386bf6ac49c..cc365ba39db7b 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java @@ -237,7 +237,12 @@ private synchronized void putConnectorConfig(String connName, } requestExecutorService.submit(() -> { - updateConnectorTasks(connName); + try { + updateConnectorTasks(connName); + } catch (Exception e) { + callback.onCompletion(e, null); + return; + } callback.onCompletion(null, new Created<>(created, createConnectorInfo(connName))); }); }); @@ -263,7 +268,11 @@ public synchronized void requestTaskReconfiguration(String connName) { log.error("Task that requested reconfiguration does not exist: {}", connName); return; } - updateConnectorTasks(connName); + try { + updateConnectorTasks(connName); + } catch (Exception e) { + log.error("Unable to generate task configs for {}", connName, e); + } } @Override @@ -427,7 +436,7 @@ private void startConnector(String connName, Callback onStart) { worker.startConnector(connName, connConfigs, new HerderConnectorContext(this, connName), this, targetState, onStart); } - private List> recomputeTaskConfigs(String connName) { + private List> recomputeTaskConfigs(String connName) throws Exception { Map config = configState.connectorConfig(connName); ConnectorConfig connConfig = worker.isSinkConnector(connName) ? @@ -483,7 +492,7 @@ private void removeConnectorTasks(String connName) { } } - private synchronized void updateConnectorTasks(String connName) { + private synchronized void updateConnectorTasks(String connName) throws Exception { if (!worker.isRunning(connName)) { log.info("Skipping update of tasks for connector {} since it is not running", connName); return; @@ -546,7 +555,13 @@ public void onConnectorTargetStateChange(String connector) { } if (newState == TargetState.STARTED) { - requestExecutorService.submit(() -> updateConnectorTasks(connector)); + requestExecutorService.submit(() -> { + try { + updateConnectorTasks(connector); + } catch (Exception e) { + log.error("Unable to generate task configs for {}", connector, e); + } + }); } }); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java index 3a2c88b0893ed..d34c708e78b45 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java @@ -19,12 +19,9 @@ import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.common.InvalidRecordException; import org.apache.kafka.common.record.RecordBatch; -import org.apache.kafka.connect.connector.Connector; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.runtime.distributed.DistributedConfig; -import org.apache.kafka.connect.sink.SinkConnector; -import org.apache.kafka.connect.source.SourceConnector; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -137,14 +134,6 @@ public static void addMetricsContextProperties(Map prop, WorkerC } } - public static boolean isSinkConnector(Connector connector) { - return SinkConnector.class.isAssignableFrom(connector.getClass()); - } - - public static boolean isSourceConnector(Connector connector) { - return SourceConnector.class.isAssignableFrom(connector.getClass()); - } - /** * Apply a specified transformation {@link Function} to every value in a Map. * @param map the Map to be transformed diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java index 62283a02771aa..2abaed6279406 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java @@ -33,6 +33,9 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.errors.NotFoundException; import org.apache.kafka.connect.runtime.distributed.SampleConnectorClientConfigOverridePolicy; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; import org.apache.kafka.connect.runtime.isolation.LoaderSwap; import org.apache.kafka.connect.runtime.isolation.PluginDesc; import org.apache.kafka.connect.runtime.isolation.PluginType; @@ -47,6 +50,7 @@ import org.apache.kafka.connect.runtime.rest.entities.ConnectorStateInfo; import org.apache.kafka.connect.runtime.rest.entities.ConnectorType; import org.apache.kafka.connect.runtime.rest.errors.BadRequestException; +import org.apache.kafka.connect.source.SourceConnector; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.storage.ConfigBackingStore; import org.apache.kafka.connect.storage.StatusBackingStore; @@ -83,8 +87,10 @@ import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertThrows; +import static org.mockito.ArgumentMatchers.any; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doReturn; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.CALLS_REAL_METHODS; import static org.mockito.Mockito.doAnswer; @@ -194,12 +200,13 @@ public void testConnectorClientConfigOverridePolicyClose() { } @Test - public void testConnectorStatus() { + public void testConnectorStatus() throws Exception { ConnectorTaskId taskId = new ConnectorTaskId(connectorName, 0); AbstractHerder herder = testHerder(); - when(plugins.newConnector(anyString())).thenReturn(new SampleSourceConnector()); + IsolatedConnector connector = mockConnector(SampleSourceConnector.class); + doReturn(connector).when(plugins).newConnector(anyString()); when(herder.plugins()).thenReturn(plugins); when(herder.rawConfig(connectorName)).thenReturn(Collections.singletonMap( @@ -261,10 +268,11 @@ public void testConnectorStatusMissingPlugin() { } @Test - public void testConnectorInfo() { + public void testConnectorInfo() throws Exception { AbstractHerder herder = testHerder(); - when(plugins.newConnector(anyString())).thenReturn(new SampleSourceConnector()); + IsolatedConnector connector = mockConnector(SampleSourceConnector.class); + doReturn(connector).when(plugins).newConnector(anyString()); when(herder.plugins()).thenReturn(plugins); when(configStore.snapshot()).thenReturn(SNAPSHOT); @@ -351,7 +359,7 @@ public void testBuildRestartPlanForUnknownConnector() { } @Test - public void testConfigValidationNullConfig() { + public void testConfigValidationNullConfig() throws Exception { AbstractHerder herder = createConfigValidationHerder(SampleSourceConnector.class, noneConnectorClientConfigOverridePolicy); Map config = new HashMap<>(); @@ -368,7 +376,7 @@ public void testConfigValidationNullConfig() { } @Test - public void testConfigValidationMultipleNullConfig() { + public void testConfigValidationMultipleNullConfig() throws Exception { AbstractHerder herder = createConfigValidationHerder(SampleSourceConnector.class, noneConnectorClientConfigOverridePolicy); Map config = new HashMap<>(); @@ -445,7 +453,7 @@ public void testBuildRestartPlanForNoRestart() { } @Test - public void testConfigValidationEmptyConfig() { + public void testConfigValidationEmptyConfig() throws Exception { AbstractHerder herder = createConfigValidationHerder(SampleSourceConnector.class, noneConnectorClientConfigOverridePolicy, 0); assertThrows(BadRequestException.class, () -> herder.validateConnectorConfig(Collections.emptyMap(), s -> null, false)); @@ -453,7 +461,7 @@ public void testConfigValidationEmptyConfig() { } @Test - public void testConfigValidationMissingName() { + public void testConfigValidationMissingName() throws Exception { final Class connectorClass = SampleSourceConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, noneConnectorClientConfigOverridePolicy); @@ -484,7 +492,7 @@ public void testConfigValidationMissingName() { } @Test - public void testConfigValidationInvalidTopics() { + public void testConfigValidationInvalidTopics() throws Exception { final Class connectorClass = SampleSinkConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, noneConnectorClientConfigOverridePolicy); @@ -507,7 +515,7 @@ public void testConfigValidationInvalidTopics() { } @Test - public void testConfigValidationTopicsWithDlq() { + public void testConfigValidationTopicsWithDlq() throws Exception { final Class connectorClass = SampleSinkConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, noneConnectorClientConfigOverridePolicy); @@ -526,7 +534,7 @@ public void testConfigValidationTopicsWithDlq() { } @Test - public void testConfigValidationTopicsRegexWithDlq() { + public void testConfigValidationTopicsRegexWithDlq() throws Exception { final Class connectorClass = SampleSinkConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, noneConnectorClientConfigOverridePolicy); @@ -545,7 +553,7 @@ public void testConfigValidationTopicsRegexWithDlq() { } @Test - public void testConfigValidationTransformsExtendResults() { + public void testConfigValidationTransformsExtendResults() throws Exception { final Class connectorClass = SampleSourceConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, noneConnectorClientConfigOverridePolicy); @@ -597,7 +605,7 @@ public void testConfigValidationTransformsExtendResults() { } @Test - public void testConfigValidationPredicatesExtendResults() { + public void testConfigValidationPredicatesExtendResults() throws Exception { final Class connectorClass = SampleSourceConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, noneConnectorClientConfigOverridePolicy); @@ -669,7 +677,7 @@ private PluginDesc> transformationPluginDesc() { } @Test - public void testConfigValidationPrincipalOnlyOverride() { + public void testConfigValidationPrincipalOnlyOverride() throws Exception { final Class connectorClass = SampleSourceConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, new PrincipalConnectorClientConfigOverridePolicy()); @@ -711,7 +719,7 @@ public void testConfigValidationPrincipalOnlyOverride() { } @Test - public void testConfigValidationAllOverride() { + public void testConfigValidationAllOverride() throws Exception { final Class connectorClass = SampleSourceConnector.class; AbstractHerder herder = createConfigValidationHerder(connectorClass, new AllConnectorClientConfigOverridePolicy()); @@ -1164,13 +1172,13 @@ private void testConfigProviderRegex(String rawConnConfig, boolean expected) { } private AbstractHerder createConfigValidationHerder(Class connectorClass, - ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { + ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) throws Exception { return createConfigValidationHerder(connectorClass, connectorClientConfigOverridePolicy, 1); } private AbstractHerder createConfigValidationHerder(Class connectorClass, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, - int countOfCallingNewConnector) { + int countOfCallingNewConnector) throws Exception { AbstractHerder herder = testHerder(connectorClientConfigOverridePolicy); @@ -1180,16 +1188,32 @@ private AbstractHerder createConfigValidationHerder(Class c final ArgumentCaptor> mapArgumentCaptor = ArgumentCaptor.forClass(Map.class); when(transformer.transform(mapArgumentCaptor.capture())).thenAnswer(invocation -> mapArgumentCaptor.getValue()); when(worker.getPlugins()).thenReturn(plugins); + IsolatedConnector isolatedConnector = mockConnector(connectorClass); + if (countOfCallingNewConnector > 0) { + mockValidationIsolation(connectorClass.getName(), isolatedConnector); + } + return herder; + } + + private IsolatedConnector mockConnector(Class connectorClass) throws Exception { final Connector connector; try { connector = connectorClass.getConstructor().newInstance(); } catch (ReflectiveOperationException e) { throw new RuntimeException("Couldn't create connector", e); } - if (countOfCallingNewConnector > 0) { - mockValidationIsolation(connectorClass.getName(), connector); + IsolatedConnector isolatedConnector; + if (SourceConnector.class.isAssignableFrom(connectorClass)) { + isolatedConnector = mock(IsolatedSourceConnector.class); + when(isolatedConnector.type()).thenReturn(PluginType.SOURCE); + } else { + isolatedConnector = mock(IsolatedSinkConnector.class); + when(isolatedConnector.type()).thenReturn(PluginType.SINK); } - return herder; + doReturn(connectorClass).when(isolatedConnector).pluginClass(); + when(isolatedConnector.config()).thenReturn(connector.config()); + when(isolatedConnector.validate(any())).thenAnswer(invocation -> connector.validate(invocation.getArgument(0))); + return isolatedConnector; } private AbstractHerder testHerder() { @@ -1202,8 +1226,9 @@ private AbstractHerder testHerder(ConnectorClientConfigOverridePolicy connectorC .defaultAnswer(CALLS_REAL_METHODS)); } - private void mockValidationIsolation(String connectorClass, Connector connector) { - when(plugins.newConnector(connectorClass)).thenReturn(connector); + @SuppressWarnings({"unchecked", "rawtypes"}) + private void mockValidationIsolation(String connectorClass, IsolatedConnector connector) { + when(plugins.newConnector(connectorClass)).thenReturn((IsolatedConnector) connector); when(plugins.connectorLoader(connectorClass)).thenReturn(classLoader); when(plugins.withClassLoader(classLoader)).thenReturn(loaderSwap); } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerConnectorTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerConnectorTest.java index f91e37ef9f9f9..f5fee8df6c8ef 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerConnectorTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerConnectorTest.java @@ -20,10 +20,12 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.health.ConnectorType; import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; +import org.apache.kafka.connect.runtime.isolation.PluginType; import org.apache.kafka.connect.runtime.isolation.Plugins; -import org.apache.kafka.connect.sink.SinkConnector; import org.apache.kafka.connect.sink.SinkConnectorContext; -import org.apache.kafka.connect.source.SourceConnector; import org.apache.kafka.connect.source.SourceConnectorContext; import org.apache.kafka.connect.storage.CloseableOffsetStorageReader; import org.apache.kafka.connect.storage.ConnectorOffsetBackingStore; @@ -43,6 +45,7 @@ import org.mockito.ArgumentCaptor; import org.mockito.InOrder; import org.mockito.Mock; +import org.mockito.Mockito; import org.mockito.junit.MockitoJUnit; import org.mockito.junit.MockitoRule; import org.mockito.quality.Strictness; @@ -86,7 +89,7 @@ public class WorkerConnectorTest { @Mock private ClassLoader classLoader; private final ConnectorType connectorType; - private final Connector connector; + private final IsolatedConnector connector; private final CloseableOffsetStorageReader offsetStorageReader; private final ConnectorOffsetBackingStore offsetStore; @@ -95,22 +98,27 @@ public static Collection parameters() { return Arrays.asList(ConnectorType.SOURCE, ConnectorType.SINK); } - public WorkerConnectorTest(ConnectorType connectorType) { + public WorkerConnectorTest(ConnectorType connectorType) throws Exception { this.connectorType = connectorType; switch (connectorType) { case SINK: - this.connector = mock(SinkConnector.class); + this.connector = mock(IsolatedSinkConnector.class); this.offsetStorageReader = null; this.offsetStore = null; + Mockito.>when(connector.pluginClass()).thenReturn(SampleSinkConnector.class); + when(connector.type()).thenReturn(PluginType.SINK); break; case SOURCE: - this.connector = mock(SourceConnector.class); + this.connector = mock(IsolatedSourceConnector.class); this.offsetStorageReader = mock(CloseableOffsetStorageReader.class); this.offsetStore = mock(ConnectorOffsetBackingStore.class); + Mockito.>when(connector.pluginClass()).thenReturn(SampleSourceConnector.class); + when(connector.type()).thenReturn(PluginType.SOURCE); break; default: throw new IllegalStateException("Unexpected connector type: " + connectorType); } + when(connector.version()).thenReturn(VERSION); } @Before @@ -125,10 +133,9 @@ public void tearDown() { } @Test - public void testInitializeFailure() { + public void testInitializeFailure() throws Exception { RuntimeException exception = new RuntimeException(); - when(connector.version()).thenReturn(VERSION); doThrow(exception).when(connector).initialize(any()); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -145,10 +152,9 @@ public void testInitializeFailure() { } @Test - public void testFailureIsFinalState() { + public void testFailureIsFinalState() throws Exception { RuntimeException exception = new RuntimeException(); - when(connector.version()).thenReturn(VERSION); doThrow(exception).when(connector).initialize(any()); Callback onStateChange = mockCallback(); @@ -172,9 +178,7 @@ public void testFailureIsFinalState() { } @Test - public void testStartupAndShutdown() { - when(connector.version()).thenReturn(VERSION); - + public void testStartupAndShutdown() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -196,9 +200,7 @@ public void testStartupAndShutdown() { } @Test - public void testStartupAndPause() { - when(connector.version()).thenReturn(VERSION); - + public void testStartupAndPause() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -225,9 +227,7 @@ public void testStartupAndPause() { } @Test - public void testStartupAndStop() { - when(connector.version()).thenReturn(VERSION); - + public void testStartupAndStop() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -255,9 +255,7 @@ public void testStartupAndStop() { } @Test - public void testOnResume() { - when(connector.version()).thenReturn(VERSION); - + public void testOnResume() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -285,9 +283,7 @@ public void testOnResume() { } @Test - public void testStartupPaused() { - when(connector.version()).thenReturn(VERSION); - + public void testStartupPaused() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -309,9 +305,7 @@ public void testStartupPaused() { } @Test - public void testStartupStopped() { - when(connector.version()).thenReturn(VERSION); - + public void testStartupStopped() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -333,10 +327,9 @@ public void testStartupStopped() { } @Test - public void testStartupFailure() { + public void testStartupFailure() throws Exception { RuntimeException exception = new RuntimeException(); - when(connector.version()).thenReturn(VERSION); doThrow(exception).when(connector).start(CONFIG); Callback onStateChange = mockCallback(); @@ -360,11 +353,9 @@ public void testStartupFailure() { } @Test - public void testStopFailure() { + public void testStopFailure() throws Exception { RuntimeException exception = new RuntimeException(); - when(connector.version()).thenReturn(VERSION); - // Fail during the first call to stop, then succeed for the next attempt doThrow(exception).doNothing().when(connector).stop(); @@ -402,11 +393,9 @@ public void testStopFailure() { } @Test - public void testShutdownFailure() { + public void testShutdownFailure() throws Exception { RuntimeException exception = new RuntimeException(); - when(connector.version()).thenReturn(VERSION); - doThrow(exception).when(connector).stop(); Callback onStateChange = mockCallback(); @@ -430,9 +419,7 @@ public void testShutdownFailure() { } @Test - public void testTransitionStartedToStarted() { - when(connector.version()).thenReturn(VERSION); - + public void testTransitionStartedToStarted() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -457,9 +444,7 @@ public void testTransitionStartedToStarted() { } @Test - public void testTransitionPausedToPaused() { - when(connector.version()).thenReturn(VERSION); - + public void testTransitionPausedToPaused() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -488,9 +473,7 @@ public void testTransitionPausedToPaused() { } @Test - public void testTransitionStoppedToStopped() { - when(connector.version()).thenReturn(VERSION); - + public void testTransitionStoppedToStopped() throws Exception { Callback onStateChange = mockCallback(); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, connector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); @@ -519,9 +502,11 @@ public void testTransitionStoppedToStopped() { } @Test - public void testFailConnectorThatIsNeitherSourceNorSink() { - Connector badConnector = mock(Connector.class); + public void testFailConnectorThatIsNeitherSourceNorSink() throws Exception { + IsolatedConnector badConnector = mock(IsolatedConnector.class); when(badConnector.version()).thenReturn(VERSION); + Mockito.>when(badConnector.pluginClass()).thenReturn(SampleSourceConnector.class); + when(badConnector.type()).thenReturn(PluginType.TRANSFORMATION); WorkerConnector workerConnector = new WorkerConnector(CONNECTOR, badConnector, connectorConfig, ctx, metrics, listener, offsetStorageReader, offsetStore, classLoader); workerConnector.initialize(); @@ -610,7 +595,8 @@ private Callback mockCallback() { return mock(Callback.class); } - private void verifyInitialize() { + private void verifyInitialize() throws Exception { + verify(connector).pluginClass(); verify(connector).version(); if (connectorType == ConnectorType.SOURCE) { verify(offsetStore).start(); @@ -620,15 +606,15 @@ private void verifyInitialize() { } } - private void verifyCleanShutdown(boolean started) { + private void verifyCleanShutdown(boolean started) throws Exception { verifyShutdown(true, started); } - private void verifyShutdown(boolean clean, boolean started) { + private void verifyShutdown(boolean clean, boolean started) throws Exception { verifyShutdown(1, clean, started); } - private void verifyShutdown(int connectorStops, boolean clean, boolean started) { + private void verifyShutdown(int connectorStops, boolean clean, boolean started) throws Exception { verify(ctx).close(); if (connectorType == ConnectorType.SOURCE) { verify(offsetStorageReader).close(); 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 4579794a2c4ff..906177b58bf2d 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 @@ -57,8 +57,13 @@ import org.apache.kafka.connect.json.JsonConverter; import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; import org.apache.kafka.connect.runtime.MockConnectMetrics.MockMetricsReporter; -import org.apache.kafka.connect.runtime.distributed.DistributedConfig; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; import org.apache.kafka.connect.runtime.isolation.LoaderSwap; +import org.apache.kafka.connect.runtime.isolation.PluginType; +import org.apache.kafka.connect.storage.ClusterConfigState; +import org.apache.kafka.connect.runtime.distributed.DistributedConfig; import org.apache.kafka.connect.runtime.isolation.PluginClassLoader; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.isolation.Plugins.ClassLoaderUsage; @@ -67,14 +72,11 @@ import org.apache.kafka.connect.runtime.rest.entities.ConnectorStateInfo; import org.apache.kafka.connect.runtime.rest.entities.Message; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; -import org.apache.kafka.connect.sink.SinkConnector; import org.apache.kafka.connect.sink.SinkRecord; import org.apache.kafka.connect.sink.SinkTask; -import org.apache.kafka.connect.source.SourceConnector; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.source.SourceTask; import org.apache.kafka.connect.storage.CloseableOffsetStorageReader; -import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.storage.ConnectorOffsetBackingStore; import org.apache.kafka.connect.storage.Converter; import org.apache.kafka.connect.storage.HeaderConverter; @@ -101,6 +103,7 @@ import org.mockito.MockitoSession; import org.mockito.invocation.InvocationOnMock; import org.mockito.quality.Strictness; +import org.mockito.stubbing.OngoingStubbing; import javax.management.MBeanServer; import javax.management.ObjectName; @@ -220,10 +223,10 @@ public class WorkerTest { private StatusBackingStore statusBackingStore; @Mock - private SourceConnector sourceConnector; + private IsolatedSourceConnector sourceConnector; @Mock - private SinkConnector sinkConnector; + private IsolatedSinkConnector sinkConnector; @Mock private CloseableConnectorContext ctx; @@ -1539,7 +1542,7 @@ public void testOffsetStoreForRegularSourceTask() { producerProps.put(BOOTSTRAP_SERVERS_CONFIG, workerBootstrapServers); // With no connector-specific offsets topic in the config, we should only use the worker-global store // Pass in a null topic admin to make sure that with these parameters, the method doesn't require a topic admin - ConnectorOffsetBackingStore connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithoutOffsetsTopic, sourceConnector.getClass(), producer, producerProps, null); + ConnectorOffsetBackingStore connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithoutOffsetsTopic, sourceConnector.pluginClass(), producer, producerProps, null); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertFalse(connectorStore.hasConnectorSpecificStore()); @@ -1549,13 +1552,13 @@ public void testOffsetStoreForRegularSourceTask() { final SourceConnectorConfig sourceConfigWithOffsetsTopic = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With a connector-specific offsets topic in the config (whose name differs from the worker's offsets topic), we should use both a // connector-specific store and the worker-global store - connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithOffsetsTopic, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithOffsetsTopic, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); assertThrows(NullPointerException.class, () -> worker.offsetStoreForRegularSourceTask( - TASK_ID, sourceConfigWithOffsetsTopic, sourceConnector.getClass(), producer, producerProps, null + TASK_ID, sourceConfigWithOffsetsTopic, sourceConnector.pluginClass(), producer, producerProps, null ) ); connectorStore.stop(); @@ -1564,14 +1567,14 @@ public void testOffsetStoreForRegularSourceTask() { final SourceConnectorConfig sourceConfigWithSameOffsetsTopicAsWorker = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With a connector-specific offsets topic in the config whose name matches the worker's offsets topic, and no overridden bootstrap.servers // for the connector, we should only use a connector-specific store - connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertFalse(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); assertThrows( NullPointerException.class, () -> worker.offsetStoreForRegularSourceTask( - TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.getClass(), producer, producerProps, null + TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.pluginClass(), producer, producerProps, null ) ); connectorStore.stop(); @@ -1579,14 +1582,14 @@ public void testOffsetStoreForRegularSourceTask() { producerProps.put(BOOTSTRAP_SERVERS_CONFIG, workerBootstrapServers); // With a connector-specific offsets topic in the config whose name matches the worker's offsets topic, and an overridden bootstrap.servers // for the connector that exactly matches the worker's, we should only use a connector-specific store - connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertFalse(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); assertThrows( NullPointerException.class, () -> worker.offsetStoreForRegularSourceTask( - TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.getClass(), producer, producerProps, null + TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.pluginClass(), producer, producerProps, null ) ); connectorStore.stop(); @@ -1594,14 +1597,14 @@ public void testOffsetStoreForRegularSourceTask() { producerProps.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:1111"); // With a connector-specific offsets topic in the config whose name matches the worker's offsets topic, and an overridden bootstrap.servers // for the connector that doesn't match the worker's, we should use both a connector-specific store and the worker-global store - connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); assertThrows( NullPointerException.class, () -> worker.offsetStoreForRegularSourceTask( - TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.getClass(), producer, producerProps, null + TASK_ID, sourceConfigWithSameOffsetsTopicAsWorker, sourceConnector.pluginClass(), producer, producerProps, null ) ); connectorStore.stop(); @@ -1610,7 +1613,7 @@ public void testOffsetStoreForRegularSourceTask() { // With no connector-specific offsets topic in the config and an overridden bootstrap.servers // for the connector that doesn't match the worker's, we should still only use the worker-global store // Pass in a null topic admin to make sure that with these parameters, the method doesn't require a topic admin - connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithoutOffsetsTopic, sourceConnector.getClass(), producer, producerProps, null); + connectorStore = worker.offsetStoreForRegularSourceTask(TASK_ID, sourceConfigWithoutOffsetsTopic, sourceConnector.pluginClass(), producer, producerProps, null); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertFalse(connectorStore.hasConnectorSpecificStore()); @@ -1653,7 +1656,7 @@ public void testOffsetStoreForExactlyOnceSourceTask() { SourceConnectorConfig sourceConfig = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); producerProps.put(BOOTSTRAP_SERVERS_CONFIG, workerBootstrapServers); // With no connector-specific offsets topic in the config, we should only use a connector-specific offsets store - ConnectorOffsetBackingStore connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.getClass(), producer, producerProps, topicAdmin); + ConnectorOffsetBackingStore connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertFalse(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); @@ -1663,7 +1666,7 @@ public void testOffsetStoreForExactlyOnceSourceTask() { sourceConfig = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With a connector-specific offsets topic in the config (whose name differs from the worker's offsets topic), we should use both a // connector-specific store and the worker-global store - connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); @@ -1673,7 +1676,7 @@ public void testOffsetStoreForExactlyOnceSourceTask() { sourceConfig = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With a connector-specific offsets topic in the config whose name matches the worker's offsets topic, and no overridden bootstrap.servers // for the connector, we should only use a connector-specific store - connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertFalse(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); @@ -1683,7 +1686,7 @@ public void testOffsetStoreForExactlyOnceSourceTask() { sourceConfig = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With a connector-specific offsets topic in the config whose name matches the worker's offsets topic, and an overridden bootstrap.servers // for the connector that exactly matches the worker's, we should only use a connector-specific store - connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertFalse(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); @@ -1693,7 +1696,7 @@ public void testOffsetStoreForExactlyOnceSourceTask() { sourceConfig = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With a connector-specific offsets topic in the config whose name matches the worker's offsets topic, and an overridden bootstrap.servers // for the connector that doesn't match the worker's, we should use both a connector-specific store and the worker-global store - connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); @@ -1703,7 +1706,7 @@ public void testOffsetStoreForExactlyOnceSourceTask() { sourceConfig = new SourceConnectorConfig(plugins, connectorProps, enableTopicCreation); // With no connector-specific offsets topic in the config and an overridden bootstrap.servers // for the connector that doesn't match the worker's, we should use both a connector-specific store and the worker-global store - connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.getClass(), producer, producerProps, topicAdmin); + connectorStore = worker.offsetStoreForExactlyOnceSourceTask(TASK_ID, sourceConfig, sourceConnector.pluginClass(), producer, producerProps, topicAdmin); connectorStore.configure(config); assertTrue(connectorStore.hasWorkerGlobalStore()); assertTrue(connectorStore.hasConnectorSpecificStore()); @@ -2015,7 +2018,8 @@ public void testGetSourceConnectorOffsetsError() { } @Test - public void testAlterOffsetsConnectorDoesNotSupportOffsetAlteration() { + @SuppressWarnings({"unchecked", "rawtypes"}) + public void testAlterOffsetsConnectorDoesNotSupportOffsetAlteration() throws Exception { mockKafkaClusterId(); mockInternalConverters(); @@ -2024,7 +2028,7 @@ public void testAlterOffsetsConnectorDoesNotSupportOffsetAlteration() { worker.start(); mockGenericIsolation(); - when(plugins.newConnector(anyString())).thenReturn(sourceConnector); + when(plugins.newConnector(anyString())).thenReturn((IsolatedConnector) sourceConnector); when(plugins.withClassLoader(any(ClassLoader.class), any(Runnable.class))).thenAnswer(AdditionalAnswers.returnsSecondArg()); when(sourceConnector.alterOffsets(eq(connectorProps), anyMap())).thenThrow(new UnsupportedOperationException("This connector doesn't " + "support altering of offsets")); @@ -2890,13 +2894,17 @@ private void verifyGenericIsolation() { verify(loaderSwap, atLeastOnce()).close(); } - private void mockConnectorIsolation(String connectorClass, Connector connector) { + private void mockConnectorIsolation(String connectorClass, IsolatedConnector connector) throws Exception { mockGenericIsolation(); - when(plugins.newConnector(connectorClass)).thenReturn(connector); + OngoingStubbing> newConnector = when(plugins.newConnector(connectorClass)); + newConnector.thenReturn(connector); + OngoingStubbing> pluginClass = when(connector.pluginClass()); + pluginClass.thenReturn(SampleSourceConnector.class); + when(connector.type()).thenReturn(PluginType.SOURCE); when(connector.version()).thenReturn("1.0"); } - private void verifyConnectorIsolation(Connector connector) { + private void verifyConnectorIsolation(IsolatedConnector connector) throws Exception { verifyGenericIsolation(); verify(plugins).newConnector(anyString()); verify(connector, atLeastOnce()).version(); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java index fa05e55015efd..2e67c8a7f6c8d 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java @@ -46,6 +46,8 @@ import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.runtime.WorkerConfigTransformer; import org.apache.kafka.connect.runtime.distributed.DistributedHerder.HerderMetrics; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.rest.InternalRequestSignature; import org.apache.kafka.connect.runtime.rest.RestClient; @@ -59,7 +61,6 @@ import org.apache.kafka.connect.runtime.rest.entities.TaskInfo; import org.apache.kafka.connect.runtime.rest.errors.BadRequestException; import org.apache.kafka.connect.runtime.rest.errors.ConnectRestException; -import org.apache.kafka.connect.sink.SinkConnector; import org.apache.kafka.connect.source.ConnectorTransactionBoundaries; import org.apache.kafka.connect.source.ExactlyOnceSupport; import org.apache.kafka.connect.source.SourceConnector; @@ -568,16 +569,16 @@ public void testRebalanceFailedConnector() throws Exception { } @Test - public void testRevoke() throws TimeoutException { + public void testRevoke() throws Exception { revokeAndReassign(false); } @Test - public void testIncompleteRebalanceBeforeRevoke() throws TimeoutException { + public void testIncompleteRebalanceBeforeRevoke() throws Exception { revokeAndReassign(true); } - public void revokeAndReassign(boolean incompleteRebalance) throws TimeoutException { + public void revokeAndReassign(boolean incompleteRebalance) throws Exception { connectProtocolVersion = CONNECT_PROTOCOL_V1; int configOffset = 1; @@ -887,7 +888,7 @@ public void testConnectorNameConflictsWithWorkerGroupId() { Map config = new HashMap<>(CONN2_CONFIG); config.put(ConnectorConfig.NAME_CONFIG, "test-group"); - SinkConnector connectorMock = mock(SinkConnector.class); + IsolatedSinkConnector connectorMock = mock(IsolatedSinkConnector.class); // CONN2 creation should fail because the worker group id (connect-test-group) conflicts with // the consumer group id we would use for this sink @@ -906,7 +907,7 @@ public void testConnectorGroupIdConflictsWithWorkerGroupId() { Map config = new HashMap<>(CONN2_CONFIG); config.put(overriddenGroupId, "connect-test-group"); - SinkConnector connectorMock = mock(SinkConnector.class); + IsolatedSinkConnector connectorMock = mock(IsolatedSinkConnector.class); // CONN2 creation should fail because the worker group id (connect-test-group) conflicts with // the consumer group id we would use for this sink @@ -918,11 +919,6 @@ public void testConnectorGroupIdConflictsWithWorkerGroupId() { Collections.singletonList("Consumer group connect-test-group conflicts with Connect worker group connect-test-group"), overriddenGroupIdConfig.errorMessages()); - ConfigValue nameConfig = validatedConfigs.get(ConnectorConfig.NAME_CONFIG); - assertEquals( - Collections.emptyList(), - nameConfig.errorMessages() - ); } @Test @@ -3184,7 +3180,7 @@ public void testHerderStopServicesClosesUponShutdown() { } @Test - public void testPollDurationOnSlowConnectorOperations() { + public void testPollDurationOnSlowConnectorOperations() throws Exception { connectProtocolVersion = CONNECT_PROTOCOL_V1; // If an operation during tick() takes some amount of time, that time should count against the rebalance delay final int rebalanceDelayMs = 20000; @@ -3250,7 +3246,7 @@ public void shouldThrowWhenStartAndStopExecutorThrowsRejectedExecutionExceptionA } @Test - public void testTaskReconfigurationRetriesWithConnectorTaskConfigsException() { + public void testTaskReconfigurationRetriesWithConnectorTaskConfigsException() throws Exception { when(member.memberId()).thenReturn("leader"); when(member.currentProtocolVersion()).thenReturn(CONNECT_PROTOCOL_V0); expectRebalance(1, Collections.emptyList(), Collections.emptyList(), true); @@ -3270,7 +3266,7 @@ public void testTaskReconfigurationRetriesWithConnectorTaskConfigsException() { } @Test - public void testTaskReconfigurationNoRetryWithTooManyTasks() { + public void testTaskReconfigurationNoRetryWithTooManyTasks() throws Exception { // initial tick when(member.memberId()).thenReturn("leader"); when(member.currentProtocolVersion()).thenReturn(CONNECT_PROTOCOL_V0); @@ -3311,7 +3307,7 @@ public void testTaskReconfigurationNoRetryWithTooManyTasks() { } @Test - public void testTaskReconfigurationRetriesWithLeaderRequestForwardingException() { + public void testTaskReconfigurationRetriesWithLeaderRequestForwardingException() throws Exception { herder = mock(DistributedHerder.class, withSettings().defaultAnswer(CALLS_REAL_METHODS).useConstructor(new DistributedConfig(HERDER_CONFIG), worker, WORKER_ID, KAFKA_CLUSTER_ID, statusBackingStore, configBackingStore, member, MEMBER_URL, restClient, metrics, time, noneConnectorClientConfigOverridePolicy, Collections.emptyList(), new MockSynchronousExecutor(), new AutoCloseable[]{})); @@ -3424,12 +3420,12 @@ public void preserveHighestImpactRestartRequest() { } @Test - public void testExactlyOnceSourceSupportValidation() { + public void testExactlyOnceSourceSupportValidation() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG, REQUIRED.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); when(connectorMock.exactlyOnceSupport(eq(config))).thenReturn(ExactlyOnceSupport.SUPPORTED); Map validatedConfigs = herder.validateSourceConnectorConfig( @@ -3440,12 +3436,12 @@ public void testExactlyOnceSourceSupportValidation() { } @Test - public void testExactlyOnceSourceSupportValidationOnUnsupportedConnector() { + public void testExactlyOnceSourceSupportValidationOnUnsupportedConnector() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG, REQUIRED.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); when(connectorMock.exactlyOnceSupport(eq(config))).thenReturn(ExactlyOnceSupport.UNSUPPORTED); Map validatedConfigs = herder.validateSourceConnectorConfig( @@ -3458,12 +3454,12 @@ public void testExactlyOnceSourceSupportValidationOnUnsupportedConnector() { } @Test - public void testExactlyOnceSourceSupportValidationOnUnknownConnector() { + public void testExactlyOnceSourceSupportValidationOnUnknownConnector() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG, REQUIRED.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); when(connectorMock.exactlyOnceSupport(eq(config))).thenReturn(null); Map validatedConfigs = herder.validateSourceConnectorConfig( @@ -3478,12 +3474,12 @@ public void testExactlyOnceSourceSupportValidationOnUnknownConnector() { } @Test - public void testExactlyOnceSourceSupportValidationHandlesConnectorErrorsGracefully() { + public void testExactlyOnceSourceSupportValidationHandlesConnectorErrorsGracefully() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG, REQUIRED.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); String errorMessage = "time to add a new unit test :)"; when(connectorMock.exactlyOnceSupport(eq(config))).thenThrow(new NullPointerException(errorMessage)); @@ -3499,11 +3495,11 @@ public void testExactlyOnceSourceSupportValidationHandlesConnectorErrorsGraceful } @Test - public void testExactlyOnceSourceSupportValidationWhenExactlyOnceNotEnabledOnWorker() { + public void testExactlyOnceSourceSupportValidationWhenExactlyOnceNotEnabledOnWorker() throws Exception { Map config = new HashMap<>(); config.put(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG, REQUIRED.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); when(connectorMock.exactlyOnceSupport(eq(config))).thenReturn(ExactlyOnceSupport.SUPPORTED); Map validatedConfigs = herder.validateSourceConnectorConfig( @@ -3521,7 +3517,7 @@ public void testExactlyOnceSourceSupportValidationHandlesInvalidValuesGracefully Map config = new HashMap<>(); config.put(SourceConnectorConfig.EXACTLY_ONCE_SUPPORT_CONFIG, "invalid"); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); Map validatedConfigs = herder.validateSourceConnectorConfig( connectorMock, SourceConnectorConfig.configDef(), config); @@ -3535,12 +3531,12 @@ public void testExactlyOnceSourceSupportValidationHandlesInvalidValuesGracefully } @Test - public void testConnectorTransactionBoundaryValidation() { + public void testConnectorTransactionBoundaryValidation() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.TRANSACTION_BOUNDARY_CONFIG, CONNECTOR.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); when(connectorMock.canDefineTransactionBoundaries(eq(config))) .thenReturn(ConnectorTransactionBoundaries.SUPPORTED); @@ -3552,12 +3548,12 @@ public void testConnectorTransactionBoundaryValidation() { } @Test - public void testConnectorTransactionBoundaryValidationOnUnsupportedConnector() { + public void testConnectorTransactionBoundaryValidationOnUnsupportedConnector() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.TRANSACTION_BOUNDARY_CONFIG, CONNECTOR.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); when(connectorMock.canDefineTransactionBoundaries(eq(config))) .thenReturn(ConnectorTransactionBoundaries.UNSUPPORTED); @@ -3573,12 +3569,12 @@ public void testConnectorTransactionBoundaryValidationOnUnsupportedConnector() { } @Test - public void testConnectorTransactionBoundaryValidationHandlesConnectorErrorsGracefully() { + public void testConnectorTransactionBoundaryValidationHandlesConnectorErrorsGracefully() throws Exception { herder = exactlyOnceHerder(); Map config = new HashMap<>(); config.put(SourceConnectorConfig.TRANSACTION_BOUNDARY_CONFIG, CONNECTOR.toString()); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); String errorMessage = "Wait I thought we tested for this?"; when(connectorMock.canDefineTransactionBoundaries(eq(config))).thenThrow(new ConnectException(errorMessage)); @@ -3599,7 +3595,7 @@ public void testConnectorTransactionBoundaryValidationHandlesInvalidValuesGracef Map config = new HashMap<>(); config.put(SourceConnectorConfig.TRANSACTION_BOUNDARY_CONFIG, "CONNECTOR.toString()"); - SourceConnector connectorMock = mock(SourceConnector.class); + IsolatedSourceConnector connectorMock = mock(IsolatedSourceConnector.class); Map validatedConfigs = herder.validateSourceConnectorConfig( connectorMock, SourceConnectorConfig.configDef(), config); @@ -4024,7 +4020,7 @@ private ClusterConfigState exactlyOnceSnapshot( Collections.emptySet()); } - private void expectExecuteTaskReconfiguration(boolean running, ConnectorConfig connectorConfig, Answer>> answer) { + private void expectExecuteTaskReconfiguration(boolean running, ConnectorConfig connectorConfig, Answer>> answer) throws Exception { when(worker.isRunning(CONN1)).thenReturn(running); if (running) { when(worker.getPlugins()).thenReturn(plugins); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginsTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginsTest.java index 8a637beee5ec0..185d50e39ad55 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginsTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginsTest.java @@ -425,7 +425,7 @@ public void newHeaderConverterShouldConfigureWithPluginClassLoader() { @Test public void newConnectorShouldInstantiateWithPluginClassLoader() { - Connector plugin = plugins.newConnector(TestPlugin.SAMPLING_CONNECTOR.className()); + Connector plugin = plugins.newConnector(TestPlugin.SAMPLING_CONNECTOR.className()).delegate; assertInstanceOf(SamplingTestPlugin.class, plugin, "Cannot collect samples"); Map samples = ((SamplingTestPlugin) plugin).flatten(); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java index bed08ffa2ebf7..607c3630d4ee8 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java @@ -40,7 +40,11 @@ import org.apache.kafka.connect.runtime.Worker; import org.apache.kafka.connect.runtime.WorkerConfigTransformer; import org.apache.kafka.connect.runtime.distributed.SampleConnectorClientConfigOverridePolicy; +import org.apache.kafka.connect.runtime.isolation.IsolatedConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSinkConnector; +import org.apache.kafka.connect.runtime.isolation.IsolatedSourceConnector; import org.apache.kafka.connect.runtime.isolation.LoaderSwap; +import org.apache.kafka.connect.runtime.isolation.PluginType; import org.apache.kafka.connect.runtime.rest.entities.Message; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.isolation.PluginClassLoader; @@ -65,6 +69,7 @@ import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; import org.mockito.Mock; +import org.mockito.Mockito; import org.mockito.junit.MockitoJUnitRunner; import java.util.ArrayList; @@ -167,18 +172,19 @@ public void testCreateSourceConnector() throws Exception { } @Test - public void testCreateConnectorFailedValidation() { + public void testCreateConnectorFailedValidation() throws Exception { // Basic validation should be performed and return an error, but should still evaluate the connector's config - Map config = connectorConfig(SourceSink.SOURCE); config.remove(ConnectorConfig.NAME_CONFIG); - Connector connectorMock = mock(SourceConnector.class); + IsolatedConnector connectorMock = mock(IsolatedSourceConnector.class); + when(connectorMock.type()).thenReturn(PluginType.SOURCE); + Mockito.>when(connectorMock.pluginClass()).thenReturn(BogusSourceConnector.class); when(worker.configTransformer()).thenReturn(transformer); final ArgumentCaptor> configCapture = ArgumentCaptor.forClass(Map.class); when(transformer.transform(configCapture.capture())).thenAnswer(invocation -> configCapture.getValue()); when(worker.getPlugins()).thenReturn(plugins); - when(plugins.newConnector(anyString())).thenReturn(connectorMock); + Mockito.>when(plugins.newConnector(anyString())).thenReturn(connectorMock); when(plugins.connectorLoader(anyString())).thenReturn(pluginLoader); when(plugins.withClassLoader(pluginLoader)).thenReturn(loaderSwap); @@ -313,7 +319,6 @@ public void testRestartConnectorNewTaskConfigs() throws Exception { expectAdd(SourceSink.SOURCE); Map config = connectorConfig(SourceSink.SOURCE); - ConnectorTaskId taskId = new ConnectorTaskId(CONNECTOR_NAME, 0); expectConfigValidation(SourceSink.SOURCE, config); herder.putConnectorConfig(CONNECTOR_NAME, config, false, createCallback); @@ -335,6 +340,7 @@ public void testRestartConnectorNewTaskConfigs() throws Exception { FutureCallback restartCallback = new FutureCallback<>(); expectStop(); herder.restartConnector(CONNECTOR_NAME, restartCallback); + ConnectorTaskId taskId = new ConnectorTaskId(CONNECTOR_NAME, 0); verify(statusBackingStore).put(new TaskStatus(taskId, TaskStatus.State.DESTROYED, WORKER_ID, 0)); restartCallback.get(WAIT_TIME_MS, TimeUnit.MILLISECONDS); } @@ -371,7 +377,6 @@ public void testRestartTask() throws Exception { expectAdd(SourceSink.SOURCE); Map connectorConfig = connectorConfig(SourceSink.SOURCE); - expectConfigValidation(SourceSink.SOURCE, connectorConfig); doNothing().when(worker).stopAndAwaitTask(taskId); @@ -650,8 +655,6 @@ public void testCreateAndStop() throws Exception { @Test public void testAccessors() throws Exception { - Map connConfig = connectorConfig(SourceSink.SOURCE); - System.out.println(connConfig); Callback> listConnectorsCb = mock(Callback.class); Callback connectorInfoCb = mock(Callback.class); @@ -666,8 +669,9 @@ public void testAccessors() throws Exception { doNothing().when(tasksConfigCb).onCompletion(any(NotFoundException.class), isNull()); doNothing().when(connectorConfigCb).onCompletion(any(NotFoundException.class), isNull()); - + // Create connector expectAdd(SourceSink.SOURCE); + Map connConfig = connectorConfig(SourceSink.SOURCE); expectConfigValidation(SourceSink.SOURCE, connConfig); // Validate accessors with 1 connector @@ -712,7 +716,6 @@ public void testPutConnectorConfig() throws Exception { Callback> connectorConfigCb = mock(Callback.class); - expectAdd(SourceSink.SOURCE); expectConfigValidation(SourceSink.SOURCE, connConfig, newConnConfig); @@ -731,7 +734,6 @@ public void testPutConnectorConfig() throws Exception { when(worker.connectorTaskConfigs(CONNECTOR_NAME, new SourceConnectorConfig(plugins, newConnConfig, true))) .thenReturn(singletonList(taskConfig(SourceSink.SOURCE))); - herder.putConnectorConfig(CONNECTOR_NAME, connConfig, false, createCallback); Herder.Created connectorInfo = createCallback.get(WAIT_TIME_MS, TimeUnit.MILLISECONDS); assertEquals(createdInfo(SourceSink.SOURCE), connectorInfo.result()); @@ -760,12 +762,14 @@ public void testPutTaskConfigs() { } @Test - public void testCorruptConfig() { + public void testCorruptConfig() throws Exception { Map config = new HashMap<>(); config.put(ConnectorConfig.NAME_CONFIG, CONNECTOR_NAME); config.put(ConnectorConfig.CONNECTOR_CLASS_CONFIG, BogusSinkConnector.class.getName()); config.put(SinkConnectorConfig.TOPICS_CONFIG, TOPICS_LIST_STR); - Connector connectorMock = mock(SinkConnector.class); + IsolatedConnector connectorMock = mock(IsolatedSinkConnector.class); + when(connectorMock.type()).thenReturn(PluginType.SINK); + Mockito.>when(connectorMock.pluginClass()).thenReturn(BogusSinkConnector.class); String error = "This is an error in your config!"; List errors = new ArrayList<>(singletonList(error)); String key = "foo.invalid.key"; @@ -782,7 +786,7 @@ public void testCorruptConfig() { when(worker.getPlugins()).thenReturn(plugins); when(plugins.connectorLoader(anyString())).thenReturn(pluginLoader); when(plugins.withClassLoader(pluginLoader)).thenReturn(loaderSwap); - when(plugins.newConnector(anyString())).thenReturn(connectorMock); + Mockito.>when(plugins.newConnector(anyString())).thenReturn(connectorMock); when(connectorMock.config()).thenReturn(configDef); herder.putConnectorConfig(CONNECTOR_NAME, config, true, createCallback); @@ -942,11 +946,10 @@ public void testRequestTaskReconfigurationDoesNotDeadlock() throws Exception { // Start the connector Map config = connectorConfig(SourceSink.SOURCE); - // Prepare for connector and task config update Map newConfig = connectorConfig(SourceSink.SOURCE); newConfig.put("dummy-connector-property", "yes"); expectConfigValidation(SourceSink.SOURCE, config, newConfig); - mockStartConnector(newConfig, TargetState.STARTED, TargetState.STARTED, null); + herder.putConnectorConfig(CONNECTOR_NAME, config, false, createCallback); // Wait on connector to start @@ -984,7 +987,7 @@ public void testRequestTaskReconfigurationDoesNotDeadlock() throws Exception { verify(statusBackingStore, times(2)).put(new TaskStatus(new ConnectorTaskId(CONNECTOR_NAME, 0), TaskStatus.State.DESTROYED, WORKER_ID, 0)); } - private void expectAdd(SourceSink sourceSink) { + private void expectAdd(SourceSink sourceSink) throws Exception { Map connectorProps = connectorConfig(sourceSink); ConnectorConfig connConfig = sourceSink == SourceSink.SOURCE ? new SourceConnectorConfig(plugins, connectorProps, true) : @@ -1090,19 +1093,27 @@ private static Map taskConfig(SourceSink sourceSink) { private void expectConfigValidation( SourceSink sourceSink, Map... configs - ) { + ) throws Exception { + + IsolatedConnector connectorMock; + if (sourceSink == SourceSink.SOURCE) { + connectorMock = mock(IsolatedSourceConnector.class); + when(connectorMock.type()).thenReturn(PluginType.SOURCE); + Mockito.>when(connectorMock.pluginClass()).thenReturn(BogusSourceConnector.class); + } else { + connectorMock = mock(IsolatedSinkConnector.class); + when(connectorMock.type()).thenReturn(PluginType.SINK); + Mockito.>when(connectorMock.pluginClass()).thenReturn(BogusSinkConnector.class); + } // config validation - Connector connectorMock = sourceSink == SourceSink.SOURCE ? mock(SourceConnector.class) : mock(SinkConnector.class); when(worker.configTransformer()).thenReturn(transformer); final ArgumentCaptor> configCapture = ArgumentCaptor.forClass(Map.class); when(transformer.transform(configCapture.capture())).thenAnswer(invocation -> configCapture.getValue()); when(worker.getPlugins()).thenReturn(plugins); when(plugins.connectorLoader(anyString())).thenReturn(pluginLoader); when(plugins.withClassLoader(pluginLoader)).thenReturn(loaderSwap); - - // Assume the connector should always be created when(worker.getPlugins()).thenReturn(plugins); - when(plugins.newConnector(anyString())).thenReturn(connectorMock); + Mockito.>when(plugins.newConnector(anyString())).thenReturn(connectorMock); when(connectorMock.config()).thenReturn(new ConfigDef()); // Set up validation for each config