From 017a69a3379cd5d5a88165ee1c5d766c88dd4a1a Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 10:28:02 +0200 Subject: [PATCH 1/9] feat: expose operations in public Pipeline API --- kpops/api/__init__.py | 94 +++++--------------------------------- kpops/pipeline/__init__.py | 94 +++++++++++++++++++++++++++++++++++++- 2 files changed, 104 insertions(+), 84 deletions(-) diff --git a/kpops/api/__init__.py b/kpops/api/__init__.py index 7d58f6ad5..9425c031b 100644 --- a/kpops/api/__init__.py +++ b/kpops/api/__init__.py @@ -1,7 +1,6 @@ from __future__ import annotations import asyncio -from collections.abc import Iterator from pathlib import Path from typing import TYPE_CHECKING @@ -14,33 +13,20 @@ from kpops.component_handlers.topic.handler import TopicHandler from kpops.component_handlers.topic.kafka_rest import KafkaRest from kpops.config import KpopsConfig -from kpops.core.exception import KpopsException from kpops.core.operation import OperationMode from kpops.core.registry import Registry -from kpops.manifests.kubernetes import KubernetesManifest from kpops.pipeline import ( Pipeline, PipelineGenerator, ) from kpops.utils.cli_commands import init_project -from kpops.utils.logging import log, log_action, log_kpops_exception +from kpops.utils.logging import log if TYPE_CHECKING: - from collections.abc import Awaitable + from collections.abc import Iterator - from kpops.components.base_components.pipeline_component import PipelineComponent from kpops.config import KpopsConfig - - -async def _run_component( - action: str, component: PipelineComponent, operation: Awaitable[None] -) -> None: - log_action(action, component) - try: - await operation - except KpopsException as e: - log_kpops_exception(e) - raise + from kpops.manifests.kubernetes import KubernetesManifest def generate( @@ -106,9 +92,7 @@ def manifest_deploy( verbose=verbose, operation_mode=operation_mode, ) - for component in pipeline.components: - resource = component.manifest_deploy() - yield resource + return pipeline.manifest_deploy() def manifest_destroy( @@ -131,9 +115,7 @@ def manifest_destroy( verbose=verbose, operation_mode=operation_mode, ) - for component in pipeline.components: - resource = component.manifest_destroy() - yield resource + return pipeline.manifest_destroy() def manifest_reset( @@ -156,9 +138,7 @@ def manifest_reset( verbose=verbose, operation_mode=operation_mode, ) - for component in pipeline.components: - resource = component.manifest_reset() - yield resource + return pipeline.manifest_reset() def manifest_clean( @@ -181,9 +161,7 @@ def manifest_clean( verbose=verbose, operation_mode=operation_mode, ) - for component in pipeline.components: - resource = component.manifest_clean() - yield resource + return pipeline.manifest_clean() def deploy( @@ -218,19 +196,7 @@ def deploy( environment=environment, verbose=verbose, ) - - async def deploy_runner(component: PipelineComponent) -> None: - await _run_component("Deploy", component, component.deploy(dry_run)) - - async def async_deploy() -> None: - if parallel: - pipeline_tasks = pipeline.build_execution_graph(deploy_runner) - await pipeline_tasks - else: - for component in pipeline.components: - await deploy_runner(component) - - asyncio.run(async_deploy()) + asyncio.run(pipeline.deploy(dry_run, parallel)) def destroy( @@ -265,21 +231,7 @@ def destroy( environment=environment, verbose=verbose, ) - - async def destroy_runner(component: PipelineComponent) -> None: - await _run_component("Destroy", component, component.destroy(dry_run)) - - async def async_destroy() -> None: - if parallel: - pipeline_tasks = pipeline.build_execution_graph( - destroy_runner, reverse=True - ) - await pipeline_tasks - else: - for component in reversed(pipeline.components): - await destroy_runner(component) - - asyncio.run(async_destroy()) + asyncio.run(pipeline.destroy(dry_run, parallel)) def reset( @@ -314,19 +266,7 @@ def reset( environment=environment, verbose=verbose, ) - - async def reset_runner(component: PipelineComponent) -> None: - await _run_component("Reset", component, component.reset(dry_run)) - - async def async_reset() -> None: - if parallel: - pipeline_tasks = pipeline.build_execution_graph(reset_runner, reverse=True) - await pipeline_tasks - else: - for component in reversed(pipeline.components): - await reset_runner(component) - - asyncio.run(async_reset()) + asyncio.run(pipeline.reset(dry_run, parallel)) def clean( @@ -361,19 +301,7 @@ def clean( environment=environment, verbose=verbose, ) - - async def clean_runner(component: PipelineComponent) -> None: - await _run_component("Clean", component, component.clean(dry_run)) - - async def async_clean() -> None: - if parallel: - pipeline_tasks = pipeline.build_execution_graph(clean_runner, reverse=True) - await pipeline_tasks - else: - for component in reversed(pipeline.components): - await clean_runner(component) - - asyncio.run(async_clean()) + asyncio.run(pipeline.clean(dry_run, parallel)) def init( diff --git a/kpops/pipeline/__init__.py b/kpops/pipeline/__init__.py index 93d00e9bc..59df17fd8 100644 --- a/kpops/pipeline/__init__.py +++ b/kpops/pipeline/__init__.py @@ -15,10 +15,11 @@ from kpops.component_handlers import ComponentHandlers from kpops.components.base_components.pipeline_component import PipelineComponent -from kpops.core.exception import ParsingException, ValidationError +from kpops.core.exception import KpopsException, ParsingException, ValidationError from kpops.core.registry import Registry from kpops.utils.dict_ops import update_nested_pair from kpops.utils.environment import ENV, PIPELINE_PATH +from kpops.utils.logging import log_action, log_kpops_exception from kpops.utils.yaml import CustomSafeDumper, load_yaml_file if TYPE_CHECKING: @@ -26,12 +27,24 @@ from pathlib import Path from kpops.config import KpopsConfig + from kpops.manifests.kubernetes import KubernetesManifest log = structlog.get_logger("PipelineGenerator") ComponentFilterPredicate: TypeAlias = Callable[[PipelineComponent], bool] +async def _run_component( + action: str, component: PipelineComponent, operation: Awaitable[None] +) -> None: + log_action(action, component) + try: + await operation + except KpopsException as e: + log_kpops_exception(e) + raise + + @dataclass class Pipeline: """Pipeline representation.""" @@ -143,6 +156,85 @@ async def run_graph_layers( return run_graph_layers(sorted_layers) + async def deploy(self, dry_run: bool, parallel: bool = False) -> None: + """Deploy pipeline steps. + + :param dry_run: Whether to dry run the command or execute it. + :param parallel: Enable or disable parallel execution of pipeline steps. + """ + await self._run_action( + "Deploy", + lambda component: component.deploy(dry_run), + parallel, + reverse=False, + ) + + async def destroy(self, dry_run: bool, parallel: bool = False) -> None: + """Destroy pipeline steps. + + :param dry_run: Whether to dry run the command or execute it. + :param parallel: Enable or disable parallel execution of pipeline steps. + """ + await self._run_action( + "Destroy", + lambda component: component.destroy(dry_run), + parallel, + reverse=True, + ) + + async def reset(self, dry_run: bool, parallel: bool = False) -> None: + """Reset pipeline steps. + + :param dry_run: Whether to dry run the command or execute it. + :param parallel: Enable or disable parallel execution of pipeline steps. + """ + await self._run_action( + "Reset", lambda component: component.reset(dry_run), parallel, reverse=True + ) + + async def clean(self, dry_run: bool, parallel: bool = False) -> None: + """Clean pipeline steps. + + :param dry_run: Whether to dry run the command or execute it. + :param parallel: Enable or disable parallel execution of pipeline steps. + """ + await self._run_action( + "Clean", lambda component: component.clean(dry_run), parallel, reverse=True + ) + + def manifest_deploy(self) -> Iterator[tuple[KubernetesManifest, ...]]: + for component in self.components: + yield component.manifest_deploy() + + def manifest_destroy(self) -> Iterator[tuple[KubernetesManifest, ...]]: + for component in self.components: + yield component.manifest_destroy() + + def manifest_reset(self) -> Iterator[tuple[KubernetesManifest, ...]]: + for component in self.components: + yield component.manifest_reset() + + def manifest_clean(self) -> Iterator[tuple[KubernetesManifest, ...]]: + for component in self.components: + yield component.manifest_clean() + + async def _run_action( + self, + action_name: str, + component_action: Callable[[PipelineComponent], Coroutine[Any, Any, None]], + parallel: bool, + reverse: bool, + ) -> None: + async def runner(component: PipelineComponent) -> None: + await _run_component(action_name, component, component_action(component)) + + if parallel: + await self.build_execution_graph(runner, reverse=reverse) + else: + components = reversed(self.components) if reverse else self.components + for component in components: + await runner(component) + def __getitem__(self, component_id: str) -> PipelineComponent: try: return self._component_index[component_id] From 5099af8c6354a8d084f80461a4c9a629c84e3195 Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 10:28:52 +0200 Subject: [PATCH 2/9] test: use Pipeline API for testing --- tests/api/test_run_component.py | 2 +- tests/components/conftest.py | 30 ++++-------- tests/components/test_kubernetes_app.py | 7 +-- tests/conftest.py | 45 +++++++++++++++++ tests/pipeline/test_clean.py | 65 ++++++++++++------------- uv.lock | 2 +- 6 files changed, 90 insertions(+), 61 deletions(-) diff --git a/tests/api/test_run_component.py b/tests/api/test_run_component.py index 4e57456e4..056be3a84 100644 --- a/tests/api/test_run_component.py +++ b/tests/api/test_run_component.py @@ -6,9 +6,9 @@ from structlog.contextvars import merge_contextvars from structlog.testing import capture_logs -from kpops.api import _run_component from kpops.components.base_components.pipeline_component import PipelineComponent from kpops.core.exception import KpopsException +from kpops.pipeline import _run_component from kpops.utils.logging import bound_service_context diff --git a/tests/components/conftest.py b/tests/components/conftest.py index 31b40b30e..1fb0538c7 100644 --- a/tests/components/conftest.py +++ b/tests/components/conftest.py @@ -1,29 +1,19 @@ -from unittest import mock +from pathlib import Path import pytest from kpops.component_handlers import ComponentHandlers -from kpops.config import KpopsConfig, TopicNameConfig, set_config +from kpops.config import KpopsConfig from tests.components import PIPELINE_BASE_DIR -@pytest.fixture(autouse=True, scope="module") -def config() -> None: - config = KpopsConfig( - topic_name_config=TopicNameConfig( - default_error_topic_name="${component.type}-error-topic", - default_output_topic_name="${component.type}-output-topic", - ), - kafka_brokers="broker:9092", - pipeline_base_dir=PIPELINE_BASE_DIR, - ) - set_config(config) +@pytest.fixture(scope="module") +def pipeline_base_dir() -> Path: + return PIPELINE_BASE_DIR -@pytest.fixture(autouse=True, scope="module") -def handlers() -> None: - ComponentHandlers( - schema_handler=mock.AsyncMock(), - connector_handler=mock.AsyncMock(), - topic_handler=mock.AsyncMock(), - ) +@pytest.fixture(autouse=True) +def _apply_config_and_handlers( + config: KpopsConfig, handlers: ComponentHandlers +) -> None: + pass diff --git a/tests/components/test_kubernetes_app.py b/tests/components/test_kubernetes_app.py index a0d1fa1d7..244772f4d 100644 --- a/tests/components/test_kubernetes_app.py +++ b/tests/components/test_kubernetes_app.py @@ -8,7 +8,6 @@ KubernetesApp, KubernetesAppValues, ) -from kpops.config import KpopsConfig HELM_RELEASE_NAME = create_helm_release_name("${pipeline.name}-test-kubernetes-app") @@ -17,12 +16,8 @@ class KubernetesTestValues(KubernetesAppValues): foo: str -@pytest.mark.usefixtures("mock_env", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env") class TestKubernetesApp: - @pytest.fixture(autouse=True) - def config(self) -> KpopsConfig: - return KpopsConfig.create(None, verbose=False) - @pytest.fixture() def app_values(self) -> KubernetesTestValues: return KubernetesTestValues(foo="foo") diff --git a/tests/conftest.py b/tests/conftest.py index fa137b814..12ecbb166 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -7,6 +7,8 @@ import pytest +from kpops.component_handlers import ComponentHandlers +from kpops.config import KpopsConfig, TopicNameConfig, set_config from kpops.utils.environment import ENV, Environment from kpops.utils.yaml import load_yaml_file @@ -69,6 +71,49 @@ def clear_kpops_config() -> Iterator[None]: yield +@pytest.fixture(scope="module") +def pipeline_base_dir() -> Path: + """Return the base directory used by the ``config`` fixture. + + Override this fixture in a more specific ``conftest.py`` to point + ``KpopsConfig.pipeline_base_dir`` elsewhere. + """ + return Path() + + +@pytest.fixture(scope="module") +def config(pipeline_base_dir: Path) -> KpopsConfig: + """Provide a ready-to-use ``KpopsConfig`` for tests that construct components directly. + + Not autouse: opt in explicitly via ``usefixtures`` or a local autouse + wrapper fixture where needed. + """ + config = KpopsConfig( + topic_name_config=TopicNameConfig( + default_error_topic_name="${component.type}-error-topic", + default_output_topic_name="${component.type}-output-topic", + ), + kafka_brokers="broker:9092", + pipeline_base_dir=pipeline_base_dir, + ) + set_config(config) + return config + + +@pytest.fixture(scope="module") +def handlers() -> ComponentHandlers: + """Provide ``ComponentHandlers`` with mocked handlers. + + Not autouse: opt in explicitly via ``usefixtures`` or a local autouse + wrapper fixture where needed. + """ + return ComponentHandlers( + schema_handler=mock.AsyncMock(), + connector_handler=mock.AsyncMock(), + topic_handler=mock.AsyncMock(), + ) + + KUBECONFIG = """ apiVersion: v1 clusters: diff --git a/tests/pipeline/test_clean.py b/tests/pipeline/test_clean.py index 9a4de24fb..989df7c4b 100644 --- a/tests/pipeline/test_clean.py +++ b/tests/pipeline/test_clean.py @@ -1,28 +1,25 @@ -from pathlib import Path from unittest.mock import MagicMock import pytest from pytest_mock import MockerFixture -from typer.testing import CliRunner -from kpops.cli.main import app from kpops.component_handlers.helm.helm import Helm from kpops.components.base_components import HelmApp +from kpops.components.base_components.helm_app import HelmAppValues +from kpops.components.streams_bootstrap.producer.model import ProducerAppValues from kpops.components.streams_bootstrap.producer.producer_app import ( ProducerApp, ProducerAppCleaner, ) +from kpops.components.streams_bootstrap.streams.model import StreamsAppValues from kpops.components.streams_bootstrap.streams.streams_app import ( StreamsApp, StreamsAppCleaner, ) +from kpops.pipeline import Pipeline -runner = CliRunner() -RESOURCE_PATH = Path(__file__).parent / "resources" - - -@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env", "config", "handlers") class TestClean: @pytest.fixture(autouse=True) def helm_mock(self, mocker: MockerFixture) -> MagicMock: @@ -32,22 +29,33 @@ def helm_mock(self, mocker: MockerFixture) -> MagicMock: ) return helm_mock - # TODO: test using public Pipeline API - # @pytest.fixture() - # def pipeline(self) -> Pipeline: - # pipeline = Pipeline() - # pipeline.add( - # ProducerApp( - # name="producer", - # namespace="test-namespace", - # values=ProducerAppValues( - # streams=ProducerStreamsConfig(brokers=get_config().kafka_brokers) - # ), - # ) - # ) - # return pipeline + @pytest.fixture() + def pipeline(self) -> Pipeline: + pipeline = Pipeline() + pipeline.add( + ProducerApp( + name="producer", + namespace="test-namespace", + values=ProducerAppValues(image="producer-image"), + ) + ) + pipeline.add( + StreamsApp( + name="streams", + namespace="test-namespace", + values=StreamsAppValues(image="streams-image"), + ) + ) + pipeline.add( + HelmApp( + name="helm-app", + namespace="test-namespace", + values=HelmAppValues(), + ) + ) + return pipeline - def test_order(self, mocker: MockerFixture) -> None: + async def test_order(self, pipeline: Pipeline, mocker: MockerFixture) -> None: # destroy producer_app_mock_destroy = mocker.patch.object(ProducerApp, "destroy") streams_app_mock_destroy = mocker.patch.object(StreamsApp, "destroy") @@ -65,16 +73,7 @@ def test_order(self, mocker: MockerFixture) -> None: async_mocker.attach_mock(producer_app_mock_clean, "producer_app_mock_clean") async_mocker.attach_mock(streams_app_mock_clean, "streams_app_mock_clean") - result = runner.invoke( - app, - [ - "clean", - str(RESOURCE_PATH / "simple-pipeline" / "pipeline.yaml"), - ], - catch_exceptions=False, - ) - - assert result.exit_code == 0, result.stdout + await pipeline.clean(dry_run=True) # check called producer_app_mock_destroy.assert_called_once_with(True) diff --git a/uv.lock b/uv.lock index 853a6ca9e..e28a867e9 100644 --- a/uv.lock +++ b/uv.lock @@ -448,7 +448,7 @@ wheels = [ [[package]] name = "kpops" -version = "10.9.0" +version = "11.0.0" source = { editable = "." } dependencies = [ { name = "anyio" }, From 81092bb534c9ba9c0406018a92f183d256bc1927 Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 10:31:54 +0200 Subject: [PATCH 3/9] test(components): set module fixture scope --- tests/components/conftest.py | 2 +- tests/components/test_kubernetes_app.py | 21 ++++++++++++--------- 2 files changed, 13 insertions(+), 10 deletions(-) diff --git a/tests/components/conftest.py b/tests/components/conftest.py index 1fb0538c7..42bd181be 100644 --- a/tests/components/conftest.py +++ b/tests/components/conftest.py @@ -12,7 +12,7 @@ def pipeline_base_dir() -> Path: return PIPELINE_BASE_DIR -@pytest.fixture(autouse=True) +@pytest.fixture(autouse=True, scope="module") def _apply_config_and_handlers( config: KpopsConfig, handlers: ComponentHandlers ) -> None: diff --git a/tests/components/test_kubernetes_app.py b/tests/components/test_kubernetes_app.py index 244772f4d..b40e23c91 100644 --- a/tests/components/test_kubernetes_app.py +++ b/tests/components/test_kubernetes_app.py @@ -12,36 +12,39 @@ HELM_RELEASE_NAME = create_helm_release_name("${pipeline.name}-test-kubernetes-app") -class KubernetesTestValues(KubernetesAppValues): +class DummyKubernetesApp(KubernetesApp): ... + + +class DummyKubernetesValues(KubernetesAppValues): foo: str @pytest.mark.usefixtures("mock_env") class TestKubernetesApp: @pytest.fixture() - def app_values(self) -> KubernetesTestValues: - return KubernetesTestValues(foo="foo") + def app_values(self) -> DummyKubernetesValues: + return DummyKubernetesValues(foo="foo") @pytest.fixture() def repo_config(self) -> HelmRepoConfig: return HelmRepoConfig(repository_name="test", url="https://bakdata.com") @pytest.fixture() - def kubernetes_app(self, app_values: KubernetesTestValues) -> KubernetesApp: - return KubernetesApp( + def kubernetes_app(self, app_values: DummyKubernetesValues) -> DummyKubernetesApp: + return DummyKubernetesApp( name="test-kubernetes-app", values=app_values, namespace="test-namespace", ) def test_should_raise_value_error_when_name_is_not_valid( - self, app_values: KubernetesTestValues + self, app_values: DummyKubernetesValues ) -> None: with pytest.raises( ValueError, match=r"The component name .* is invalid for Kubernetes\.", ): - KubernetesApp( + DummyKubernetesApp( name="Not-Compatible*", values=app_values, namespace="test-namespace", @@ -51,13 +54,13 @@ def test_should_raise_value_error_when_name_is_not_valid( ValueError, match=r"The component name .* is invalid for Kubernetes\.", ): - KubernetesApp( + DummyKubernetesApp( name="snake_case*", values=app_values, namespace="test-namespace", ) - assert KubernetesApp( + assert DummyKubernetesApp( name="valid-name", values=app_values, namespace="test-namespace", From 936579beef5f41bc5cfb2af7c187e5cacc60496c Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 10:34:11 +0200 Subject: [PATCH 4/9] style: cleanup comments --- tests/conftest.py | 15 --------------- 1 file changed, 15 deletions(-) diff --git a/tests/conftest.py b/tests/conftest.py index 12ecbb166..67f8aa5e0 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -73,21 +73,11 @@ def clear_kpops_config() -> Iterator[None]: @pytest.fixture(scope="module") def pipeline_base_dir() -> Path: - """Return the base directory used by the ``config`` fixture. - - Override this fixture in a more specific ``conftest.py`` to point - ``KpopsConfig.pipeline_base_dir`` elsewhere. - """ return Path() @pytest.fixture(scope="module") def config(pipeline_base_dir: Path) -> KpopsConfig: - """Provide a ready-to-use ``KpopsConfig`` for tests that construct components directly. - - Not autouse: opt in explicitly via ``usefixtures`` or a local autouse - wrapper fixture where needed. - """ config = KpopsConfig( topic_name_config=TopicNameConfig( default_error_topic_name="${component.type}-error-topic", @@ -102,11 +92,6 @@ def config(pipeline_base_dir: Path) -> KpopsConfig: @pytest.fixture(scope="module") def handlers() -> ComponentHandlers: - """Provide ``ComponentHandlers`` with mocked handlers. - - Not autouse: opt in explicitly via ``usefixtures`` or a local autouse - wrapper fixture where needed. - """ return ComponentHandlers( schema_handler=mock.AsyncMock(), connector_handler=mock.AsyncMock(), From 205820639d80ca701201b8ce547bed26b219d3cb Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 11:11:22 +0200 Subject: [PATCH 5/9] test: move --- tests/{api => pipeline}/test_run_component.py | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename tests/{api => pipeline}/test_run_component.py (100%) diff --git a/tests/api/test_run_component.py b/tests/pipeline/test_run_component.py similarity index 100% rename from tests/api/test_run_component.py rename to tests/pipeline/test_run_component.py From aa6a6814748c6dd91cd6c9717df81348b6060167 Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 11:36:59 +0200 Subject: [PATCH 6/9] test(pipeline): rewrite using new Pipeline API --- tests/pipeline/conftest.py | 49 +++++++++++++++++++ .../resources/simple-pipeline/pipeline.yaml | 6 --- tests/pipeline/test_clean.py | 40 --------------- tests/pipeline/test_deploy.py | 32 ++---------- tests/pipeline/test_destroy.py | 32 ++---------- tests/pipeline/test_reset.py | 34 ++----------- 6 files changed, 61 insertions(+), 132 deletions(-) create mode 100644 tests/pipeline/conftest.py delete mode 100644 tests/pipeline/resources/simple-pipeline/pipeline.yaml diff --git a/tests/pipeline/conftest.py b/tests/pipeline/conftest.py new file mode 100644 index 000000000..eb2725b93 --- /dev/null +++ b/tests/pipeline/conftest.py @@ -0,0 +1,49 @@ +from unittest.mock import MagicMock + +import pytest +from pytest_mock import MockerFixture + +from kpops.component_handlers.helm.helm import Helm +from kpops.components.base_components import HelmApp +from kpops.components.base_components.helm_app import HelmAppValues +from kpops.components.streams_bootstrap.producer.model import ProducerAppValues +from kpops.components.streams_bootstrap.producer.producer_app import ProducerApp +from kpops.components.streams_bootstrap.streams.model import StreamsAppValues +from kpops.components.streams_bootstrap.streams.streams_app import StreamsApp +from kpops.pipeline import Pipeline + + +@pytest.fixture() +def helm_mock(mocker: MockerFixture) -> MagicMock: + helm_mock = mocker.MagicMock(Helm) + mocker.patch( + "kpops.components.base_components.helm_app.Helm", return_value=helm_mock + ) + return helm_mock + + +@pytest.fixture() +def pipeline(helm_mock: MagicMock) -> Pipeline: + pipeline = Pipeline() + pipeline.add( + ProducerApp( + name="producer", + namespace="test-namespace", + values=ProducerAppValues(image="producer-image"), + ) + ) + pipeline.add( + StreamsApp( + name="streams", + namespace="test-namespace", + values=StreamsAppValues(image="streams-image"), + ) + ) + pipeline.add( + HelmApp( + name="helm-app", + namespace="test-namespace", + values=HelmAppValues(), + ) + ) + return pipeline diff --git a/tests/pipeline/resources/simple-pipeline/pipeline.yaml b/tests/pipeline/resources/simple-pipeline/pipeline.yaml deleted file mode 100644 index f78d0c385..000000000 --- a/tests/pipeline/resources/simple-pipeline/pipeline.yaml +++ /dev/null @@ -1,6 +0,0 @@ -- type: producer-app - -- type: streams-app - -- type: helm-app - values: {} diff --git a/tests/pipeline/test_clean.py b/tests/pipeline/test_clean.py index 989df7c4b..a942d116e 100644 --- a/tests/pipeline/test_clean.py +++ b/tests/pipeline/test_clean.py @@ -1,17 +1,11 @@ -from unittest.mock import MagicMock - import pytest from pytest_mock import MockerFixture -from kpops.component_handlers.helm.helm import Helm from kpops.components.base_components import HelmApp -from kpops.components.base_components.helm_app import HelmAppValues -from kpops.components.streams_bootstrap.producer.model import ProducerAppValues from kpops.components.streams_bootstrap.producer.producer_app import ( ProducerApp, ProducerAppCleaner, ) -from kpops.components.streams_bootstrap.streams.model import StreamsAppValues from kpops.components.streams_bootstrap.streams.streams_app import ( StreamsApp, StreamsAppCleaner, @@ -21,40 +15,6 @@ @pytest.mark.usefixtures("mock_env", "config", "handlers") class TestClean: - @pytest.fixture(autouse=True) - def helm_mock(self, mocker: MockerFixture) -> MagicMock: - helm_mock = mocker.MagicMock(Helm) - mocker.patch( - "kpops.components.base_components.helm_app.Helm", return_value=helm_mock - ) - return helm_mock - - @pytest.fixture() - def pipeline(self) -> Pipeline: - pipeline = Pipeline() - pipeline.add( - ProducerApp( - name="producer", - namespace="test-namespace", - values=ProducerAppValues(image="producer-image"), - ) - ) - pipeline.add( - StreamsApp( - name="streams", - namespace="test-namespace", - values=StreamsAppValues(image="streams-image"), - ) - ) - pipeline.add( - HelmApp( - name="helm-app", - namespace="test-namespace", - values=HelmAppValues(), - ) - ) - return pipeline - async def test_order(self, pipeline: Pipeline, mocker: MockerFixture) -> None: # destroy producer_app_mock_destroy = mocker.patch.object(ProducerApp, "destroy") diff --git a/tests/pipeline/test_deploy.py b/tests/pipeline/test_deploy.py index 0be9589c7..4fc791765 100644 --- a/tests/pipeline/test_deploy.py +++ b/tests/pipeline/test_deploy.py @@ -1,29 +1,14 @@ -from pathlib import Path -from unittest.mock import AsyncMock - import pytest from pytest_mock import MockerFixture -from typer.testing import CliRunner -from kpops.cli.main import app from kpops.components.base_components import HelmApp from kpops.components.streams_bootstrap import ProducerApp, StreamsApp - -runner = CliRunner() - -RESOURCE_PATH = Path(__file__).parent / "resources" +from kpops.pipeline import Pipeline -@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env", "config", "handlers") class TestDeploy: - @pytest.fixture(autouse=True) - def mock_helm(self, mocker: MockerFixture) -> AsyncMock: - return mocker.patch( - "kpops.components.base_components.helm_app.Helm", - return_value=AsyncMock(), - ).return_value - - def test_order(self, mocker: MockerFixture) -> None: + async def test_order(self, pipeline: Pipeline, mocker: MockerFixture) -> None: producer_app_mock_deploy = mocker.patch.object(ProducerApp, "deploy") streams_app_mock_deploy = mocker.patch.object(StreamsApp, "deploy") helm_app_mock_deploy = mocker.patch.object(HelmApp, "deploy") @@ -32,16 +17,7 @@ def test_order(self, mocker: MockerFixture) -> None: mock_deploy.attach_mock(streams_app_mock_deploy, "streams_app_mock_deploy") mock_deploy.attach_mock(helm_app_mock_deploy, "helm_app_mock_deploy") - result = runner.invoke( - app, - [ - "deploy", - str(RESOURCE_PATH / "simple-pipeline" / "pipeline.yaml"), - ], - catch_exceptions=False, - ) - - assert result.exit_code == 0, result.stdout + await pipeline.deploy(dry_run=True) # check called producer_app_mock_deploy.assert_called_once_with(True) diff --git a/tests/pipeline/test_destroy.py b/tests/pipeline/test_destroy.py index f39f0779f..842d3642b 100644 --- a/tests/pipeline/test_destroy.py +++ b/tests/pipeline/test_destroy.py @@ -1,29 +1,14 @@ -from pathlib import Path -from unittest.mock import AsyncMock - import pytest from pytest_mock import MockerFixture -from typer.testing import CliRunner -from kpops.cli.main import app from kpops.components.base_components import HelmApp from kpops.components.streams_bootstrap import ProducerApp, StreamsApp - -runner = CliRunner() - -RESOURCE_PATH = Path(__file__).parent / "resources" +from kpops.pipeline import Pipeline -@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env", "config", "handlers") class TestDestroy: - @pytest.fixture(autouse=True) - def mock_helm(self, mocker: MockerFixture) -> AsyncMock: - return mocker.patch( - "kpops.components.base_components.helm_app.Helm", - return_value=AsyncMock(), - ).return_value - - def test_order(self, mocker: MockerFixture) -> None: + async def test_order(self, pipeline: Pipeline, mocker: MockerFixture) -> None: producer_app_mock_destroy = mocker.patch.object(ProducerApp, "destroy") streams_app_mock_destroy = mocker.patch.object(StreamsApp, "destroy") helm_app_mock_destroy = mocker.patch.object(HelmApp, "destroy") @@ -32,16 +17,7 @@ def test_order(self, mocker: MockerFixture) -> None: mock_destroy.attach_mock(streams_app_mock_destroy, "streams_app_mock_destroy") mock_destroy.attach_mock(helm_app_mock_destroy, "helm_app_mock_destroy") - result = runner.invoke( - app, - [ - "destroy", - str(RESOURCE_PATH / "simple-pipeline" / "pipeline.yaml"), - ], - catch_exceptions=False, - ) - - assert result.exit_code == 0, result.stdout + await pipeline.destroy(dry_run=True) # check called producer_app_mock_destroy.assert_called_once_with(True) diff --git a/tests/pipeline/test_reset.py b/tests/pipeline/test_reset.py index 4ff284025..8a69bc412 100644 --- a/tests/pipeline/test_reset.py +++ b/tests/pipeline/test_reset.py @@ -1,35 +1,18 @@ -from pathlib import Path -from unittest.mock import MagicMock - import pytest from pytest_mock import MockerFixture -from typer.testing import CliRunner -from kpops.cli.main import app -from kpops.component_handlers.helm.helm import Helm from kpops.components.base_components import HelmApp from kpops.components.streams_bootstrap import ProducerApp, StreamsApp from kpops.components.streams_bootstrap.producer.producer_app import ( ProducerAppCleaner, ) from kpops.components.streams_bootstrap.streams.streams_app import StreamsAppCleaner - -runner = CliRunner() - -RESOURCE_PATH = Path(__file__).parent / "resources" +from kpops.pipeline import Pipeline -@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env", "config", "handlers") class TestReset: - @pytest.fixture(autouse=True) - def helm_mock(self, mocker: MockerFixture) -> MagicMock: - helm_mock = mocker.MagicMock(Helm) - mocker.patch( - "kpops.components.base_components.helm_app.Helm", return_value=helm_mock - ) - return helm_mock - - def test_order(self, mocker: MockerFixture) -> None: + async def test_order(self, pipeline: Pipeline, mocker: MockerFixture) -> None: # destroy producer_app_mock_destroy = mocker.patch.object(ProducerApp, "destroy") streams_app_mock_destroy = mocker.patch.object(StreamsApp, "destroy") @@ -47,16 +30,7 @@ def test_order(self, mocker: MockerFixture) -> None: async_mocker.attach_mock(producer_app_mock_reset, "producer_app_mock_reset") async_mocker.attach_mock(streams_app_mock_reset, "streams_app_mock_reset") - result = runner.invoke( - app, - [ - "reset", - str(RESOURCE_PATH / "simple-pipeline" / "pipeline.yaml"), - ], - catch_exceptions=False, - ) - - assert result.exit_code == 0, result.stdout + await pipeline.reset(dry_run=True) # check called producer_app_mock_destroy.assert_called_once_with(True) From 49a7e5d154d0ba7723eeb9df31372231aa7ae3ef Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 11:55:07 +0200 Subject: [PATCH 7/9] test: improve isolation --- tests/api/test_handlers.py | 1 + tests/conftest.py | 6 ++++++ 2 files changed, 7 insertions(+) diff --git a/tests/api/test_handlers.py b/tests/api/test_handlers.py index 82a7c77f9..1b458e2d3 100644 --- a/tests/api/test_handlers.py +++ b/tests/api/test_handlers.py @@ -30,6 +30,7 @@ def handlers() -> Generator[ComponentHandlers, None, None]: ComponentHandlers._instance = None +@pytest.mark.usefixtures("clear_handlers") def test_global_handlers_not_initialized() -> None: with pytest.raises( RuntimeError, match="ComponentHandlers has not been initialized" diff --git a/tests/conftest.py b/tests/conftest.py index 67f8aa5e0..b3eb5d5a7 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -71,6 +71,12 @@ def clear_kpops_config() -> Iterator[None]: yield +@pytest.fixture(scope="module") +def clear_handlers() -> Iterator[None]: + ComponentHandlers._instance = None + yield + + @pytest.fixture(scope="module") def pipeline_base_dir() -> Path: return Path() From dbd5325baa377ac72d30dd7321c23ffbb05be65a Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 11:56:25 +0200 Subject: [PATCH 8/9] test: cleanup fixtures --- tests/cli/test_init.py | 2 +- tests/conftest.py | 4 +--- tests/pipeline/test_example.py | 2 +- tests/pipeline/test_generate.py | 2 +- tests/test_kpops_config.py | 2 +- 5 files changed, 5 insertions(+), 7 deletions(-) diff --git a/tests/cli/test_init.py b/tests/cli/test_init.py index 9b1285a44..9ff7a8618 100644 --- a/tests/cli/test_init.py +++ b/tests/cli/test_init.py @@ -23,7 +23,7 @@ def test_create_config(tmp_path: Path) -> None: assert len(opt_conf.read_text()) > len(req_conf.read_text()) -@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_config") def test_init_project_exclude_optional(tmp_path: Path, snapshot: Snapshot) -> None: req_path = tmp_path / "req" req_path.mkdir() diff --git a/tests/conftest.py b/tests/conftest.py index b3eb5d5a7..6cf60f9fe 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -64,9 +64,7 @@ def custom_components() -> Iterator[None]: @pytest.fixture(scope="module") -def clear_kpops_config() -> Iterator[None]: - from kpops.config import KpopsConfig - +def clear_config() -> Iterator[None]: KpopsConfig._instance = None yield diff --git a/tests/pipeline/test_example.py b/tests/pipeline/test_example.py index 6ef21761f..0040542cc 100644 --- a/tests/pipeline/test_example.py +++ b/tests/pipeline/test_example.py @@ -13,7 +13,7 @@ EXAMPLES_PATH = Path("examples").absolute() -@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_kpops_config") +@pytest.mark.usefixtures("mock_env", "load_yaml_file_clear_cache", "clear_config") class TestExample: @pytest.fixture(scope="class", autouse=True) @classmethod diff --git a/tests/pipeline/test_generate.py b/tests/pipeline/test_generate.py index 988609092..45d094fa9 100644 --- a/tests/pipeline/test_generate.py +++ b/tests/pipeline/test_generate.py @@ -29,7 +29,7 @@ @pytest.mark.usefixtures( - "mock_env", "load_yaml_file_clear_cache", "custom_components", "clear_kpops_config" + "mock_env", "load_yaml_file_clear_cache", "custom_components", "clear_config" ) class TestGenerate: def test_python_api(self) -> None: diff --git a/tests/test_kpops_config.py b/tests/test_kpops_config.py index 954beac9f..a15e82579 100644 --- a/tests/test_kpops_config.py +++ b/tests/test_kpops_config.py @@ -73,7 +73,7 @@ def test_kpops_config_with_different_invalid_urls() -> None: ) -@pytest.mark.usefixtures("clear_kpops_config") +@pytest.mark.usefixtures("clear_config") def test_global_kpops_config_not_initialized_error() -> None: with pytest.raises( RuntimeError, From a2800bfe2d135140bd0657e209a1a4ad9a7d7ab5 Mon Sep 17 00:00:00 2001 From: Salomon Popp Date: Thu, 13 Aug 2026 14:44:40 +0200 Subject: [PATCH 9/9] refactor: apply review feedback --- kpops/api/__init__.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/kpops/api/__init__.py b/kpops/api/__init__.py index 9425c031b..f440d7830 100644 --- a/kpops/api/__init__.py +++ b/kpops/api/__init__.py @@ -92,7 +92,7 @@ def manifest_deploy( verbose=verbose, operation_mode=operation_mode, ) - return pipeline.manifest_deploy() + yield from pipeline.manifest_deploy() def manifest_destroy( @@ -115,7 +115,7 @@ def manifest_destroy( verbose=verbose, operation_mode=operation_mode, ) - return pipeline.manifest_destroy() + yield from pipeline.manifest_destroy() def manifest_reset( @@ -138,7 +138,7 @@ def manifest_reset( verbose=verbose, operation_mode=operation_mode, ) - return pipeline.manifest_reset() + yield from pipeline.manifest_reset() def manifest_clean( @@ -161,7 +161,7 @@ def manifest_clean( verbose=verbose, operation_mode=operation_mode, ) - return pipeline.manifest_clean() + yield from pipeline.manifest_clean() def deploy(