diff --git a/docs/docs/resources/pipeline-components/dependencies/defaults_pipeline_component_dependencies.yaml b/docs/docs/resources/pipeline-components/dependencies/defaults_pipeline_component_dependencies.yaml index 0e3be65f8..8b58f9f16 100644 --- a/docs/docs/resources/pipeline-components/dependencies/defaults_pipeline_component_dependencies.yaml +++ b/docs/docs/resources/pipeline-components/dependencies/defaults_pipeline_component_dependencies.yaml @@ -13,11 +13,9 @@ kafka-connector.yaml: - from_.yaml - to.yaml - config-kafka-connector.yaml -- resetter_values.yaml kafka-sink-connector.yaml: [] kafka-source-connector.yaml: - from_-kafka-source-connector.yaml -- offset_topic-kafka-source-connector.yaml kubernetes-app.yaml: - prefix.yaml - from_.yaml diff --git a/docs/docs/resources/pipeline-components/dependencies/kpops_structure.yaml b/docs/docs/resources/pipeline-components/dependencies/kpops_structure.yaml index 594bae28d..ed2cb7ded 100644 --- a/docs/docs/resources/pipeline-components/dependencies/kpops_structure.yaml +++ b/docs/docs/resources/pipeline-components/dependencies/kpops_structure.yaml @@ -37,8 +37,6 @@ kpops_components_fields: - to - config - state - - resetter_namespace - - resetter_values kafka-sink-connector: - name - enabled @@ -47,8 +45,6 @@ kpops_components_fields: - to - config - state - - resetter_namespace - - resetter_values kafka-source-connector: - name - enabled @@ -57,9 +53,6 @@ kpops_components_fields: - to - config - state - - resetter_namespace - - resetter_values - - offset_topic kubernetes-app: - name - enabled diff --git a/docs/docs/resources/pipeline-components/dependencies/pipeline_component_dependencies.yaml b/docs/docs/resources/pipeline-components/dependencies/pipeline_component_dependencies.yaml index c517d87e3..23969c9d0 100644 --- a/docs/docs/resources/pipeline-components/dependencies/pipeline_component_dependencies.yaml +++ b/docs/docs/resources/pipeline-components/dependencies/pipeline_component_dependencies.yaml @@ -23,20 +23,16 @@ kafka-connector.yaml: - from_.yaml - to.yaml - config-kafka-connector.yaml -- resetter_values.yaml kafka-sink-connector.yaml: - prefix.yaml - from_.yaml - to.yaml - config-kafka-connector.yaml -- resetter_values.yaml kafka-source-connector.yaml: - prefix.yaml - from_-kafka-source-connector.yaml - to.yaml - config-kafka-connector.yaml -- resetter_values.yaml -- offset_topic-kafka-source-connector.yaml kubernetes-app.yaml: - prefix.yaml - from_.yaml diff --git a/docs/docs/resources/pipeline-components/kafka-connector.yaml b/docs/docs/resources/pipeline-components/kafka-connector.yaml index 4a64cdf3a..dc809d45b 100644 --- a/docs/docs/resources/pipeline-components/kafka-connector.yaml +++ b/docs/docs/resources/pipeline-components/kafka-connector.yaml @@ -45,7 +45,3 @@ # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" diff --git a/docs/docs/resources/pipeline-components/kafka-sink-connector.yaml b/docs/docs/resources/pipeline-components/kafka-sink-connector.yaml index 9253db0e0..dae7ef1e4 100644 --- a/docs/docs/resources/pipeline-components/kafka-sink-connector.yaml +++ b/docs/docs/resources/pipeline-components/kafka-sink-connector.yaml @@ -46,7 +46,3 @@ # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" diff --git a/docs/docs/resources/pipeline-components/kafka-source-connector.yaml b/docs/docs/resources/pipeline-components/kafka-source-connector.yaml index 8fa6c51d5..a470311c9 100644 --- a/docs/docs/resources/pipeline-components/kafka-source-connector.yaml +++ b/docs/docs/resources/pipeline-components/kafka-source-connector.yaml @@ -27,10 +27,3 @@ # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" - # offset.storage.topic - # https://kafka.apache.org/documentation/#connect_running - offset_topic: offset_topic diff --git a/docs/docs/resources/pipeline-components/pipeline.yaml b/docs/docs/resources/pipeline-components/pipeline.yaml index b6ad9a726..ae3766a50 100644 --- a/docs/docs/resources/pipeline-components/pipeline.yaml +++ b/docs/docs/resources/pipeline-components/pipeline.yaml @@ -248,10 +248,6 @@ # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" # Kafka source connector - type: kafka-source-connector # required name: kafka-source-connector # required @@ -281,13 +277,6 @@ # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" - # offset.storage.topic - # https://kafka.apache.org/documentation/#connect_running - offset_topic: offset_topic # Base Kubernetes App - type: kubernetes-app name: kubernetes-app # required diff --git a/docs/docs/resources/pipeline-components/sections/offset_topic-kafka-source-connector.yaml b/docs/docs/resources/pipeline-components/sections/offset_topic-kafka-source-connector.yaml deleted file mode 100644 index bae3a033e..000000000 --- a/docs/docs/resources/pipeline-components/sections/offset_topic-kafka-source-connector.yaml +++ /dev/null @@ -1,3 +0,0 @@ - # offset.storage.topic - # https://kafka.apache.org/documentation/#connect_running - offset_topic: offset_topic diff --git a/docs/docs/resources/pipeline-components/sections/resetter_values.yaml b/docs/docs/resources/pipeline-components/sections/resetter_values.yaml deleted file mode 100644 index 376c472d3..000000000 --- a/docs/docs/resources/pipeline-components/sections/resetter_values.yaml +++ /dev/null @@ -1,4 +0,0 @@ - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" diff --git a/docs/docs/resources/pipeline-defaults/defaults-kafka-connector.yaml b/docs/docs/resources/pipeline-defaults/defaults-kafka-connector.yaml index 9ace5173b..dc06b447d 100644 --- a/docs/docs/resources/pipeline-defaults/defaults-kafka-connector.yaml +++ b/docs/docs/resources/pipeline-defaults/defaults-kafka-connector.yaml @@ -48,7 +48,3 @@ kafka-connector: # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" diff --git a/docs/docs/resources/pipeline-defaults/defaults-kafka-source-connector.yaml b/docs/docs/resources/pipeline-defaults/defaults-kafka-source-connector.yaml index 8d634f456..fcd203910 100644 --- a/docs/docs/resources/pipeline-defaults/defaults-kafka-source-connector.yaml +++ b/docs/docs/resources/pipeline-defaults/defaults-kafka-source-connector.yaml @@ -4,6 +4,3 @@ kafka-source-connector: # The source connector has no `from` section # from: - # offset.storage.topic - # https://kafka.apache.org/documentation/#connect_running - offset_topic: offset_topic diff --git a/docs/docs/resources/pipeline-defaults/defaults.yaml b/docs/docs/resources/pipeline-defaults/defaults.yaml index e1e16c56c..b5f776590 100644 --- a/docs/docs/resources/pipeline-defaults/defaults.yaml +++ b/docs/docs/resources/pipeline-defaults/defaults.yaml @@ -172,10 +172,6 @@ kafka-connector: # Full documentation on connectors: https://kafka.apache.org/documentation/#connectconfigs config: # required tasks.max: 1 - # Overriding Kafka Connect Resetter Helm values. E.g. to override the - # Image Tag etc. - resetter_values: - imageTag: "1.2.3" # Kafka sink connector # # Child of: KafkaConnector @@ -187,9 +183,6 @@ kafka-sink-connector: kafka-source-connector: # The source connector has no `from` section # from: - # offset.storage.topic - # https://kafka.apache.org/documentation/#connect_running - offset_topic: offset_topic # Base Kubernetes App # # Parent of: HelmApp diff --git a/docs/docs/schema/defaults.json b/docs/docs/schema/defaults.json index 55394a99d..690a0628f 100644 --- a/docs/docs/schema/defaults.json +++ b/docs/docs/schema/defaults.json @@ -1397,23 +1397,6 @@ "title": "Prefix", "type": "string" }, - "resetter_namespace": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "Kubernetes namespace in which the Kafka Connect resetter shall be deployed", - "title": "Resetter Namespace" - }, - "resetter_values": { - "$ref": "#/$defs/HelmAppValues", - "description": "Overriding Kafka Connect resetter Helm values, e.g. to override the image tag etc." - }, "state": { "anyOf": [ { @@ -1548,23 +1531,6 @@ "title": "Prefix", "type": "string" }, - "resetter_namespace": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "Kubernetes namespace in which the Kafka Connect resetter shall be deployed", - "title": "Resetter Namespace" - }, - "resetter_values": { - "$ref": "#/$defs/HelmAppValues", - "description": "Overriding Kafka Connect resetter Helm values, e.g. to override the image tag etc." - }, "state": { "anyOf": [ { @@ -1635,42 +1601,12 @@ "title": "Name", "type": "string" }, - "offset_topic": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "`offset.storage.topic`, more info: https://kafka.apache.org/documentation/#connect_running", - "title": "Offset Topic" - }, "prefix": { "default": "${pipeline.name}-", "description": "Pipeline prefix that will prefix every component name. If you wish to not have any prefix you can specify an empty string.", "title": "Prefix", "type": "string" }, - "resetter_namespace": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "Kubernetes namespace in which the Kafka Connect resetter shall be deployed", - "title": "Resetter Namespace" - }, - "resetter_values": { - "$ref": "#/$defs/HelmAppValues", - "description": "Overriding Kafka Connect resetter Helm values, e.g. to override the image tag etc." - }, "state": { "anyOf": [ { diff --git a/docs/docs/schema/pipeline.json b/docs/docs/schema/pipeline.json index 9f7f9a771..96f297feb 100644 --- a/docs/docs/schema/pipeline.json +++ b/docs/docs/schema/pipeline.json @@ -1353,23 +1353,6 @@ "title": "Prefix", "type": "string" }, - "resetter_namespace": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "Kubernetes namespace in which the Kafka Connect resetter shall be deployed", - "title": "Resetter Namespace" - }, - "resetter_values": { - "$ref": "#/$defs/HelmAppValues", - "description": "Overriding Kafka Connect resetter Helm values, e.g. to override the image tag etc." - }, "state": { "anyOf": [ { @@ -1440,42 +1423,12 @@ "title": "Name", "type": "string" }, - "offset_topic": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "`offset.storage.topic`, more info: https://kafka.apache.org/documentation/#connect_running", - "title": "Offset Topic" - }, "prefix": { "default": "${pipeline.name}-", "description": "Pipeline prefix that will prefix every component name. If you wish to not have any prefix you can specify an empty string.", "title": "Prefix", "type": "string" }, - "resetter_namespace": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "default": null, - "description": "Kubernetes namespace in which the Kafka Connect resetter shall be deployed", - "title": "Resetter Namespace" - }, - "resetter_values": { - "$ref": "#/$defs/HelmAppValues", - "description": "Overriding Kafka Connect resetter Helm values, e.g. to override the image tag etc." - }, "state": { "anyOf": [ { diff --git a/docs/docs/user/core-concepts/defaults.md b/docs/docs/user/core-concepts/defaults.md index d19c9d26d..5a8b92e93 100644 --- a/docs/docs/user/core-concepts/defaults.md +++ b/docs/docs/user/core-concepts/defaults.md @@ -10,8 +10,6 @@ KPOps has a very efficient way of dealing with repeating settings which manifest An important mechanic of KPOps is that `defaults` set for a component apply to all components that inherit from it. -It is possible, although not recommended, to add settings that are specific to a component's subclass. An example would be configuring `offset_topic` under `kafka-connector` instead of `kafka-source-connector`. - ### Configuration KPOps allows using multiple default values. The `defaults.yaml` (or `defaults_.yaml`) files can be distributed across multiple files. These will be picked up by KPOps and get merged into a single `pipeline.yaml` file. diff --git a/kpops/component_handlers/kafka_connect/connect_wrapper.py b/kpops/component_handlers/kafka_connect/connect_wrapper.py index 97755b37b..06703ada2 100644 --- a/kpops/component_handlers/kafka_connect/connect_wrapper.py +++ b/kpops/component_handlers/kafka_connect/connect_wrapper.py @@ -63,11 +63,11 @@ async def create_connector( response = await self._client.post( "/connectors", json=payload.model_dump(exclude_none=True) ) - if response.status_code == httpx.codes.CREATED: + if response.status_code == httpx.codes.CREATED.value: log.info(f"Connector {connector_config.name} created.") log.debug(response.json()) return ConnectorResponse.model_validate(response.json()) - if response.status_code == httpx.codes.CONFLICT: + if response.status_code == httpx.codes.CONFLICT.value: log.warning( "Rebalancing in progress while creating a connector... Retrying..." ) @@ -83,12 +83,12 @@ async def get_connector(self, connector_name: str) -> ConnectorResponse: :return: Information about the connector. """ response = await self._client.get(f"/connectors/{connector_name}") - if response.status_code == httpx.codes.OK: + if response.is_success: log.debug(response.json()) return ConnectorResponse.model_validate(response.json()) - if response.status_code == httpx.codes.NOT_FOUND: + if response.status_code == httpx.codes.NOT_FOUND.value: raise ConnectorNotFoundException - if response.status_code == httpx.codes.CONFLICT: + if response.status_code == httpx.codes.CONFLICT.value: log.warning( "Rebalancing in progress while getting a connector... Retrying..." ) @@ -106,10 +106,10 @@ async def get_connector_status( :return: Status of the connector. """ response = await self._client.get(f"/connectors/{connector_name}/status") - if response.status_code == httpx.codes.OK: + if response.is_success: log.debug(response.json()) return ConnectorStatusResponse.model_validate(response.json()) - if response.status_code == httpx.codes.NOT_FOUND: + if response.status_code == httpx.codes.NOT_FOUND.value: raise ConnectorNotFoundException raise KafkaConnectError(response) @@ -120,9 +120,10 @@ async def pause_connector(self, connector_name: str) -> None: :param connector_name: Name of the connector """ response = await self._client.put(f"/connectors/{connector_name}/pause") - if response.status_code != httpx.codes.ACCEPTED: - raise KafkaConnectError(response) - log.info(f"Connector {connector_name} paused.") + if response.is_success: + log.info(f"Connector {connector_name} paused.") + return + raise KafkaConnectError(response) async def resume_connector(self, connector_name: str) -> None: """Resume connector. @@ -131,9 +132,10 @@ async def resume_connector(self, connector_name: str) -> None: :param connector_name: Name of the connector """ response = await self._client.put(f"/connectors/{connector_name}/resume") - if response.status_code != httpx.codes.ACCEPTED: - raise KafkaConnectError(response) - log.info(f"Connector {connector_name} resumed.") + if response.is_success: + log.info(f"Connector {connector_name} resumed.") + return + raise KafkaConnectError(response) async def stop_connector(self, connector_name: str) -> None: """Stop connector. @@ -142,9 +144,10 @@ async def stop_connector(self, connector_name: str) -> None: :param connector_name: Name of the connector """ response = await self._client.put(f"/connectors/{connector_name}/stop") - if response.status_code != httpx.codes.NO_CONTENT: - raise KafkaConnectError(response) - log.info(f"Connector {connector_name} stopped.") + if response.is_success: + log.info(f"Connector {connector_name} stopped.") + return + raise KafkaConnectError(response) async def update_connector_config( self, connector_config: KafkaConnectorConfig @@ -165,15 +168,15 @@ async def update_connector_config( ) data: dict[str, Any] = response.json() - if response.status_code == httpx.codes.OK: + if response.status_code == httpx.codes.OK.value: log.info(f"Config for connector {connector_name} updated.") log.debug(data) return ConnectorResponse.model_validate(data) - if response.status_code == httpx.codes.CREATED: + if response.status_code == httpx.codes.CREATED.value: log.info(f"Connector {connector_name} created.") log.debug(data) return ConnectorResponse.model_validate(data) - if response.status_code == httpx.codes.CONFLICT: + if response.status_code == httpx.codes.CONFLICT.value: log.warning( "Rebalancing in progress while updating a connector... Retrying..." ) @@ -195,7 +198,7 @@ async def validate_connector_config( json=connector_config.model_dump(), ) - if response.status_code == httpx.codes.OK: + if response.status_code == httpx.codes.OK.value: kafka_connect_error_response = KafkaConnectConfigErrorResponse( **response.json() ) @@ -220,15 +223,31 @@ async def delete_connector(self, connector_name: str) -> None: :raises ConnectorNotFoundException: Connector not found """ response = await self._client.delete(f"/connectors/{connector_name}") - if response.status_code == httpx.codes.NO_CONTENT: + if response.status_code == httpx.codes.NO_CONTENT.value: log.info(f"Connector {connector_name} deleted.") return None - if response.status_code == httpx.codes.NOT_FOUND: + if response.status_code == httpx.codes.NOT_FOUND.value: raise ConnectorNotFoundException - if response.status_code == httpx.codes.CONFLICT: + if response.status_code == httpx.codes.CONFLICT.value: log.warning( "Rebalancing in progress while deleting a connector... Retrying..." ) await asyncio.sleep(1) return await self.delete_connector(connector_name) raise KafkaConnectError(response) + + async def reset_offset(self, connector_name: str) -> None: + """Reset the offsets for a connector; the connector must exist, and must be in the STOPPED state. + + API Reference: + https://docs.confluent.io/platform/current/connect/references/restapi.html#delete--connectors-connector-offsets + :param connector_name: Configuration parameters for the connector. + :raises ConnectorNotFoundException: Connector not found + """ + response = await self._client.delete(f"/connectors/{connector_name}/offsets") + if response.is_success: + log.info(f"Connector {connector_name} offsets reset.") + return + if response.status_code == httpx.codes.NOT_FOUND.value: + raise ConnectorNotFoundException + raise KafkaConnectError(response) diff --git a/kpops/component_handlers/kafka_connect/kafka_connect_handler.py b/kpops/component_handlers/kafka_connect/kafka_connect_handler.py index d56081724..28740736b 100644 --- a/kpops/component_handlers/kafka_connect/kafka_connect_handler.py +++ b/kpops/component_handlers/kafka_connect/kafka_connect_handler.py @@ -87,6 +87,36 @@ async def destroy_connector(self, connector_name: str, *, dry_run: bool) -> None f"Connector Destruction: the connector {connector_name} does not exist. Skipping." ) + async def reset_connector( + self, connector_config: KafkaConnectorConfig, *, dry_run: bool + ) -> None: + """Reset connector offsets. + + If the connector does not exist, it is created temporarily in a + paused state so its offsets can be reset. + + :param connector_config: The connector config. + :param dry_run: Whether the connector reset should be run in dry run mode. + """ + connector_name = connector_config.name + try: + await self._connect_wrapper.get_connector(connector_name) + except ConnectorNotFoundException: + await self.create_connector( + connector_config, state=ConnectorNewState.PAUSED, dry_run=dry_run + ) + + if dry_run: + await self.__dry_run_connector_reset(connector_name) + else: + try: + await self._connect_wrapper.stop_connector(connector_name) + await self._connect_wrapper.reset_offset(connector_name) + except ConnectorNotFoundException: + log.warning( + f"Connector reset: the connector {connector_name} does not exist. Skipping." + ) + async def __dry_run_connector_creation( self, connector_config: KafkaConnectorConfig, @@ -142,6 +172,17 @@ async def __dry_run_connector_creation( f"Connector Creation: connector config for {connector_name} is valid!" ) + async def __dry_run_connector_reset(self, connector_name: str) -> None: + log.info( + magentaify( + f"Connector reset: resetting offsets for connector {connector_name}." + ) + ) + log.debug(f"PUT /connectors/{connector_name}/stop HTTP/1.1") + log.debug(f"HOST: {self._connect_wrapper.url}") + log.debug(f"DELETE /connectors/{connector_name}/offsets HTTP/1.1") + log.debug(f"HOST: {self._connect_wrapper.url}") + async def __dry_run_connector_deletion(self, connector_name: str) -> None: try: await self._connect_wrapper.get_connector(connector_name) diff --git a/kpops/components/base_components/kafka_connector.py b/kpops/components/base_components/kafka_connector.py index fd68ffae2..fdc60780c 100644 --- a/kpops/components/base_components/kafka_connector.py +++ b/kpops/components/base_components/kafka_connector.py @@ -2,101 +2,24 @@ import logging from abc import ABC -from functools import cached_property -from typing import Any, Literal, NoReturn, Self +from typing import Any, NoReturn -import pydantic -from pydantic import Field, PrivateAttr, ValidationInfo, field_validator +from pydantic import PrivateAttr, ValidationInfo, field_validator from typing_extensions import override from kpops.component_handlers import get_handlers -from kpops.component_handlers.helm_wrapper.model import ( - HelmRepoConfig, -) from kpops.component_handlers.kafka_connect.model import ( ConnectorNewState, KafkaConnectorConfig, KafkaConnectorType, ) -from kpops.components.base_components.cleaner import Cleaner -from kpops.components.base_components.helm_app import HelmAppValues from kpops.components.base_components.models.from_section import FromTopic from kpops.components.base_components.pipeline_component import PipelineComponent from kpops.components.common.topic import KafkaTopic -from kpops.config import get_config -from kpops.utils.colorify import magentaify -from kpops.utils.pydantic import CamelCaseConfigModel, SkipGenerate log = logging.getLogger("KafkaConnector") -class KafkaConnectorResetterConfig(CamelCaseConfigModel): - brokers: str - connector: str - delete_consumer_group: bool | None = None - offset_topic: str | None = None - - -class KafkaConnectorResetterValues(HelmAppValues): - connector_type: Literal["source", "sink"] - config: KafkaConnectorResetterConfig - - -class KafkaConnectorResetter(Cleaner, ABC): - """Helm app for resetting and cleaning a Kafka Connector. - - :param repo_config: Configuration of the Helm chart repo to be used for - deploying the component, defaults to kafka-connect-resetter Helm repo - :param version: Helm chart version, defaults to "1.0.4" - """ - - from_: None = None # pyright: ignore[reportIncompatibleVariableOverride] - to: None = None # pyright: ignore[reportIncompatibleVariableOverride] - values: KafkaConnectorResetterValues # pyright: ignore[reportIncompatibleVariableOverride] - repo_config: SkipGenerate[HelmRepoConfig] = HelmRepoConfig( # pyright: ignore[reportIncompatibleVariableOverride] - repository_name="bakdata-kafka-connect-resetter", - url="https://bakdata.github.io/kafka-connect-resetter/", - ) - version: str | None = "1.0.4" - - @property - @override - def helm_chart(self) -> str: - return f"{self.repo_config.repository_name}/kafka-connect-resetter" - - @override - async def reset(self, dry_run: bool) -> None: - """Reset connector. - - At first, it deletes the previous cleanup job (connector resetter) - to make sure that there is no running clean job in the cluster. Then it releases a cleanup job. - If retain_clean_jobs config is set to false the cleanup job will be deleted subsequently. - - :param dry_run: If the cleanup should be run in dry run mode or not - """ - log.info( - magentaify( - f"Connector Cleanup: uninstalling cleanup job Helm release from previous runs for {self.values.config.connector}" - ) - ) - await self.destroy(dry_run) - - log.info( - magentaify( - f"Connector Cleanup: deploy Connect {self.values.connector_type} resetter for {self.values.config.connector}" - ) - ) - await self.deploy(dry_run) - - if not get_config().retain_clean_jobs: - log.info(magentaify("Connector Cleanup: uninstall Kafka Resetter.")) - await self.destroy(dry_run) - - @override - async def clean(self, dry_run: bool) -> None: - await self.reset(dry_run) - - class KafkaConnector(PipelineComponent, ABC): """Base class for all Kafka connectors. @@ -104,15 +27,10 @@ class KafkaConnector(PipelineComponent, ABC): :param config: Connector config :param state: Connector state - :param resetter_namespace: Kubernetes namespace in which the Kafka Connect resetter shall be deployed - :param resetter_values: Overriding Kafka Connect resetter Helm values, e.g. to override the image tag etc., - defaults to empty HelmAppValues """ config: KafkaConnectorConfig state: ConnectorNewState | None = None - resetter_namespace: str | None = None - resetter_values: HelmAppValues = Field(default_factory=HelmAppValues) _connector_type: KafkaConnectorType = PrivateAttr() @field_validator("config", mode="before") @@ -132,35 +50,6 @@ def connector_config_should_have_component_name( config["name"] = component_name return KafkaConnectorConfig.model_validate(config) - @cached_property - def _resetter(self) -> KafkaConnectorResetter: - kwargs: dict[str, Any] = {} - if self.resetter_namespace: - kwargs["namespace"] = self.resetter_namespace - return KafkaConnectorResetter( - **kwargs, - **self.model_dump( - by_alias=True, - exclude={ - "_resetter", - "resetter_values", - "resetter_namespace", - "values", - "config", - "from_", - "to", - }, - ), - values=KafkaConnectorResetterValues( - connector_type=self._connector_type.value, - config=KafkaConnectorResetterConfig( - connector=self.full_name, - brokers=get_config().kafka_brokers, - ), - **self.resetter_values.model_dump(), - ), - ) - @override async def deploy(self, dry_run: bool) -> None: """Deploy Kafka Connector (Source/Sink). Create output topics and register schemas if configured.""" @@ -177,15 +66,23 @@ async def deploy(self, dry_run: bool) -> None: @override async def destroy(self, dry_run: bool) -> None: - """Delete Kafka Connector (Source/Sink) from the Kafka connect cluster.""" + """Delete connector.""" await get_handlers().connector_handler.destroy_connector( self.full_name, dry_run=dry_run ) + @override + async def reset(self, dry_run: bool) -> None: + """Reset connector offsets. Delete connector afterwards.""" + await get_handlers().connector_handler.reset_connector( + self.config, dry_run=dry_run + ) + await super().reset(dry_run) + @override async def clean(self, dry_run: bool) -> None: """Delete Kafka Connector. If schema handler is enabled, then remove schemas. Delete all the output topics.""" - await super().clean(dry_run) + await self.reset(dry_run) if self.to: if schema_handler := get_handlers().schema_handler: await schema_handler.delete_schemas(to_section=self.to, dry_run=dry_run) @@ -194,40 +91,15 @@ async def clean(self, dry_run: bool) -> None: class KafkaSourceConnector(KafkaConnector): - """Kafka source connector model. - - :param offset_topic: `offset.storage.topic`, - more info: https://kafka.apache.org/documentation/#connect_running, - defaults to None - """ - - offset_topic: str | None = None + """Kafka source connector model.""" _connector_type: KafkaConnectorType = PrivateAttr(KafkaConnectorType.SOURCE) - @pydantic.model_validator(mode="after") - def populate_offset_topic(self) -> Self: - if self.offset_topic: - self._resetter.values.config.offset_topic = self.offset_topic - return self - @override def apply_from_inputs(self, name: str, topic: FromTopic) -> NoReturn: msg = "Kafka source connector doesn't support FromSection" raise NotImplementedError(msg) - @override - async def reset(self, dry_run: bool) -> None: - """Reset state. Keep connector.""" - await super().reset(dry_run) - await self._resetter.reset(dry_run) - - @override - async def clean(self, dry_run: bool) -> None: - """Delete connector and reset state.""" - await super().clean(dry_run) - await self._resetter.clean(dry_run) - class KafkaSinkConnector(KafkaConnector): """Kafka sink connector model.""" @@ -250,17 +122,3 @@ def set_input_pattern(self, name: str) -> None: @override def set_error_topic(self, topic: KafkaTopic) -> None: self.config.errors_deadletterqueue_topic_name = topic - - @override - async def reset(self, dry_run: bool) -> None: - """Reset state. Keep consumer group and connector.""" - await super().reset(dry_run) - self._resetter.values.config.delete_consumer_group = False - await self._resetter.reset(dry_run) - - @override - async def clean(self, dry_run: bool) -> None: - """Delete connector and consumer group.""" - await super().clean(dry_run) - self._resetter.values.config.delete_consumer_group = True - await self._resetter.clean(dry_run) diff --git a/tests/component_handlers/kafka_connect/test_connect_handler.py b/tests/component_handlers/kafka_connect/test_connect_handler.py index 0ab445db9..874c86234 100644 --- a/tests/component_handlers/kafka_connect/test_connect_handler.py +++ b/tests/component_handlers/kafka_connect/test_connect_handler.py @@ -460,3 +460,105 @@ async def test_print_correct_warning_log_when_destroying_connector_and_connector log_warning_mock.assert_called_once_with( f"Connector Destruction: the connector {CONNECTOR_NAME} does not exist. Skipping." ) + + async def test_reset_connector_dry_run( + self, + connect_wrapper: AsyncMock, + handler: KafkaConnectHandler, + connector_config: KafkaConnectorConfig, + log_info_mock: MagicMock, + ) -> None: + await handler.reset_connector(connector_config, dry_run=True) + + connect_wrapper.get_connector.assert_called_once_with(CONNECTOR_NAME) + connect_wrapper.create_connector.assert_not_called() + connect_wrapper.stop_connector.assert_not_called() + connect_wrapper.reset_offset.assert_not_called() + log_info_mock.assert_called_once_with( + magentaify( + f"Connector reset: resetting offsets for connector {CONNECTOR_NAME}." + ) + ) + + async def test_reset_connector_dry_run_connector_missing( + self, + connect_wrapper: AsyncMock, + handler: KafkaConnectHandler, + connector_config: KafkaConnectorConfig, + renderer_diff_mock: MagicMock, + log_info_mock: MagicMock, + ) -> None: + connect_wrapper.get_connector.side_effect = ConnectorNotFoundException() + renderer_diff_mock.return_value = "" + + await handler.reset_connector(connector_config, dry_run=True) + + assert connect_wrapper.get_connector.call_count == 2 + connect_wrapper.validate_connector_config.assert_called_once_with( + connector_config + ) + connect_wrapper.create_connector.assert_not_called() + connect_wrapper.stop_connector.assert_not_called() + connect_wrapper.reset_offset.assert_not_called() + + assert log_info_mock.mock_calls == [ + mock.call( + f"Connector Creation: connector {CONNECTOR_NAME} does not exist. Creating connector in paused state with config:\n" + ), + mock.call( + f"Connector Creation: connector config for {CONNECTOR_NAME} is valid!" + ), + mock.call( + magentaify( + f"Connector reset: resetting offsets for connector {CONNECTOR_NAME}." + ) + ), + ] + + async def test_reset_connector( + self, + connect_wrapper: AsyncMock, + handler: KafkaConnectHandler, + connector_config: KafkaConnectorConfig, + ) -> None: + await handler.reset_connector(connector_config, dry_run=False) + assert connect_wrapper.mock_calls == [ + mock.call.get_connector(CONNECTOR_NAME), + mock.call.stop_connector(CONNECTOR_NAME), + mock.call.reset_offset(CONNECTOR_NAME), + ] + + async def test_reset_connector_when_connector_missing( + self, + connect_wrapper: AsyncMock, + handler: KafkaConnectHandler, + connector_config: KafkaConnectorConfig, + ) -> None: + connect_wrapper.get_connector.side_effect = ConnectorNotFoundException() + + await handler.reset_connector(connector_config, dry_run=False) + + assert connect_wrapper.mock_calls == [ + mock.call.get_connector(CONNECTOR_NAME), + mock.call.get_connector(CONNECTOR_NAME), + mock.call.create_connector(connector_config, ConnectorNewState.PAUSED), + mock.call.stop_connector(CONNECTOR_NAME), + mock.call.reset_offset(CONNECTOR_NAME), + ] + + async def test_reset_connector_disappears_during_reset( + self, + connect_wrapper: AsyncMock, + handler: KafkaConnectHandler, + connector_config: KafkaConnectorConfig, + log_warning_mock: MagicMock, + ) -> None: + """Connector existed but got deleted concurrently before it could be stopped.""" + connect_wrapper.stop_connector.side_effect = ConnectorNotFoundException() + + await handler.reset_connector(connector_config, dry_run=False) + + connect_wrapper.reset_offset.assert_not_called() + log_warning_mock.assert_called_once_with( + f"Connector reset: the connector {CONNECTOR_NAME} does not exist. Skipping." + ) diff --git a/tests/component_handlers/kafka_connect/test_connect_wrapper.py b/tests/component_handlers/kafka_connect/test_connect_wrapper.py index 7d67df8c8..e8941e471 100644 --- a/tests/component_handlers/kafka_connect/test_connect_wrapper.py +++ b/tests/component_handlers/kafka_connect/test_connect_wrapper.py @@ -301,7 +301,8 @@ async def test_pause_connector( url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}/pause", status_code=httpx.codes.ACCEPTED, ) - await connect_wrapper.pause_connector(CONNECTOR_NAME) + with caplog.at_level(logging.INFO): + await connect_wrapper.pause_connector(CONNECTOR_NAME) assert len(caplog.records) == 1 assert caplog.records[0].message == f"Connector {CONNECTOR_NAME} paused." assert caplog.records[0].levelname == "INFO" @@ -328,7 +329,8 @@ async def test_resume_connector( url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}/resume", status_code=httpx.codes.ACCEPTED, ) - await connect_wrapper.resume_connector(CONNECTOR_NAME) + with caplog.at_level(logging.INFO): + await connect_wrapper.resume_connector(CONNECTOR_NAME) assert len(caplog.records) == 1 assert caplog.records[0].message == f"Connector {CONNECTOR_NAME} resumed." assert caplog.records[0].levelname == "INFO" @@ -355,7 +357,8 @@ async def test_stop_connector( url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}/stop", status_code=httpx.codes.NO_CONTENT, ) - await connect_wrapper.stop_connector(CONNECTOR_NAME) + with caplog.at_level(logging.INFO): + await connect_wrapper.stop_connector(CONNECTOR_NAME) assert len(caplog.records) == 1 assert caplog.records[0].message == f"Connector {CONNECTOR_NAME} stopped." assert caplog.records[0].levelname == "INFO" @@ -485,7 +488,8 @@ async def test_delete_connector( url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}", status_code=httpx.codes.NO_CONTENT, ) - await connect_wrapper.delete_connector(CONNECTOR_NAME) + with caplog.at_level(logging.INFO): + await connect_wrapper.delete_connector(CONNECTOR_NAME) assert len(caplog.records) == 1 assert caplog.records[0].message == f"Connector {CONNECTOR_NAME} deleted." assert caplog.records[0].levelname == "INFO" @@ -541,6 +545,49 @@ async def test_delete_connector_retry( assert caplog.records[1].message == "Connector test-connector deleted." assert caplog.records[1].levelname == "INFO" + async def test_reset_offset( + self, + connect_wrapper: ConnectWrapper, + httpx_mock: HTTPXMock, + caplog: pytest.LogCaptureFixture, + ) -> None: + httpx_mock.add_response( + method="DELETE", + url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}/offsets", + status_code=httpx.codes.NO_CONTENT, + ) + with caplog.at_level(logging.INFO): + await connect_wrapper.reset_offset(CONNECTOR_NAME) + assert len(caplog.records) == 1 + assert caplog.records[0].message == f"Connector {CONNECTOR_NAME} offsets reset." + assert caplog.records[0].levelname == "INFO" + + async def test_reset_offset_not_found( + self, + connect_wrapper: ConnectWrapper, + httpx_mock: HTTPXMock, + ) -> None: + httpx_mock.add_response( + method="DELETE", + url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}/offsets", + headers=HEADERS, + status_code=httpx.codes.NOT_FOUND, + json={}, + ) + with pytest.raises(ConnectorNotFoundException): + await connect_wrapper.reset_offset(CONNECTOR_NAME) + + async def test_reset_offset_error( + self, connect_wrapper: ConnectWrapper, httpx_mock: HTTPXMock + ) -> None: + httpx_mock.add_response( + method="DELETE", + url=f"{DEFAULT_HOST}/connectors/{CONNECTOR_NAME}/offsets", + status_code=httpx.codes.INTERNAL_SERVER_ERROR, + ) + with pytest.raises(KafkaConnectError): + await connect_wrapper.reset_offset(CONNECTOR_NAME) + @pytest.fixture() def file_stream_connector_config(self) -> KafkaConnectorConfig: return KafkaConnectorConfig.model_validate( diff --git a/tests/components/test_kafka_connector.py b/tests/components/test_kafka_connector.py index 86d005206..c7fdcf903 100644 --- a/tests/components/test_kafka_connector.py +++ b/tests/components/test_kafka_connector.py @@ -19,7 +19,6 @@ "${pipeline.name}-" + "test-connector-with-long-612f3-clean" ) CONNECTOR_CLASS = "com.bakdata.connect.TestConnector" -RESETTER_NAMESPACE = "test-namespace" @pytest.mark.usefixtures("mock_env") @@ -49,7 +48,6 @@ def connector(self, connector_config: KafkaConnectorConfig) -> KafkaConnector: return KafkaConnector( # HACK: not supposed to be instantiated, because ABC name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, ) def test_connector_config_name_override(self, connector: KafkaConnector) -> None: @@ -58,7 +56,6 @@ def test_connector_config_name_override(self, connector: KafkaConnector) -> None connector = KafkaConnector( name=CONNECTOR_NAME, config={"connector.class": CONNECTOR_CLASS}, # pyright: ignore[reportArgumentType], gets enriched - resetter_namespace=RESETTER_NAMESPACE, ) assert connector.config.name == CONNECTOR_FULL_NAME diff --git a/tests/components/test_kafka_sink_connector.py b/tests/components/test_kafka_sink_connector.py index 8663614fe..5b7d78aae 100644 --- a/tests/components/test_kafka_sink_connector.py +++ b/tests/components/test_kafka_sink_connector.py @@ -1,21 +1,16 @@ -from unittest.mock import ANY, MagicMock, call +from unittest.mock import MagicMock import pytest from pytest_mock import MockerFixture from typing_extensions import override from kpops.component_handlers import get_handlers -from kpops.component_handlers.helm_wrapper.model import ( - HelmUpgradeInstallFlags, - RepoAuthFlags, -) from kpops.component_handlers.kafka_connect.model import ( ConnectorNewState, KafkaConnectorConfig, KafkaConnectorType, ) from kpops.components.base_components.kafka_connector import ( - KafkaConnectorResetter, KafkaSinkConnector, ) from kpops.components.base_components.models import TopicName @@ -32,14 +27,9 @@ OutputTopicTypes, TopicConfig, ) -from kpops.utils.colorify import magentaify from tests.components.test_kafka_connector import ( - CONNECTOR_CLEAN_FULL_NAME, - CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - CONNECTOR_CLEAN_RELEASE_NAME, CONNECTOR_FULL_NAME, CONNECTOR_NAME, - RESETTER_NAMESPACE, TestKafkaConnector, ) @@ -57,7 +47,6 @@ def connector(self, connector_config: KafkaConnectorConfig) -> KafkaSinkConnecto return KafkaSinkConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, to=ToSection( topics={ TopicName("${output_topic_name}"): TopicConfig( @@ -67,41 +56,6 @@ def connector(self, connector_config: KafkaConnectorConfig) -> KafkaSinkConnecto ), ) - def test_resetter(self, connector: KafkaSinkConnector) -> None: - resetter = connector._resetter - assert isinstance(resetter, KafkaConnectorResetter) - assert resetter.full_name == CONNECTOR_CLEAN_FULL_NAME - - def test_resetter_release_name(self, connector: KafkaSinkConnector) -> None: - assert connector.config.name == CONNECTOR_FULL_NAME - assert connector._resetter.helm_release_name == CONNECTOR_CLEAN_RELEASE_NAME - - def test_resetter_helm_name_override(self, connector: KafkaSinkConnector) -> None: - assert ( - connector._resetter.to_helm_values()["nameOverride"] - == CONNECTOR_CLEAN_HELM_NAMEOVERRIDE - ) - assert ( - connector._resetter.to_helm_values()["fullnameOverride"] - == CONNECTOR_CLEAN_HELM_NAMEOVERRIDE - ) - - def test_resetter_inheritance(self, connector: KafkaSinkConnector) -> None: - setattr(connector.resetter_values, "testKey", "foo") - resetter = connector._resetter - assert resetter - assert not hasattr(resetter, "_resetter") - - assert not hasattr(resetter, "resetter_namespace") - assert resetter.namespace == connector.resetter_namespace - - assert not hasattr(resetter, "resetter_values") - # check that resetter values are contained in resetter app values - assert ( - connector.resetter_values.model_dump().items() - <= resetter.values.model_dump().items() - ) - def test_connector_config_parsing( self, connector_config: KafkaConnectorConfig ) -> None: @@ -114,7 +68,6 @@ def test_connector_config_parsing( "topics.regex": topic_pattern, } ), - resetter_namespace=RESETTER_NAMESPACE, ) assert connector.config.topics_regex == topic_pattern assert connector.config.model_dump()["topics.regex"] == topic_pattern @@ -127,7 +80,6 @@ def test_from_section_parsing_input_topic( connector = KafkaSinkConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, from_=FromSection( topics={ topic1: FromTopic(type=InputTopicTypes.INPUT), @@ -162,7 +114,6 @@ def test_from_section_parsing_input_pattern( connector = KafkaSinkConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, from_=FromSection( topics={topic_pattern: FromTopic(type=InputTopicTypes.PATTERN)} ), @@ -241,88 +192,37 @@ async def test_reset_when_dry_run_is_true( mocker: MockerFixture, ) -> None: mock_destroy = mocker.patch.object(connector, "destroy") - mock_resetter_reset = mocker.spy(connector._resetter, "reset") + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" + ) dry_run = True await connector.reset(dry_run=dry_run) mock_destroy.assert_called_once_with(dry_run) - mock_resetter_reset.assert_called_once_with(dry_run) - dry_run_handler_mock.print_helm_diff.assert_called_once() + dry_run_handler_mock.print_helm_diff.assert_not_called() + mock_reset_connector.assert_called_once_with(connector.config, dry_run=dry_run) async def test_reset_when_dry_run_is_false( self, connector: KafkaSinkConnector, dry_run_handler_mock: MagicMock, - helm_mock: MagicMock, mocker: MockerFixture, ) -> None: mock_destroy = mocker.patch.object(connector, "destroy") mock_delete_topic = mocker.patch.object( get_handlers().topic_handler, "delete_topic" ) - mock_clean_connector = mocker.patch.object( - get_handlers().connector_handler, "clean_connector" + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" ) - mock_resetter_reset = mocker.spy(connector._resetter, "reset") - - mock = mocker.MagicMock() - mock.attach_mock(mock_destroy, "destroy_connector") - mock.attach_mock(mock_clean_connector, "mock_clean_connector") - mock.attach_mock(helm_mock, "helm") dry_run = False - await connector.reset(dry_run=dry_run) - mock_resetter_reset.assert_called_once_with(dry_run) - - mock.assert_has_calls( - [ - mocker.call.destroy_connector(dry_run), - mocker.call.helm.add_repo( - "bakdata-kafka-connect-resetter", - "https://bakdata.github.io/kafka-connect-resetter/", - RepoAuthFlags(), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - mocker.call.helm.upgrade_install( - CONNECTOR_CLEAN_RELEASE_NAME, - "bakdata-kafka-connect-resetter/kafka-connect-resetter", - dry_run, - RESETTER_NAMESPACE, - { - "nameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "fullnameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "connectorType": CONNECTOR_TYPE, - "config": { - "brokers": "broker:9092", - "connector": CONNECTOR_FULL_NAME, - "deleteConsumerGroup": False, - }, - }, - HelmUpgradeInstallFlags( - version="1.0.4", - wait=True, - wait_for_jobs=True, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - ] - ) - dry_run_handler_mock.print_helm_diff.assert_not_called() + mock_reset_connector.assert_called_once_with(connector.config, dry_run=dry_run) + mock_destroy.assert_called_once_with(dry_run) mock_delete_topic.assert_not_called() + dry_run_handler_mock.print_helm_diff.assert_not_called() async def test_clean_when_dry_run_is_true( self, @@ -332,95 +232,38 @@ async def test_clean_when_dry_run_is_true( dry_run = True await connector.clean(dry_run=dry_run) - dry_run_handler_mock.print_helm_diff.assert_called_once() + dry_run_handler_mock.print_helm_diff.assert_not_called() async def test_clean_when_dry_run_is_false( self, connector: KafkaSinkConnector, - helm_mock: MagicMock, - log_info_mock: MagicMock, dry_run_handler_mock: MagicMock, mocker: MockerFixture, ) -> None: mock_destroy = mocker.patch.object(connector, "destroy") - mock_delete_topic = mocker.patch.object( get_handlers().topic_handler, "delete_topic" ) - mock_clean_connector = mocker.patch.object( - get_handlers().connector_handler, "clean_connector" + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" ) mock = mocker.MagicMock() + mock.attach_mock(mock_reset_connector, "mock_reset_connector") mock.attach_mock(mock_destroy, "destroy_connector") mock.attach_mock(mock_delete_topic, "mock_delete_topic") - mock.attach_mock(mock_clean_connector, "mock_clean_connector") - mock.attach_mock(helm_mock, "helm") dry_run = False await connector.clean(dry_run=dry_run) - assert log_info_mock.mock_calls == [ - call.log_info( - magentaify( - f"Connector Cleanup: uninstalling cleanup job Helm release from previous runs for {CONNECTOR_FULL_NAME}" - ) - ), - call.log_info( - magentaify( - f"Connector Cleanup: deploy Connect {KafkaConnectorType.SINK.value} resetter for {CONNECTOR_FULL_NAME}" - ) - ), - call.log_info(magentaify("Connector Cleanup: uninstall Kafka Resetter.")), - ] - assert connector.to assert mock.mock_calls == [ + mocker.call.mock_reset_connector(connector.config, dry_run=dry_run), mocker.call.destroy_connector(dry_run), *( mocker.call.mock_delete_topic(topic, dry_run=dry_run) for topic in connector.to.kafka_topics ), - mocker.call.helm.add_repo( - "bakdata-kafka-connect-resetter", - "https://bakdata.github.io/kafka-connect-resetter/", - RepoAuthFlags(), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - mocker.call.helm.upgrade_install( - CONNECTOR_CLEAN_RELEASE_NAME, - "bakdata-kafka-connect-resetter/kafka-connect-resetter", - dry_run, - RESETTER_NAMESPACE, - { - "nameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "fullnameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "connectorType": CONNECTOR_TYPE, - "config": { - "brokers": "broker:9092", - "connector": CONNECTOR_FULL_NAME, - "deleteConsumerGroup": True, - }, - }, - HelmUpgradeInstallFlags( - version="1.0.4", - wait=True, - wait_for_jobs=True, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ ] dry_run_handler_mock.print_helm_diff.assert_not_called() @@ -432,13 +275,12 @@ async def test_clean_without_to_when_dry_run_is_true( connector = KafkaSinkConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, ) dry_run = True await connector.clean(dry_run) - dry_run_handler_mock.print_helm_diff.assert_called_once() + dry_run_handler_mock.print_helm_diff.assert_not_called() async def test_clean_without_to_when_dry_run_is_false( self, @@ -450,7 +292,6 @@ async def test_clean_without_to_when_dry_run_is_false( connector = KafkaSinkConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, ) mock_destroy = mocker.patch.object(connector, "destroy") @@ -472,51 +313,6 @@ async def test_clean_without_to_when_dry_run_is_false( assert mock.mock_calls == [ mocker.call.destroy_connector(dry_run), - mocker.call.helm.add_repo( - "bakdata-kafka-connect-resetter", - "https://bakdata.github.io/kafka-connect-resetter/", - RepoAuthFlags( - username=None, - password=None, - ca_file=None, - insecure_skip_tls_verify=False, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - mocker.call.helm.upgrade_install( - CONNECTOR_CLEAN_RELEASE_NAME, - "bakdata-kafka-connect-resetter/kafka-connect-resetter", - dry_run, - RESETTER_NAMESPACE, - { - "nameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "fullnameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "connectorType": CONNECTOR_TYPE, - "config": { - "brokers": "broker:9092", - "connector": CONNECTOR_FULL_NAME, - "deleteConsumerGroup": True, - }, - }, - HelmUpgradeInstallFlags( - version="1.0.4", - wait=True, - wait_for_jobs=True, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ ] dry_run_handler_mock.print_helm_diff.assert_not_called() diff --git a/tests/components/test_kafka_source_connector.py b/tests/components/test_kafka_source_connector.py index 76acfe1f9..60ab5af14 100644 --- a/tests/components/test_kafka_source_connector.py +++ b/tests/components/test_kafka_source_connector.py @@ -1,21 +1,16 @@ -from unittest.mock import ANY, MagicMock +from unittest.mock import MagicMock import pytest from pytest_mock import MockerFixture from typing_extensions import override from kpops.component_handlers import get_handlers -from kpops.component_handlers.helm_wrapper.model import ( - HelmUpgradeInstallFlags, - RepoAuthFlags, -) from kpops.component_handlers.kafka_connect.model import ( ConnectorNewState, KafkaConnectorConfig, KafkaConnectorType, ) from kpops.components.base_components.kafka_connector import ( - KafkaConnectorResetter, KafkaSourceConnector, ) from kpops.components.base_components.models import TopicName @@ -28,13 +23,9 @@ ToSection, ) from kpops.components.common.topic import OutputTopicTypes, TopicConfig -from kpops.utils.environment import ENV from tests.components.test_kafka_connector import ( - CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - CONNECTOR_CLEAN_RELEASE_NAME, CONNECTOR_FULL_NAME, CONNECTOR_NAME, - RESETTER_NAMESPACE, TestKafkaConnector, ) @@ -53,7 +44,6 @@ def connector( return KafkaSourceConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, to=ToSection( topics={ TopicName("${output_topic_name}"): TopicConfig( @@ -61,18 +51,8 @@ def connector( ), } ), - offset_topic=OFFSETS_TOPIC, ) - def test_resetter_release_name(self, connector: KafkaSourceConnector) -> None: - assert connector.config.name == CONNECTOR_FULL_NAME - resetter = connector._resetter - assert isinstance(resetter, KafkaConnectorResetter) - assert connector._resetter.helm_release_name == CONNECTOR_CLEAN_RELEASE_NAME - - def test_resetter_offset_topic(self, connector: KafkaSourceConnector) -> None: - assert connector._resetter.values.config.offset_topic == OFFSETS_TOPIC - def test_from_section_raises_exception( self, connector_config: KafkaConnectorConfig, @@ -81,7 +61,6 @@ def test_from_section_raises_exception( KafkaSourceConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, from_=FromSection( topics={ TopicName("connector-topic"): FromTopic( @@ -147,7 +126,6 @@ async def test_destroy( connector: KafkaSourceConnector, mocker: MockerFixture, ) -> None: - ENV["KPOPS_KAFKA_CONNECT_RESETTER_OFFSET_TOPIC"] = OFFSETS_TOPIC assert get_handlers().connector_handler mock_destroy_connector = mocker.patch.object( @@ -167,81 +145,35 @@ async def test_reset_when_dry_run_is_true( mocker: MockerFixture, ) -> None: mock_destroy = mocker.patch.object(connector, "destroy") - mock_resetter_reset = mocker.spy(connector._resetter, "reset") + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" + ) dry_run = True await connector.reset(dry_run=dry_run) mock_destroy.assert_called_once_with(dry_run) - mock_resetter_reset.assert_called_once_with(dry_run) - dry_run_handler_mock.print_helm_diff.assert_called_once() + dry_run_handler_mock.print_helm_diff.assert_not_called() + mock_reset_connector.assert_called_once_with(connector.config, dry_run=dry_run) async def test_reset_when_dry_run_is_false( self, connector: KafkaSourceConnector, dry_run_handler_mock: MagicMock, - helm_mock: MagicMock, mocker: MockerFixture, ) -> None: mock_destroy = mocker.patch.object(connector, "destroy") - mock_delete_topic = mocker.patch.object( get_handlers().topic_handler, "delete_topic" ) - mock_clean_connector = mocker.spy( - get_handlers().connector_handler, "clean_connector" + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" ) - mock = mocker.MagicMock() - mock.attach_mock(mock_destroy, "destroy_connector") - mock.attach_mock(mock_clean_connector, "mock_clean_connector") - mock.attach_mock(helm_mock, "helm") - dry_run = False await connector.reset(dry_run) - assert mock.mock_calls == [ - mocker.call.destroy_connector(dry_run), - mocker.call.helm.add_repo( - "bakdata-kafka-connect-resetter", - "https://bakdata.github.io/kafka-connect-resetter/", - RepoAuthFlags(), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - mocker.call.helm.upgrade_install( - CONNECTOR_CLEAN_RELEASE_NAME, - "bakdata-kafka-connect-resetter/kafka-connect-resetter", - dry_run, - RESETTER_NAMESPACE, - { - "connectorType": CONNECTOR_TYPE, - "config": { - "brokers": "broker:9092", - "connector": CONNECTOR_FULL_NAME, - "offsetTopic": OFFSETS_TOPIC, - }, - "nameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "fullnameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - }, - HelmUpgradeInstallFlags( - version="1.0.4", - wait=True, - wait_for_jobs=True, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - ] + mock_reset_connector.assert_called_once_with(connector.config, dry_run=dry_run) + mock_destroy.assert_called_once_with(dry_run) mock_delete_topic.assert_not_called() dry_run_handler_mock.print_helm_diff.assert_not_called() @@ -252,87 +184,43 @@ async def test_clean_when_dry_run_is_true( ) -> None: await connector.clean(dry_run=True) - dry_run_handler_mock.print_helm_diff.assert_called_once() + dry_run_handler_mock.print_helm_diff.assert_not_called() async def test_clean_when_dry_run_is_false( self, connector: KafkaSourceConnector, - helm_mock: MagicMock, dry_run_handler_mock: MagicMock, mocker: MockerFixture, ) -> None: mock_destroy = mocker.patch.object(connector, "destroy") - mock_delete_topic = mocker.patch.object( get_handlers().topic_handler, "delete_topic" ) - mock_clean_connector = mocker.spy( - get_handlers().connector_handler, "clean_connector" + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" ) mock = mocker.MagicMock() + mock.attach_mock(mock_reset_connector, "mock_reset_connector") mock.attach_mock(mock_destroy, "destroy_connector") mock.attach_mock(mock_delete_topic, "mock_delete_topic") - mock.attach_mock(mock_clean_connector, "mock_clean_connector") - mock.attach_mock(helm_mock, "helm") dry_run = False await connector.clean(dry_run) assert connector.to assert mock.mock_calls == [ + mocker.call.mock_reset_connector(connector.config, dry_run=dry_run), mocker.call.destroy_connector(dry_run), *( mocker.call.mock_delete_topic(topic, dry_run=dry_run) for topic in connector.to.kafka_topics ), - mocker.call.helm.add_repo( - "bakdata-kafka-connect-resetter", - "https://bakdata.github.io/kafka-connect-resetter/", - RepoAuthFlags(), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - mocker.call.helm.upgrade_install( - CONNECTOR_CLEAN_RELEASE_NAME, - "bakdata-kafka-connect-resetter/kafka-connect-resetter", - dry_run, - RESETTER_NAMESPACE, - { - "nameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "fullnameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "connectorType": CONNECTOR_TYPE, - "config": { - "brokers": "broker:9092", - "connector": CONNECTOR_FULL_NAME, - "offsetTopic": OFFSETS_TOPIC, - }, - }, - HelmUpgradeInstallFlags( - version="1.0.4", - wait=True, - wait_for_jobs=True, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ ] - dry_run_handler_mock.print_helm_diff.assert_not_called() async def test_clean_without_to_when_dry_run_is_false( self, - helm_mock: MagicMock, dry_run_handler_mock: MagicMock, mocker: MockerFixture, connector_config: KafkaConnectorConfig, @@ -340,75 +228,29 @@ async def test_clean_without_to_when_dry_run_is_false( connector = KafkaSourceConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, - offset_topic=OFFSETS_TOPIC, ) assert connector.to is None - assert get_handlers().connector_handler - mock_destroy = mocker.patch.object(connector, "destroy") - mock_delete_topic = mocker.patch.object( get_handlers().topic_handler, "delete_topic" ) - mock_clean_connector = mocker.spy( - get_handlers().connector_handler, "clean_connector" + mock_reset_connector = mocker.patch.object( + get_handlers().connector_handler, "reset_connector" ) mock = mocker.MagicMock() + mock.attach_mock(mock_reset_connector, "mock_reset_connector") mock.attach_mock(mock_destroy, "destroy_connector") mock.attach_mock(mock_delete_topic, "mock_delete_topic") - mock.attach_mock(mock_clean_connector, "mock_clean_connector") - mock.attach_mock(helm_mock, "helm") dry_run = False await connector.clean(dry_run) assert mock.mock_calls == [ + mocker.call.mock_reset_connector(connector.config, dry_run=dry_run), mocker.call.destroy_connector(dry_run), - mocker.call.helm.add_repo( - "bakdata-kafka-connect-resetter", - "https://bakdata.github.io/kafka-connect-resetter/", - RepoAuthFlags(), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ - mocker.call.helm.upgrade_install( - CONNECTOR_CLEAN_RELEASE_NAME, - "bakdata-kafka-connect-resetter/kafka-connect-resetter", - dry_run, - RESETTER_NAMESPACE, - { - "nameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "fullnameOverride": CONNECTOR_CLEAN_HELM_NAMEOVERRIDE, - "connectorType": CONNECTOR_TYPE, - "config": { - "brokers": "broker:9092", - "connector": CONNECTOR_FULL_NAME, - "offsetTopic": OFFSETS_TOPIC, - }, - }, - HelmUpgradeInstallFlags( - version="1.0.4", - wait=True, - wait_for_jobs=True, - ), - ), - mocker.call.helm.uninstall( - RESETTER_NAMESPACE, - CONNECTOR_CLEAN_RELEASE_NAME, - dry_run, - ), - ANY, # __bool__ - ANY, # __str__ ] - mock_delete_topic.assert_not_called() dry_run_handler_mock.print_helm_diff.assert_not_called() @@ -420,8 +262,6 @@ async def test_clean_without_to_when_dry_run_is_true( connector = KafkaSourceConnector( name=CONNECTOR_NAME, config=connector_config, - resetter_namespace=RESETTER_NAMESPACE, - offset_topic=OFFSETS_TOPIC, ) assert connector.to is None @@ -429,4 +269,4 @@ async def test_clean_without_to_when_dry_run_is_true( await connector.clean(dry_run=True) - dry_run_handler_mock.print_helm_diff.assert_called_once() + dry_run_handler_mock.print_helm_diff.assert_not_called() diff --git a/tests/pipeline/resources/resetter_values/defaults.yaml b/tests/pipeline/resources/resetter_values/defaults.yaml deleted file mode 100644 index ee97e5e1b..000000000 --- a/tests/pipeline/resources/resetter_values/defaults.yaml +++ /dev/null @@ -1,13 +0,0 @@ -helm-app: - name: "${component.type}" - namespace: "namespace" - values: - label: ${component.name} - kafka: - bootstrapServers: ${config.kafka_brokers} - -kafka-sink-connector: - config: - connector.class: "io.confluent.connect.jdbc.JdbcSinkConnector" - resetter_values: - imageTag: override-default-image-tag diff --git a/tests/pipeline/resources/resetter_values/pipeline.yaml b/tests/pipeline/resources/resetter_values/pipeline.yaml deleted file mode 100644 index ef9d05c76..000000000 --- a/tests/pipeline/resources/resetter_values/pipeline.yaml +++ /dev/null @@ -1 +0,0 @@ -- type: simple-inflate-connectors diff --git a/tests/pipeline/resources/resetter_values/pipeline_connector_only.yaml b/tests/pipeline/resources/resetter_values/pipeline_connector_only.yaml deleted file mode 100644 index af131751d..000000000 --- a/tests/pipeline/resources/resetter_values/pipeline_connector_only.yaml +++ /dev/null @@ -1,4 +0,0 @@ -- type: kafka-sink-connector - name: es-sink-connector - config: - connector.class: io.confluent.connect.elasticsearch.ElasticsearchSinkConnector diff --git a/tests/pipeline/snapshots/test_example/test_generate/atm-fraud/pipeline.yaml b/tests/pipeline/snapshots/test_example/test_generate/atm-fraud/pipeline.yaml index 7acef86c9..67bd46f73 100644 --- a/tests/pipeline/snapshots/test_example/test_generate/atm-fraud/pipeline.yaml +++ b/tests/pipeline/snapshots/test_example/test_generate/atm-fraud/pipeline.yaml @@ -224,4 +224,3 @@ errors.deadletterqueue.topic.replication.factor: '1' errors.tolerance: all pk.mode: record_value - resetter_values: {} diff --git a/tests/pipeline/snapshots/test_example/test_generate/word-count/pipeline.yaml b/tests/pipeline/snapshots/test_example/test_generate/word-count/pipeline.yaml index 056264a29..6628995c0 100644 --- a/tests/pipeline/snapshots/test_example/test_generate/word-count/pipeline.yaml +++ b/tests/pipeline/snapshots/test_example/test_generate/word-count/pipeline.yaml @@ -73,4 +73,3 @@ tasks.max: '1' key.converter: org.apache.kafka.connect.storage.StringConverter value.converter: org.apache.kafka.connect.storage.StringConverter - resetter_values: {} diff --git a/tests/pipeline/snapshots/test_generate/test_inflate_pipeline/pipeline.yaml b/tests/pipeline/snapshots/test_generate/test_inflate_pipeline/pipeline.yaml index 1a9f20984..cb8db80f6 100644 --- a/tests/pipeline/snapshots/test_generate/test_inflate_pipeline/pipeline.yaml +++ b/tests/pipeline/snapshots/test_generate/test_inflate_pipeline/pipeline.yaml @@ -149,7 +149,6 @@ max.buffered.records: '20000' read.timeout.ms: '120000' tasks.max: '1' - resetter_values: {} - type: streams-app name: should-inflate-inflated-streams-app enabled: true diff --git a/tests/pipeline/snapshots/test_generate/test_kafka_connect_sink_weave_from_topics/pipeline.yaml b/tests/pipeline/snapshots/test_generate/test_kafka_connect_sink_weave_from_topics/pipeline.yaml index 2f055caee..fba490286 100644 --- a/tests/pipeline/snapshots/test_generate/test_kafka_connect_sink_weave_from_topics/pipeline.yaml +++ b/tests/pipeline/snapshots/test_generate/test_kafka_connect_sink_weave_from_topics/pipeline.yaml @@ -52,5 +52,4 @@ max.buffered.records: '20000' read.timeout.ms: '120000' tasks.max: '1' - resetter_values: {} diff --git a/tests/pipeline/snapshots/test_generate/test_read_from_component/pipeline.yaml b/tests/pipeline/snapshots/test_generate/test_read_from_component/pipeline.yaml index 71474a107..66a1beac0 100644 --- a/tests/pipeline/snapshots/test_generate/test_read_from_component/pipeline.yaml +++ b/tests/pipeline/snapshots/test_generate/test_read_from_component/pipeline.yaml @@ -107,7 +107,6 @@ max.buffered.records: '20000' read.timeout.ms: '120000' tasks.max: '1' - resetter_values: {} - type: streams-app name: inflate-step-inflated-streams-app enabled: true @@ -209,7 +208,6 @@ max.buffered.records: '20000' read.timeout.ms: '120000' tasks.max: '1' - resetter_values: {} - type: streams-app name: inflate-step-without-prefix-inflated-streams-app enabled: true diff --git a/tests/pipeline/snapshots/test_generate/test_with_env_defaults/pipeline.yaml b/tests/pipeline/snapshots/test_generate/test_with_env_defaults/pipeline.yaml index cf59070dd..3150480d3 100644 --- a/tests/pipeline/snapshots/test_generate/test_with_env_defaults/pipeline.yaml +++ b/tests/pipeline/snapshots/test_generate/test_with_env_defaults/pipeline.yaml @@ -52,5 +52,4 @@ max.buffered.records: '20000' read.timeout.ms: '120000' tasks.max: '1' - resetter_values: {} diff --git a/tests/pipeline/test_generate.py b/tests/pipeline/test_generate.py index 400a875d5..ef6e2e645 100644 --- a/tests/pipeline/test_generate.py +++ b/tests/pipeline/test_generate.py @@ -886,32 +886,6 @@ def test_temp_trim_release_name(self) -> None: == "in-order-to-have-len-fifty-two-name-should-end--here" ) - def test_substitution_in_inflated_component(self) -> None: - pipeline = kpops.generate(RESOURCE_PATH / "resetter_values" / PIPELINE_YAML) - assert isinstance(pipeline.components[1], KafkaSinkConnector) - assert ( - pipeline.components[1]._resetter.values.label == "inflated-connector-name" # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType] - ) - assert ( - pipeline.components[1]._resetter.values.imageTag # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType] - == "override-default-image-tag" - ) - - def test_substitution_in_resetter(self) -> None: - pipeline = kpops.generate( - RESOURCE_PATH - / "resetter_values" - / KpopsFileType.PIPELINE.as_yaml_file(suffix="_connector_only"), - ) - assert isinstance(pipeline.components[0], KafkaSinkConnector) - assert pipeline.components[0].name == "es-sink-connector" - assert pipeline.components[0]._resetter.name == "es-sink-connector" - assert pipeline.components[0]._resetter.values.label == "es-sink-connector" # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType] - assert ( - pipeline.components[0]._resetter.values.imageTag # pyright: ignore[reportAttributeAccessIssue,reportUnknownMemberType] - == "override-default-image-tag" - ) - def test_streams_bootstrap(self, snapshot: Snapshot) -> None: pipeline = kpops.generate( RESOURCE_PATH / "streams-bootstrap" / PIPELINE_YAML,