diff --git a/argo_rollouts/tests/conftest.py b/argo_rollouts/tests/conftest.py index 2ddc79b1fa8ac..e8f84ac610ec2 100644 --- a/argo_rollouts/tests/conftest.py +++ b/argo_rollouts/tests/conftest.py @@ -2,13 +2,11 @@ # All rights reserved # Licensed under a 3-clause BSD style license (see LICENSE) import os -from contextlib import ExitStack import pytest from datadog_checks.dev import get_here from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward from datadog_checks.dev.subprocess import run_command HERE = get_here() @@ -27,10 +25,11 @@ def setup_argo_rollouts(): @pytest.fixture(scope='session') def dd_environment(): - with kind_run(conditions=[setup_argo_rollouts], sleep=30) as kubeconfig, ExitStack() as stack: - argo_rollouts_host, argo_rollouts_port = stack.enter_context( - port_forward(kubeconfig, 'argo-rollouts', 8090, 'deployment', 'argo-rollouts') - ) - instances = [{'openmetrics_endpoint': f'http://{argo_rollouts_host}:{argo_rollouts_port}/metrics'}] + with kind_run(conditions=[setup_argo_rollouts], sleep=30) as kubeconfig: + instances = [ + {'openmetrics_endpoint': 'http://argo-rollouts-metrics.argo-rollouts.svc.cluster.local:8090/metrics'} + ] - yield {'instances': instances} + metadata = {'agent_type': 'kubernetes', 'kubernetes': {'kubeconfig': kubeconfig}} + + yield {'instances': instances}, metadata diff --git a/argocd/tests/conftest.py b/argocd/tests/conftest.py index 4d70e9428d1e2..cc4ab0e104325 100644 --- a/argocd/tests/conftest.py +++ b/argocd/tests/conftest.py @@ -8,15 +8,8 @@ from datadog_checks.dev import get_here from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward from datadog_checks.dev.subprocess import run_command -try: - from contextlib import ExitStack -except ImportError: - from contextlib2 import ExitStack - - HERE = get_here() opj = os.path.join @@ -40,39 +33,25 @@ def setup_argocd(): @pytest.fixture(scope='session') def dd_environment(dd_save_state): with kind_run(conditions=[setup_argocd]) as kubeconfig: - with ExitStack() as stack: - app_controller_host, app_controller_port = stack.enter_context( - port_forward(kubeconfig, 'argocd', 8082, 'service', 'argocd-metrics') - ) - appset_controller_host, appset_controller_port = stack.enter_context( - port_forward(kubeconfig, 'argocd', 8080, 'service', 'argocd-applicationset-controller') - ) - api_server_host, api_server_port = stack.enter_context( - port_forward(kubeconfig, 'argocd', 8083, 'service', 'argocd-server-metrics') - ) - repo_server_host, repo_server_port = stack.enter_context( - port_forward(kubeconfig, 'argocd', 8084, 'service', 'argocd-repo-server') - ) - notifications_controller_host, notifications_controller_port = stack.enter_context( - port_forward(kubeconfig, 'argocd', 9001, 'service', 'argocd-notifications-controller-metrics') - ) - app_controller_endpoint = 'http://{}:{}/metrics'.format(app_controller_host, app_controller_port) - appset_controller_endpoint = 'http://{}:{}/metrics'.format(appset_controller_host, appset_controller_port) - api_server_endpoint = 'http://{}:{}/metrics'.format(api_server_host, api_server_port) - repo_server_endpoint = 'http://{}:{}/metrics'.format(repo_server_host, repo_server_port) - notifications_controller_endpoint = 'http://{}:{}/metrics'.format( - notifications_controller_host, notifications_controller_port - ) - - instance = { - 'app_controller_endpoint': app_controller_endpoint, - 'appset_controller_endpoint': appset_controller_endpoint, - 'api_server_endpoint': api_server_endpoint, - 'repo_server_endpoint': repo_server_endpoint, - 'notifications_controller_endpoint': notifications_controller_endpoint, - } - - # save this instance to use for openmetrics_v2 instance, since the endpoint is different each run - dd_save_state("argocd_instance", instance) - - yield instance + instance = { + 'app_controller_endpoint': 'http://argocd-metrics.argocd.svc.cluster.local:8082/metrics', + 'appset_controller_endpoint': ( + 'http://argocd-applicationset-controller.argocd.svc.cluster.local:8080/metrics' + ), + 'api_server_endpoint': 'http://argocd-server-metrics.argocd.svc.cluster.local:8083/metrics', + 'repo_server_endpoint': 'http://argocd-repo-server.argocd.svc.cluster.local:8084/metrics', + 'notifications_controller_endpoint': ( + 'http://argocd-notifications-controller-metrics.argocd.svc.cluster.local:9001/metrics' + ), + } + metadata = { + 'agent_type': 'kubernetes', + 'kubernetes': { + 'kubeconfig': kubeconfig, + }, + } + + # Save this instance to use for the openmetrics_v2 instance. + dd_save_state("argocd_instance", instance) + + yield instance, metadata diff --git a/calico/tests/conftest.py b/calico/tests/conftest.py index 7f7453240259c..58679e1beb528 100644 --- a/calico/tests/conftest.py +++ b/calico/tests/conftest.py @@ -5,15 +5,15 @@ import pytest -from datadog_checks.dev.conditions import CheckEndpoints, WaitFor +from datadog_checks.dev.conditions import WaitFor from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward from datadog_checks.dev.subprocess import run_command from .common import EXTRA_METRICS NAMESPACE = "calico" HERE = path.dirname(path.abspath(__file__)) +FELIX_METRICS_ENDPOINT = 'http://felix-metrics-svc.kube-system.svc.cluster.local:9091/metrics' def _felix_config_default_exists(): @@ -21,6 +21,33 @@ def _felix_config_default_exists(): return result.code == 0 +def felix_metrics_available() -> bool: + result = run_command( + [ + "kubectl", + "run", + "felix-metrics-readiness", + "--namespace", + "kube-system", + "--image=busybox:1.36.1", + "--restart=Never", + "--attach", + "--rm", + "--quiet", + "--", + "wget", + "-q", + "-T", + "2", + "-O", + "/dev/null", + FELIX_METRICS_ENDPOINT, + ], + capture='both', + ) + return result.code == 0 + + def setup_calico(): # Deploy calico run_command(["kubectl", "apply", "-f", path.join(HERE, 'kind', 'calico.yaml')]) @@ -60,29 +87,23 @@ def setup_calico(): check=True, ) + # Check from a temporary pod because the host cannot resolve Kubernetes Service DNS. + WaitFor(felix_metrics_available, wait=2, attempts=100)() + @pytest.fixture(scope='session') def dd_environment(): - with ( - kind_run( - conditions=[setup_calico], kind_config=path.join(HERE, 'kind', 'kind-calico.yaml'), sleep=10 - ) as kubeconfig, - port_forward(kubeconfig, 'kube-system', 9091, 'service', 'felix-metrics-svc') as ( - calico_host, - calico_port, - ), - ): - endpoint = 'http://{}:{}/metrics'.format(calico_host, calico_port) - - # We can't add this to `kind_run` because we don't know the URL at this moment - condition = CheckEndpoints(endpoint, wait=2, attempts=100) - condition() - - yield { - "openmetrics_endpoint": endpoint, + with kind_run( + conditions=[setup_calico], kind_config=path.join(HERE, 'kind', 'kind-calico.yaml'), sleep=10 + ) as kubeconfig: + instance = { + "openmetrics_endpoint": FELIX_METRICS_ENDPOINT, "namespace": NAMESPACE, "extra_metrics": EXTRA_METRICS, } + metadata = {'agent_type': 'kubernetes', 'kubernetes': {'kubeconfig': kubeconfig}} + + yield instance, metadata @pytest.fixture diff --git a/datadog_checks_dev/changelog.d/24680.added b/datadog_checks_dev/changelog.d/24680.added new file mode 100644 index 0000000000000..e2df056d40f0f --- /dev/null +++ b/datadog_checks_dev/changelog.d/24680.added @@ -0,0 +1 @@ +Add backend-neutral Agent log retrieval for E2E test diagnostics. diff --git a/datadog_checks_dev/datadog_checks/dev/plugin/pytest.py b/datadog_checks_dev/datadog_checks/dev/plugin/pytest.py index 5ad3e146c9234..9980634edc1d2 100644 --- a/datadog_checks_dev/datadog_checks/dev/plugin/pytest.py +++ b/datadog_checks_dev/datadog_checks/dev/plugin/pytest.py @@ -219,7 +219,10 @@ def run_check(config=None, **kwargs): if not matches: message_parts = [] - debug_result = run_command(['docker', 'logs', 'dd_{}_{}'.format(check, env)], capture=True) + debug_result = run_command( + [python_path, '-m', 'ddev', 'env', 'logs', check, env], + capture=True, + ) if not debug_result.code: message_parts.append(debug_result.stdout + debug_result.stderr) diff --git a/ddev/changelog.d/24680.added b/ddev/changelog.d/24680.added new file mode 100644 index 0000000000000..d811757c468b6 --- /dev/null +++ b/ddev/changelog.d/24680.added @@ -0,0 +1 @@ +Add a Kubernetes Agent interface for running the Agent in Kind-based E2E test clusters, alongside the existing Docker and Vagrant interfaces. diff --git a/ddev/src/ddev/cli/env/__init__.py b/ddev/src/ddev/cli/env/__init__.py index 6f38d7bf9730c..769a4079ef916 100644 --- a/ddev/src/ddev/cli/env/__init__.py +++ b/ddev/src/ddev/cli/env/__init__.py @@ -6,6 +6,7 @@ from ddev.cli.env.agent import agent from ddev.cli.env.check import check from ddev.cli.env.config import config +from ddev.cli.env.logs import logs from ddev.cli.env.reload import reload_command from ddev.cli.env.shell import shell from ddev.cli.env.show import show @@ -24,6 +25,7 @@ def env(): env.add_command(agent) env.add_command(check) env.add_command(config) +env.add_command(logs) env.add_command(reload_command) env.add_command(shell) env.add_command(show) diff --git a/ddev/src/ddev/cli/env/agent.py b/ddev/src/ddev/cli/env/agent.py index d67192f00e03e..864c2e758ec56 100644 --- a/ddev/src/ddev/cli/env/agent.py +++ b/ddev/src/ddev/cli/env/agent.py @@ -49,6 +49,18 @@ def _validate_env_vars(ctx: click.Context, param: click.Parameter, value: tuple[ return env_vars or None +def _sync_restored_config(app: Application, agent: AgentInterface) -> None: + import sys + + original_error = sys.exception() + try: + agent.sync_config() + except Exception as e: + if original_error is None: + raise + app.display_warning(f'Unable to restore the Agent configuration: {e}') + + @click.command( short_help='Invoke the Agent', context_settings={'help_option_names': [], 'ignore_unknown_options': True} ) @@ -79,9 +91,8 @@ def agent( """ import subprocess - from ddev.e2e.agent import get_agent_interface + from ddev.e2e.agent import create_agent_interface from ddev.e2e.config import EnvDataStorage - from ddev.e2e.constants import DEFAULT_AGENT_TYPE, E2EMetadata from ddev.utils.fs import Path integration = app.repo.integrations.get(intg_name) @@ -91,8 +102,7 @@ def agent( app.abort(f'Environment `{environment}` for integration `{integration.name}` is not running') metadata = env_data.read_metadata() - agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) - agent = get_agent_interface(agent_type)(app, integration, environment, metadata, env_data.config_file) + agent = create_agent_interface(app, integration, environment, metadata, env_data.config_file) full_args = list(args) trigger_run = False @@ -131,6 +141,7 @@ def agent( app.abort(str(e)) finally: env_data.config_file.unlink() + _sync_restored_config(app, agent) else: temp_config_file = env_data.config_file.parent / f'{env_data.config_file.name}.bak.example' env_data.config_file.replace(temp_config_file) @@ -141,3 +152,4 @@ def agent( app.abort(str(e)) finally: temp_config_file.replace(env_data.config_file) + _sync_restored_config(app, agent) diff --git a/ddev/src/ddev/cli/env/logs.py b/ddev/src/ddev/cli/env/logs.py new file mode 100644 index 0000000000000..c72c3d5280278 --- /dev/null +++ b/ddev/src/ddev/cli/env/logs.py @@ -0,0 +1,35 @@ +# (C) Datadog, Inc. 2026-present +# All rights reserved +# Licensed under a 3-clause BSD style license (see LICENSE) +from __future__ import annotations + +from typing import TYPE_CHECKING + +import click + +if TYPE_CHECKING: + from ddev.cli.application import Application + + +@click.command('logs', short_help='Show logs for the Agent') +@click.argument('intg_name', metavar='INTEGRATION') +@click.argument('environment') +@click.pass_obj +def logs(app: Application, *, intg_name: str, environment: str): + """Show backend-specific diagnostics for the Agent.""" + from ddev.e2e.agent import create_agent_interface + from ddev.e2e.config import EnvDataStorage + + integration = app.repo.integrations.get(intg_name) + env_data = EnvDataStorage(app.data_dir).get(integration.name, environment) + + if not env_data.exists(): + app.abort(f'Environment `{environment}` for integration `{integration.name}` is not running') + + metadata = env_data.read_metadata() + agent = create_agent_interface(app, integration, environment, metadata, env_data.config_file) + + try: + agent.show_logs() + except Exception as e: + app.abort(str(e)) diff --git a/ddev/src/ddev/cli/env/reload.py b/ddev/src/ddev/cli/env/reload.py index 284c8c1e50563..e302c4f33fc29 100644 --- a/ddev/src/ddev/cli/env/reload.py +++ b/ddev/src/ddev/cli/env/reload.py @@ -19,9 +19,8 @@ def reload_command(app: Application, *, intg_name: str, environment: str): """ Restart the Agent to detect environment changes. """ - from ddev.e2e.agent import get_agent_interface + from ddev.e2e.agent import create_agent_interface from ddev.e2e.config import EnvDataStorage - from ddev.e2e.constants import DEFAULT_AGENT_TYPE, E2EMetadata integration = app.repo.integrations.get(intg_name) env_data = EnvDataStorage(app.data_dir).get(integration.name, environment) @@ -30,8 +29,7 @@ def reload_command(app: Application, *, intg_name: str, environment: str): app.abort(f'Environment `{environment}` for integration `{integration.name}` is not running') metadata = env_data.read_metadata() - agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) - agent = get_agent_interface(agent_type)(app, integration, environment, metadata, env_data.config_file) + agent = create_agent_interface(app, integration, environment, metadata, env_data.config_file) try: agent.restart() diff --git a/ddev/src/ddev/cli/env/shell.py b/ddev/src/ddev/cli/env/shell.py index 21ce3f864cf89..7a4bc5947a70b 100644 --- a/ddev/src/ddev/cli/env/shell.py +++ b/ddev/src/ddev/cli/env/shell.py @@ -21,9 +21,8 @@ def shell(app: Application, *, intg_name: str, environment: str): """ import subprocess - from ddev.e2e.agent import get_agent_interface + from ddev.e2e.agent import create_agent_interface from ddev.e2e.config import EnvDataStorage - from ddev.e2e.constants import DEFAULT_AGENT_TYPE, E2EMetadata integration = app.repo.integrations.get(intg_name) env_data = EnvDataStorage(app.data_dir).get(integration.name, environment) @@ -32,8 +31,7 @@ def shell(app: Application, *, intg_name: str, environment: str): app.abort(f'Environment `{environment}` for integration `{integration.name}` is not running') metadata = env_data.read_metadata() - agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) - agent = get_agent_interface(agent_type)(app, integration, environment, metadata, env_data.config_file) + agent = create_agent_interface(app, integration, environment, metadata, env_data.config_file) try: agent.enter_shell() diff --git a/ddev/src/ddev/cli/env/show.py b/ddev/src/ddev/cli/env/show.py index fca3946d8d23a..5d50c5ec39821 100644 --- a/ddev/src/ddev/cli/env/show.py +++ b/ddev/src/ddev/cli/env/show.py @@ -72,7 +72,7 @@ def show(app: Application, *, intg_name: str | None, environment: str | None, fo app.display_table('Available', available_columns, show_lines=True, force_ascii=force_ascii) # Display information about a specific environment else: - from ddev.e2e.agent import get_agent_interface + from ddev.e2e.agent import create_agent_interface integration = app.repo.integrations.get(intg_name) env_data = storage.get(integration.name, environment) @@ -82,7 +82,7 @@ def show(app: Application, *, intg_name: str | None, environment: str | None, fo metadata = env_data.read_metadata() agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) - agent = get_agent_interface(agent_type)(app, integration, environment, metadata, env_data.config_file) + agent = create_agent_interface(app, integration, environment, metadata, env_data.config_file) app.display_pair('Agent type', agent_type) app.display_pair('Agent ID', agent.get_id()) diff --git a/ddev/src/ddev/cli/env/start.py b/ddev/src/ddev/cli/env/start.py index fd0c1577c8cfb..0a76d849d3f32 100644 --- a/ddev/src/ddev/cli/env/start.py +++ b/ddev/src/ddev/cli/env/start.py @@ -132,6 +132,7 @@ def start( result = json.loads(result_file.read_text()) metadata = result['metadata'] + agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) # TODO Remove once we have migrated the `docker_run` function if serialized_volumes := metadata.get(E2EMetadata.ENV_VARS, {}).get(E2EEnvVars.DOCKER_VOLUMES): @@ -144,17 +145,17 @@ def start( config = result['config'] env_data.write_config(config) - agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) - - if agent_type == "vagrant" and running_on_ci(): - app.abort(text="Vagrant is not supported on CI", code=0) + agent_class = get_agent_interface(agent_type) + if running_on_ci() and not agent_class.supports_ci: + app.abort(text=f'{agent_type.capitalize()} is not supported on CI', code=0) - agent = get_agent_interface(agent_type)(app, integration, environment, metadata, env_data.config_file) + agent = agent_class(app, integration, environment, metadata, env_data.config_file) if not agent_build: + configured_agent_build = agent.get_configured_build(app.config.agent.config) agent_build = ( os.getenv(E2EEnvVars.AGENT_BUILD_PY2 if agent.python_version[0] == 2 else E2EEnvVars.AGENT_BUILD) - or app.config.agent.config.get(agent_type) + or configured_agent_build or '' ) @@ -162,6 +163,8 @@ def start( try: agent.start(agent_build=agent_build, local_packages=local_packages, env_vars=agent_env_vars) + # Backends may add runtime metadata needed by later ddev processes. + env_data.write_metadata(metadata) except Exception as e: from ddev.cli.env.stop import stop diff --git a/ddev/src/ddev/cli/env/stop.py b/ddev/src/ddev/cli/env/stop.py index 62e8e3a19266f..a73632e51bec0 100644 --- a/ddev/src/ddev/cli/env/stop.py +++ b/ddev/src/ddev/cli/env/stop.py @@ -20,9 +20,9 @@ def stop(app: Application, *, intg_name: str, environment: str, ignore_state: bo """ Stop environments. To stop all the running environments, use `all` as the integration name and the environment. """ - from ddev.e2e.agent import get_agent_interface + from ddev.e2e.agent import create_agent_interface from ddev.e2e.config import EnvDataStorage - from ddev.e2e.constants import DEFAULT_AGENT_TYPE, E2EEnvVars, E2EMetadata + from ddev.e2e.constants import E2EEnvVars, E2EMetadata from ddev.e2e.run import E2EEnvironmentRunner from ddev.utils.fs import temp_directory @@ -63,8 +63,7 @@ def stop(app: Application, *, intg_name: str, environment: str, ignore_state: bo metadata = env_data.read_metadata() env_vars.update(metadata.get(E2EMetadata.ENV_VARS, {})) - agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) - agent = get_agent_interface(agent_type)(app, integration, env, metadata, env_data.config_file) + agent = create_agent_interface(app, integration, env, metadata, env_data.config_file) try: agent.stop() diff --git a/ddev/src/ddev/e2e/agent/__init__.py b/ddev/src/ddev/e2e/agent/__init__.py index f0dcfe5a29556..522d5abfc7553 100644 --- a/ddev/src/ddev/e2e/agent/__init__.py +++ b/ddev/src/ddev/e2e/agent/__init__.py @@ -3,10 +3,26 @@ # Licensed under a 3-clause BSD style license (see LICENSE) from __future__ import annotations -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, Any if TYPE_CHECKING: + from ddev.cli.application import Application from ddev.e2e.agent.interface import AgentInterface + from ddev.integration.core import Integration + from ddev.utils.fs import Path + + +def create_agent_interface( + app: Application, + integration: Integration, + environment: str, + metadata: dict[str, Any], + config_file: Path, +) -> AgentInterface: + from ddev.e2e.constants import DEFAULT_AGENT_TYPE, E2EMetadata + + agent_type = metadata.get(E2EMetadata.AGENT_TYPE, DEFAULT_AGENT_TYPE) + return get_agent_interface(agent_type)(app, integration, environment, metadata, config_file) def get_agent_interface(agent_type: str) -> type[AgentInterface]: @@ -20,4 +36,9 @@ def get_agent_interface(agent_type: str) -> type[AgentInterface]: return VagrantAgent + if agent_type == "kubernetes": + from ddev.e2e.agent.kubernetes import KubernetesAgent + + return KubernetesAgent + raise NotImplementedError(f"Unsupported Agent type: {agent_type}") diff --git a/ddev/src/ddev/e2e/agent/docker.py b/ddev/src/ddev/e2e/agent/docker.py index a362b0f92f19d..a3f4978a6b496 100644 --- a/ddev/src/ddev/e2e/agent/docker.py +++ b/ddev/src/ddev/e2e/agent/docker.py @@ -4,7 +4,6 @@ from __future__ import annotations import os -import re import sys from contextlib import AbstractContextManager, contextmanager, nullcontext from functools import cache, cached_property, partial @@ -12,6 +11,7 @@ import stamina +from ddev.e2e.agent.image import normalize_agent_image_name from ddev.e2e.agent.interface import AgentInterface from ddev.utils.structures import EnvVars @@ -20,16 +20,6 @@ from ddev.utils.fs import Path -AGENT_IMAGE_REGEX = r'^([^/]+)/([^:]+):(.*)$' -AGENT_VERSION_REGEX = ( - # Main version: 7, 7.69, 7.69.0 ... - r"^(?P\d+(?:\.[\dx]+)*|latest|main|master|nightly)" - # rcs: rc.1, rc ... - r"(?P-rc(?:\.\d+)?)?" - # Anny suffixes: -jmx, -linux, -full... - r"(?P(?:-[a-zA-Z0-9]+)*)$" -) - @contextmanager def disable_integration_before_install(config_file): @@ -45,43 +35,9 @@ def disable_integration_before_install(config_file): new.rename(config_file.parent / old) -def _normalize_agent_image_name(agent_build: str | None, python_major: int, use_jmx: bool) -> str: - if not agent_build: - return 'registry.datadoghq.com/agent-dev:master-py3' - - if match := re.match(AGENT_IMAGE_REGEX, agent_build): - org, image, tag = match.groups() - - if org != 'datadog' and org != 'registry.datadoghq.com': - # Some non datadog image has been selected - return agent_build - - version_match = re.match(AGENT_VERSION_REGEX, tag) - - if version_match is None: - # Not sure how to extract information for a version of this shape - return agent_build - - version = version_match.group('version') - rc = version_match.group('rc') - rc = rc if rc else '' - suffixes = version_match.group('suffixes') - - # Add -py suffix if missing, only in agent-dev the agent does not have py3 suffix after agent 6 - if image == 'agent-dev': - if not (rc != '' or any(suffix in suffixes for suffix in ['py', 'fips']) or version[0].isdigit()): - suffixes = f'-py{python_major}{suffixes}' - - # Add jmx suffix if missing - if use_jmx and '-jmx' not in suffixes: - suffixes += '-jmx' - - return f'{org}/{image}:{version}{rc}{suffixes}' - - return agent_build - - class DockerAgent(AgentInterface): + build_config_key = 'docker' + @cached_property def _isatty(self) -> bool: isatty: Callable[[], bool] | None = getattr(sys.stdout, 'isatty', None) @@ -173,8 +129,11 @@ def _execute_lifecycle_command( self._show_logs() raise RuntimeError(error_message) - def _show_logs(self) -> None: - self._run_command(['docker', 'logs', self._container_name]) + def _show_logs(self, *, check: bool = False) -> None: + self._run_command(['docker', 'logs', self._container_name], check=check) + + def show_logs(self) -> None: + self._show_logs(check=True) def get_id(self) -> str: return self._container_name @@ -182,7 +141,7 @@ def get_id(self) -> str: def start(self, *, agent_build: str | None, local_packages: dict[Path, str], env_vars: dict[str, str]) -> None: from ddev.e2e.agent.constants import AgentEnvVars - agent_build = _normalize_agent_image_name( + agent_build = normalize_agent_image_name( agent_build, self.python_version[0], self.metadata.get('use_jmx', False) ) diff --git a/ddev/src/ddev/e2e/agent/image.py b/ddev/src/ddev/e2e/agent/image.py new file mode 100644 index 0000000000000..6648c09e61904 --- /dev/null +++ b/ddev/src/ddev/e2e/agent/image.py @@ -0,0 +1,50 @@ +# (C) Datadog, Inc. 2026-present +# All rights reserved +# Licensed under a 3-clause BSD style license (see LICENSE) +from __future__ import annotations + +import re + +AGENT_IMAGE_REGEX = r'^([^/]+)/([^:]+):(.*)$' +AGENT_VERSION_REGEX = ( + # Main version: 7, 7.69, 7.69.0 ... + r"^(?P\d+(?:\.[\dx]+)*|latest|main|master|nightly)" + # rcs: rc.1, rc ... + r"(?P-rc(?:\.\d+)?)?" + # Any suffixes: -jmx, -linux, -full... + r"(?P(?:-[a-zA-Z0-9]+)*)$" +) + + +def normalize_agent_image_name(agent_build: str | None, python_major: int, use_jmx: bool) -> str: + if not agent_build: + return 'registry.datadoghq.com/agent-dev:master-py3' + + if match := re.match(AGENT_IMAGE_REGEX, agent_build): + org, image, tag = match.groups() + + if org not in {'datadog', 'registry.datadoghq.com'}: + # Some non-Datadog image has been selected. + return agent_build + + version_match = re.match(AGENT_VERSION_REGEX, tag) + if version_match is None: + # The tag does not follow a recognized Agent version format. + return agent_build + + version = version_match.group('version') + rc = version_match.group('rc') or '' + suffixes = version_match.group('suffixes') + + # Add a Python suffix when required by a development Agent image. + if image == 'agent-dev': + has_python_variant = any(suffix in suffixes for suffix in ('py', 'fips')) + if not (rc or has_python_variant or version[0].isdigit()): + suffixes = f'-py{python_major}{suffixes}' + + if use_jmx and '-jmx' not in suffixes: + suffixes += '-jmx' + + return f'{org}/{image}:{version}{rc}{suffixes}' + + return agent_build diff --git a/ddev/src/ddev/e2e/agent/interface.py b/ddev/src/ddev/e2e/agent/interface.py index 135d832a5e6d6..13b0a8089bd67 100644 --- a/ddev/src/ddev/e2e/agent/interface.py +++ b/ddev/src/ddev/e2e/agent/interface.py @@ -8,6 +8,8 @@ from typing import TYPE_CHECKING, Any if TYPE_CHECKING: + from collections.abc import Mapping + from ddev.cli.application import Application from ddev.integration.core import Integration from ddev.utils.fs import Path @@ -15,6 +17,9 @@ class AgentInterface(ABC): + build_config_key: str | None = None + supports_ci = True + def __init__( self, app: Application, integration: Integration, env: str, metadata: dict[str, Any], config_file: Path ) -> None: @@ -63,6 +68,12 @@ def python_version(self) -> tuple[int, int]: def get_id(self) -> str: return f'{self.integration.name}_{self.env}' + def get_configured_build(self, config: Mapping[str, str]) -> str | None: + if self.build_config_key is None: + return None + + return config.get(self.build_config_key) + @abstractmethod def start(self, *, agent_build: str, local_packages: dict[Path, str], env_vars: dict[str, str]) -> None: ... @@ -77,3 +88,10 @@ def invoke(self, args: list[str], *, env_vars: dict[str, str] | None = None) -> @abstractmethod def enter_shell(self) -> None: ... + + def sync_config(self) -> None: + """Synchronize the persisted host configuration with the Agent.""" + + def show_logs(self) -> None: + """Show backend-specific diagnostics for the running Agent.""" + self.invoke(['status']) diff --git a/ddev/src/ddev/e2e/agent/kubernetes.py b/ddev/src/ddev/e2e/agent/kubernetes.py new file mode 100644 index 0000000000000..545d14786b806 --- /dev/null +++ b/ddev/src/ddev/e2e/agent/kubernetes.py @@ -0,0 +1,459 @@ +# (C) Datadog, Inc. 2026-present +# All rights reserved +# Licensed under a 3-clause BSD style license (see LICENSE) +from __future__ import annotations + +import json +import time +from typing import TYPE_CHECKING, Any, cast + +from ddev.e2e.agent.constants import AgentEnvVars +from ddev.e2e.agent.image import normalize_agent_image_name +from ddev.e2e.agent.interface import AgentInterface + +if TYPE_CHECKING: + import subprocess + + from ddev.utils.fs import Path + +POD_NAME = 'ddev-agent' +CONTAINER_NAME = 'agent' +NAMESPACE = 'ddev-agent' +CLUSTER_RESOURCE_NAME = 'ddev-agent' +DEFAULT_WAIT_TIMEOUT = 120 +LOCAL_PACKAGES_METADATA = 'local_packages' +PREPARED_MARKER = '/home/.ddev-agent-prepared' + + +class KubernetesAgent(AgentInterface): + """Run the E2E Agent in a disposable, fixture-owned Kubernetes cluster. + + The backend currently supports clusters provisioned by ``kind_run``. It requires + exactly one schedulable node and runs one Agent pod in the cluster. + """ + + build_config_key = 'docker' + + @property + def _kubernetes_metadata(self) -> dict[str, Any]: + metadata = self.metadata.get('kubernetes') + if not isinstance(metadata, dict): + raise ValueError( + 'Agent type `kubernetes` requires a `kubernetes` metadata mapping with a `kubeconfig` path' + ) + return cast(dict[str, Any], metadata) + + @property + def _kubeconfig(self) -> str: + kubeconfig = self._kubernetes_metadata.get('kubeconfig') + if not isinstance(kubeconfig, str) or not kubeconfig: + raise ValueError('Kubernetes Agent metadata must define a non-empty `kubeconfig` path') + return kubeconfig + + @property + def _namespace(self) -> str: + return NAMESPACE + + @property + def _cluster_resource_name(self) -> str: + return CLUSTER_RESOURCE_NAME + + @property + def _resource_labels(self) -> dict[str, str]: + return {'app.kubernetes.io/managed-by': 'ddev'} + + @property + def _config_dir(self) -> str: + return f'/etc/datadog-agent/conf.d/{self.integration.name}.d' + + @property + def _python_path(self) -> str: + return f'/opt/datadog-agent/embedded/bin/python{self.python_version[0]}' + + @property + def _kubectl_prefix(self) -> list[str]: + return ['kubectl', '--kubeconfig', self._kubeconfig] + + @property + def _wait_timeout(self) -> int: + timeout = self._kubernetes_metadata.get('wait_timeout', DEFAULT_WAIT_TIMEOUT) + if isinstance(timeout, bool) or not isinstance(timeout, int) or timeout <= 0: + raise ValueError('Kubernetes Agent `wait_timeout` must be a positive integer') + return timeout + + def _kubectl(self, args: list[str], **kwargs) -> subprocess.CompletedProcess: + return self.platform.run_command([*self._kubectl_prefix, *args], **kwargs) + + def _captured_kubectl(self, args: list[str], *, merge_stderr: bool = True, **kwargs) -> subprocess.CompletedProcess: + return self._kubectl( + args, + stdout=self.platform.modules.subprocess.PIPE, + stderr=(self.platform.modules.subprocess.STDOUT if merge_stderr else self.platform.modules.subprocess.PIPE), + **kwargs, + ) + + @staticmethod + def _process_output(process: subprocess.CompletedProcess, *, include_stderr: bool = False) -> str: + streams = (process.stdout, process.stderr) if include_stderr else (process.stdout,) + return ''.join( + stream.decode('utf-8', errors='replace') if isinstance(stream, bytes) else stream or '' + for stream in streams + ) + + def _exec( + self, + command: list[str], + *, + env_vars: dict[str, str] | None = None, + check: bool = True, + capture: bool = False, + ) -> subprocess.CompletedProcess: + args = ['exec', '--namespace', self._namespace, f'pod/{POD_NAME}', '--container', CONTAINER_NAME, '--'] + if env_vars: + args.append('env') + args.extend(f'{key}={value}' for key, value in sorted(env_vars.items())) + args.extend(command) + if capture: + return self._captured_kubectl(args, check=check) + return self._kubectl(args, check=check) + + def _validate_context(self) -> None: + process = self._captured_kubectl(['config', 'current-context'], merge_stderr=False) + if process.returncode: + raise RuntimeError( + f'Unable to inspect Kubernetes context: {self._process_output(process, include_stderr=True)}' + ) + + context = self._process_output(process).strip() + if not context.startswith('kind-'): + raise RuntimeError(f'Refusing to use non-Kind Kubernetes context `{context}`') + + def _validate_topology(self) -> None: + process = self._captured_kubectl(['get', 'nodes', '-o', 'json'], merge_stderr=False) + if process.returncode: + raise RuntimeError( + f'Unable to inspect Kubernetes nodes: {self._process_output(process, include_stderr=True)}' + ) + + try: + node_data = json.loads(self._process_output(process)) + schedulable_nodes = [node for node in node_data['items'] if not node.get('spec', {}).get('unschedulable')] + except (KeyError, TypeError, json.JSONDecodeError) as e: + raise RuntimeError(f'Unable to parse Kubernetes node data: {e}') from e + + if len(schedulable_nodes) != 1: + raise NotImplementedError( + 'KubernetesAgent currently requires exactly one schedulable node; ' + f'found {len(schedulable_nodes)}. Multi-node execution needs an explicit Agent targeting policy.' + ) + + def _manifest(self, agent_build: str, env_vars: dict[str, str]) -> dict[str, Any]: + env_vars = env_vars.copy() + env_vars.setdefault('DD_API_KEY', 'a' * 32) + env_vars.setdefault('DD_APM_ENABLED', 'false') + env_vars.setdefault('DD_AUTOCONFIG_FROM_ENVIRONMENT', 'true') + env_vars.setdefault('DD_HOSTNAME', self._namespace) + env_vars.setdefault('DD_KUBELET_TLS_VERIFY', 'false') + + container_env: list[dict[str, Any]] = [{'name': key, 'value': value} for key, value in sorted(env_vars.items())] + container_env.extend( + [ + { + 'name': 'DD_KUBERNETES_KUBELET_HOST', + 'valueFrom': {'fieldRef': {'fieldPath': 'status.hostIP'}}, + }, + { + 'name': 'DD_KUBERNETES_KUBELET_NODENAME', + 'valueFrom': {'fieldRef': {'fieldPath': 'spec.nodeName'}}, + }, + ] + ) + + labels = {**self._resource_labels, 'app.kubernetes.io/name': POD_NAME} + + image_pull_policy = self._kubernetes_metadata.get('image_pull_policy', 'Always') + if image_pull_policy not in {'Always', 'IfNotPresent', 'Never'}: + raise ValueError('Kubernetes Agent `image_pull_policy` must be Always, IfNotPresent, or Never') + + role_rules = [ + {'apiGroups': [''], 'resources': ['nodes'], 'verbs': ['get', 'list', 'watch']}, + { + 'apiGroups': [''], + 'resources': ['nodes/metrics', 'nodes/spec', 'nodes/stats', 'nodes/proxy'], + 'verbs': ['get'], + }, + { + 'apiGroups': [''], + 'resources': ['pods', 'endpoints', 'services'], + 'verbs': ['get', 'list', 'watch'], + }, + ] + + return { + 'apiVersion': 'v1', + 'kind': 'List', + 'items': [ + { + 'apiVersion': 'v1', + 'kind': 'Namespace', + 'metadata': {'name': self._namespace, 'labels': self._resource_labels}, + }, + { + 'apiVersion': 'v1', + 'kind': 'ServiceAccount', + 'metadata': {'name': POD_NAME, 'namespace': self._namespace, 'labels': self._resource_labels}, + }, + { + 'apiVersion': 'rbac.authorization.k8s.io/v1', + 'kind': 'ClusterRole', + 'metadata': {'name': self._cluster_resource_name, 'labels': self._resource_labels}, + 'rules': role_rules, + }, + { + 'apiVersion': 'rbac.authorization.k8s.io/v1', + 'kind': 'ClusterRoleBinding', + 'metadata': {'name': self._cluster_resource_name, 'labels': self._resource_labels}, + 'roleRef': { + 'apiGroup': 'rbac.authorization.k8s.io', + 'kind': 'ClusterRole', + 'name': self._cluster_resource_name, + }, + 'subjects': [ + {'kind': 'ServiceAccount', 'name': POD_NAME, 'namespace': self._namespace}, + ], + }, + { + 'apiVersion': 'v1', + 'kind': 'Pod', + 'metadata': {'name': POD_NAME, 'namespace': self._namespace, 'labels': labels}, + 'spec': { + 'serviceAccountName': POD_NAME, + 'restartPolicy': 'Always', + 'terminationGracePeriodSeconds': 0, + 'tolerations': [{'operator': 'Exists'}], + 'containers': [ + { + 'name': CONTAINER_NAME, + 'image': agent_build, + 'imagePullPolicy': image_pull_policy, + 'env': container_env, + } + ], + }, + }, + ], + } + + def _create_manifest(self, manifest: dict[str, Any]) -> None: + process = self._captured_kubectl(['create', '-f', '-'], input=json.dumps(manifest).encode()) + if process.returncode: + raise RuntimeError(f'Unable to create Kubernetes Agent resources: {self._process_output(process)}') + + def _wait_for_pod(self) -> None: + process = self._captured_kubectl( + [ + 'wait', + '--namespace', + self._namespace, + '--for=condition=Ready', + f'pod/{POD_NAME}', + f'--timeout={self._wait_timeout}s', + ] + ) + if process.returncode: + self._show_logs() + raise RuntimeError(f'Kubernetes Agent pod did not become ready: {self._process_output(process)}') + + def _wait_for_agent(self) -> None: + deadline = time.monotonic() + self._wait_timeout + last_output = '' + while time.monotonic() < deadline: + process = self._exec(['agent', 'status'], check=False, capture=True) + if process.returncode == 0: + return + last_output = self._process_output(process) + time.sleep(1) + + self._show_logs() + raise RuntimeError(f'Kubernetes Agent did not become ready: {last_output}') + + def _copy_file(self, source: str, destination: str) -> None: + from ddev.utils.fs import Path + + source_path = Path(source).resolve() + self._kubectl( + [ + 'cp', + '--container', + CONTAINER_NAME, + source_path.name, + f'{self._namespace}/{POD_NAME}:{destination}', + ], + check=True, + cwd=source_path.parent, + ) + + def _sync_config(self) -> None: + self._exec(['mkdir', '-p', self._config_dir]) + destination = f'{self._config_dir}/conf.yaml' + if self.config_file.is_file(): + self._copy_file(str(self.config_file), destination) + else: + self._exec(['rm', '-f', destination]) + + def _sync_auto_conf(self) -> None: + auto_conf = self._kubernetes_metadata.get('auto_conf') + if auto_conf is None: + return + if not isinstance(auto_conf, str) or not auto_conf: + raise ValueError('Kubernetes Agent `auto_conf` must be a non-empty path') + self._exec(['mkdir', '-p', self._config_dir]) + self._copy_file(auto_conf, f'{self._config_dir}/auto_conf.yaml') + + def _remember_local_packages(self, local_packages: dict[Path, str]) -> None: + self._kubernetes_metadata[LOCAL_PACKAGES_METADATA] = { + str(local_package): features for local_package, features in local_packages.items() + } + + def _local_packages(self) -> dict[Path, str]: + from ddev.utils.fs import Path + + local_packages = cast(dict[str, str], self._kubernetes_metadata.get(LOCAL_PACKAGES_METADATA, {})) + return {Path(path): features for path, features in local_packages.items()} + + def _sync_local_packages(self, *, install: bool = False) -> None: + for local_package, features in self._local_packages().items(): + name = local_package.name + if not name or name in {'.', '..'}: + raise ValueError(f'Invalid Kubernetes Agent local package path: {local_package}') + + destination = f'/home/{name}' + self._exec(['rm', '-rf', destination]) + self._copy_file(str(local_package), destination) + if install: + self._exec( + [ + self._python_path, + '-m', + 'pip', + 'install', + '--disable-pip-version-check', + '-e', + f'{destination}{features}', + ] + ) + + def _run_metadata_commands(self, key: str) -> None: + commands = self.metadata.get(key, []) + if not isinstance(commands, list) or not all(isinstance(command, str) for command in commands): + raise ValueError(f'Kubernetes Agent `{key}` must be a list of commands') + for command in commands: + self._exec(self.platform.modules.shlex.split(command)) + + @staticmethod + def _validate_options(env_vars: dict[str, str]) -> None: + if AgentEnvVars.DOGSTATSD_PORT in env_vars: + raise NotImplementedError('The Kubernetes Agent backend does not currently support DogStatsD exposure') + if env_vars.get(AgentEnvVars.LOGS_ENABLED, '').lower() == 'true': + raise NotImplementedError( + 'The Kubernetes Agent backend does not currently support shared-file log collection' + ) + + def _require_prepared(self) -> None: + process = self._exec(['test', '-f', PREPARED_MARKER], check=False) + if process.returncode: + raise RuntimeError( + 'The Kubernetes Agent container is no longer prepared and may have restarted. ' + 'Recreate the environment with `ddev env stop` followed by `ddev env start`.' + ) + + def start(self, *, agent_build: str | None, local_packages: dict[Path, str], env_vars: dict[str, str]) -> None: + self._validate_options(env_vars) + agent_build = normalize_agent_image_name( + agent_build, self.python_version[0], self.metadata.get('use_jmx', False) + ) + _ = self._wait_timeout + self._validate_context() + self._validate_topology() + self._create_manifest(self._manifest(agent_build, env_vars)) + self._wait_for_pod() + self._run_metadata_commands('start_commands') + self._remember_local_packages(local_packages) + self._sync_local_packages(install=True) + self._sync_config() + self._sync_auto_conf() + self._run_metadata_commands('post_install_commands') + # A replacement during restart must lose this marker so the + # post-restart check rejects the unprepared container. + self._exec(['touch', PREPARED_MARKER]) + self._restart_agent_process() + self._require_prepared() + + def stop(self) -> None: + """Leave cleanup to the fixture that deletes the disposable cluster.""" + + def _restart_agent_process(self) -> None: + # Pod readiness does not guarantee that s6 has started the Agent process. + self._wait_for_agent() + + # The Agent image's s6 finish handler normally shuts down the whole + # service tree when the main Agent exits. Remove it before killing the + # process so s6 starts a fresh Agent in the same container, preserving + # copied configuration and installed packages. + restart_command = ( + 'old_pid=$(pidof agent) || exit 1; ' + 'set -- $old_pid; old_pid=$1; ' + 'rm -f /var/run/s6/services/agent/finish; ' + 'kill "$old_pid" || exit 1; ' + 'elapsed=0; ' + 'while kill -0 "$old_pid" 2>/dev/null; do ' + f'[ "$elapsed" -ge {self._wait_timeout} ] && exit 1; ' + 'sleep 1; elapsed=$((elapsed + 1)); ' + 'done' + ) + self._exec(['sh', '-c', restart_command]) + self._wait_for_agent() + + def restart(self) -> None: + self._sync_local_packages() + self._require_prepared() + self._sync_config() + self._sync_auto_conf() + self._restart_agent_process() + self._require_prepared() + + def sync_config(self) -> None: + self._sync_config() + + def invoke(self, args: list[str], *, env_vars: dict[str, str] | None = None) -> None: + self._sync_local_packages() + self._require_prepared() + self._sync_config() + self._sync_auto_conf() + self._exec(['agent', *args], env_vars=env_vars) + self._require_prepared() + + def enter_shell(self) -> None: + self._kubectl( + [ + 'exec', + '-it', + '--namespace', + self._namespace, + f'pod/{POD_NAME}', + '--container', + CONTAINER_NAME, + '--', + 'bash', + ], + check=True, + ) + + def _show_logs(self, *, check: bool = False) -> None: + self._kubectl( + ['logs', '--namespace', self._namespace, f'pod/{POD_NAME}', '--container', CONTAINER_NAME], + check=check, + ) + + def show_logs(self) -> None: + self._show_logs(check=True) diff --git a/ddev/src/ddev/e2e/agent/vagrant.py b/ddev/src/ddev/e2e/agent/vagrant.py index f79e7e1a911dd..32b4218de1df5 100644 --- a/ddev/src/ddev/e2e/agent/vagrant.py +++ b/ddev/src/ddev/e2e/agent/vagrant.py @@ -57,6 +57,8 @@ def disable_integration_before_install(config_file: Path): class VagrantAgent(AgentInterface): + build_config_key = 'vagrant' + supports_ci = False VM_HOST_IP = "172.30.1.5" def __init__( diff --git a/ddev/tests/cli/env/test_agent.py b/ddev/tests/cli/env/test_agent.py index 73d911e4c994b..4c2ad7d06d5c5 100644 --- a/ddev/tests/cli/env/test_agent.py +++ b/ddev/tests/cli/env/test_agent.py @@ -1,8 +1,11 @@ # (C) Datadog, Inc. 2023-present # All rights reserved # Licensed under a 3-clause BSD style license (see LICENSE) +import subprocess + import pytest +from ddev.cli.env.agent import _sync_restored_config from ddev.e2e.config import EnvDataStorage @@ -171,3 +174,81 @@ def test_trigger_run_inject_integration(ddev, data_dir, mocker): assert not result.output invoke.assert_called_once_with(['check', integration, '-l', 'debug'], env_vars=None) + + +def test_temporary_config_is_restored_and_synchronized(ddev, data_dir, temp_dir, mocker): + invoke = mocker.patch('ddev.e2e.agent.docker.DockerAgent.invoke') + sync_config = mocker.patch('ddev.e2e.agent.docker.DockerAgent.sync_config') + + integration = 'postgres' + environment = 'py3.12' + env_data = EnvDataStorage(data_dir).get(integration, environment) + env_data.write_metadata({}) + original_config = {'instances': [{'host': 'original'}]} + env_data.write_config(original_config) + temporary_config = temp_dir / 'temporary.json' + temporary_config.write_text('{"instances": []}') + + result = ddev( + 'env', + 'agent', + integration, + environment, + '--config-file', + str(temporary_config), + 'check', + ) + + assert result.exit_code == 0, result.output + assert env_data.read_config() == original_config + invoke.assert_called_once_with(['check', integration], env_vars=None) + sync_config.assert_called_once_with() + + +def test_sync_failure_does_not_mask_check_failure(mocker): + check_error = subprocess.CalledProcessError(7, ['agent', 'check']) + sync_error = subprocess.CalledProcessError(8, ['kubectl', 'cp']) + app = mocker.Mock() + agent = mocker.Mock() + agent.sync_config.side_effect = sync_error + + with pytest.raises(subprocess.CalledProcessError) as exc_info: + try: + raise check_error + finally: + _sync_restored_config(app, agent) + + assert exc_info.value is check_error + app.display_warning.assert_called_once() + + with pytest.raises(subprocess.CalledProcessError) as exc_info: + _sync_restored_config(app, agent) + + assert exc_info.value is sync_error + + +def test_temporary_config_absence_is_restored_and_synchronized(ddev, data_dir, temp_dir, mocker): + invoke = mocker.patch('ddev.e2e.agent.docker.DockerAgent.invoke') + sync_config = mocker.patch('ddev.e2e.agent.docker.DockerAgent.sync_config') + + integration = 'postgres' + environment = 'py3.12' + env_data = EnvDataStorage(data_dir).get(integration, environment) + env_data.write_metadata({}) + temporary_config = temp_dir / 'temporary.json' + temporary_config.write_text('{"instances": []}') + + result = ddev( + 'env', + 'agent', + integration, + environment, + '--config-file', + str(temporary_config), + 'check', + ) + + assert result.exit_code == 0, result.output + assert not env_data.config_file.exists() + invoke.assert_called_once_with(['check', integration], env_vars=None) + sync_config.assert_called_once_with() diff --git a/ddev/tests/cli/env/test_logs.py b/ddev/tests/cli/env/test_logs.py new file mode 100644 index 0000000000000..6c2e207eab46a --- /dev/null +++ b/ddev/tests/cli/env/test_logs.py @@ -0,0 +1,52 @@ +# (C) Datadog, Inc. 2026-present +# All rights reserved +# Licensed under a 3-clause BSD style license (see LICENSE) +from ddev.e2e.config import EnvDataStorage + + +def test_nonexistent(ddev, helpers, mocker): + show_logs = mocker.patch('ddev.e2e.agent.docker.DockerAgent.show_logs') + + integration = 'postgres' + environment = 'py3.13' + + result = ddev('env', 'logs', integration, environment) + + assert result.exit_code == 1, result.output + assert result.output == helpers.dedent( + f""" + Environment `{environment}` for integration `{integration}` is not running + """ + ) + show_logs.assert_not_called() + + +def test_basic(ddev, data_dir, mocker): + show_logs = mocker.patch('ddev.e2e.agent.docker.DockerAgent.show_logs') + + integration = 'postgres' + environment = 'py3.13' + env_data = EnvDataStorage(data_dir).get(integration, environment) + env_data.write_metadata({}) + + result = ddev('env', 'logs', integration, environment) + + assert result.exit_code == 0, result.output + assert not result.output + show_logs.assert_called_once_with() + + +def test_log_failure_is_propagated(ddev, data_dir, mocker): + show_logs = mocker.patch( + 'ddev.e2e.agent.docker.DockerAgent.show_logs', side_effect=RuntimeError('Agent logs are unavailable') + ) + + integration = 'postgres' + environment = 'py3.13' + EnvDataStorage(data_dir).get(integration, environment).write_metadata({}) + + result = ddev('env', 'logs', integration, environment) + + assert result.exit_code == 1 + assert 'Agent logs are unavailable' in result.output + show_logs.assert_called_once_with() diff --git a/ddev/tests/cli/env/test_start.py b/ddev/tests/cli/env/test_start.py index 1ecab0a0256ad..7264886eff496 100644 --- a/ddev/tests/cli/env/test_start.py +++ b/ddev/tests/cli/env/test_start.py @@ -98,6 +98,29 @@ def test_stop_on_error(ddev, helpers, data_dir, write_result_file, mocker): stop.assert_called_once() +def test_unsupported_agent_on_ci_persists_environment_state(ddev, data_dir, write_result_file, mocker): + metadata = {'agent_type': 'vagrant'} + config = {} + run = mocker.patch('subprocess.run', side_effect=write_result_file({'metadata': metadata, 'config': config})) + mocker.patch('ddev.cli.env.start.running_on_ci', return_value=True) + start_agent = mocker.patch('ddev.e2e.agent.vagrant.VagrantAgent.start') + + integration = 'postgres' + environment = 'py3.12' + env_data = EnvDataStorage(data_dir).get(integration, environment) + + result = ddev('env', 'start', integration, environment) + + assert result.exit_code == 0, result.output + assert 'Starting: py3.12' in result.output + assert 'Stopping:' not in result.output + assert 'Vagrant is not supported on CI' in result.output + assert env_data.read_metadata() == metadata + assert env_data.read_config() == {'instances': [config]} + assert run.call_count == 1 + start_agent.assert_not_called() + + def test_basic(ddev, helpers, data_dir, write_result_file, mocker): metadata = {} config = {} @@ -175,6 +198,30 @@ def test_agent_build_config(ddev, config_file, helpers, data_dir, write_result_f ) +def test_kubernetes_agent_build_uses_configured_docker_image(ddev, config_file, data_dir, write_result_file, mocker): + config_file.model.agent = '7' + config_file.save() + + metadata = {'agent_type': 'kubernetes', 'kubernetes': {'kubeconfig': '/tmp/kubeconfig'}} + config = {} + mocker.patch('subprocess.run', side_effect=write_result_file({'metadata': metadata, 'config': config})) + start = mocker.patch('ddev.e2e.agent.kubernetes.KubernetesAgent.start') + + integration = 'postgres' + environment = 'py3.12' + env_data = EnvDataStorage(data_dir).get(integration, environment) + + result = ddev('env', 'start', integration, environment) + + assert result.exit_code == 0, result.output + assert env_data.read_metadata() == metadata + start.assert_called_once_with( + agent_build='registry.datadoghq.com/agent:7', + local_packages={}, + env_vars={'DD_DD_URL': 'https://app.datadoghq.com', 'DD_SITE': 'datadoghq.com'}, + ) + + def test_agent_build_env_var(ddev, config_file, helpers, data_dir, write_result_file, mocker): config_file.model.agent = '7' config_file.save() diff --git a/ddev/tests/cli/env/test_stop.py b/ddev/tests/cli/env/test_stop.py index 3e132c822a9e4..aa253df83776d 100644 --- a/ddev/tests/cli/env/test_stop.py +++ b/ddev/tests/cli/env/test_stop.py @@ -1,6 +1,8 @@ # (C) Datadog, Inc. 2023-present # All rights reserved # Licensed under a 3-clause BSD style license (see LICENSE) +import pytest + from ddev.e2e.config import EnvDataStorage @@ -46,6 +48,42 @@ def test_basic(ddev, helpers, data_dir, mocker): stop.assert_called_once() +def test_failed_agent_cleanup_removes_environment_state(ddev, data_dir, mocker): + teardown = mocker.patch('subprocess.run', return_value=mocker.MagicMock(returncode=0)) + stop = mocker.patch('ddev.e2e.agent.docker.DockerAgent.stop', side_effect=RuntimeError('cleanup failed')) + + integration = 'postgres' + environment = 'py3.12' + env_data = EnvDataStorage(data_dir).get(integration, environment) + metadata = {'owner': 'retry-me'} + env_data.write_metadata(metadata) + + with pytest.raises(RuntimeError, match='cleanup failed'): + ddev('env', 'stop', integration, environment) + + assert not env_data.exists() + stop.assert_called_once_with() + teardown.assert_called_once() + + +def test_failed_kubernetes_fixture_teardown_removes_environment_state(ddev, data_dir, mocker): + teardown = mocker.patch('subprocess.run', return_value=mocker.MagicMock(returncode=1)) + stop = mocker.patch('ddev.e2e.agent.kubernetes.KubernetesAgent.stop') + + integration = 'postgres' + environment = 'py3.12' + env_data = EnvDataStorage(data_dir).get(integration, environment) + metadata = {'agent_type': 'kubernetes', 'kubernetes': {'kubeconfig': '/tmp/kubeconfig'}} + env_data.write_metadata(metadata) + + result = ddev('env', 'stop', integration, environment) + + assert result.exit_code == 1, result.output + assert not env_data.exists() + stop.assert_called_once_with() + teardown.assert_called_once() + + def test_stop_all(ddev, helpers, data_dir, mocker): mocker.patch('subprocess.run', return_value=mocker.MagicMock(returncode=0)) stop = mocker.patch('ddev.e2e.agent.docker.DockerAgent.stop') diff --git a/ddev/tests/e2e/agent/test_kubernetes.py b/ddev/tests/e2e/agent/test_kubernetes.py new file mode 100644 index 0000000000000..d8fdffcb05354 --- /dev/null +++ b/ddev/tests/e2e/agent/test_kubernetes.py @@ -0,0 +1,520 @@ +# (C) Datadog, Inc. 2026-present +# All rights reserved +# Licensed under a 3-clause BSD style license (see LICENSE) +import json +import subprocess + +import pytest + +from ddev.e2e.agent.constants import AgentEnvVars +from ddev.e2e.agent.kubernetes import KubernetesAgent +from ddev.integration.core import Integration +from ddev.repo.config import RepositoryConfig + +TEST_KUBECONFIG = '/tmp/kubeconfig' +TEST_NAMESPACE = 'ddev-agent' +TEST_POD = 'ddev-agent' +TEST_CONTAINER = 'agent' +TEST_PREPARED_MARKER = '/home/.ddev-agent-prepared' + + +@pytest.fixture(scope='module') +def get_integration(local_repo): + def _get_integration(name): + return Integration(local_repo / name, local_repo, RepositoryConfig(local_repo / '.ddev' / 'config.toml')) + + return _get_integration + + +@pytest.fixture +def config_file(temp_dir): + path = temp_dir / 'config' / 'velero.yaml' + path.parent.ensure_dir_exists() + path.write_text('instances: []\n') + return path + + +@pytest.fixture +def auto_conf(temp_dir): + path = temp_dir / 'auto_conf.yaml' + path.write_text('ad_identifiers:\n - velero\n') + return path + + +@pytest.fixture +def metadata(auto_conf): + return { + 'kubernetes': { + 'kubeconfig': TEST_KUBECONFIG, + 'auto_conf': str(auto_conf), + }, + } + + +@pytest.fixture +def agent(app, get_integration, metadata, config_file): + return KubernetesAgent(app, get_integration('velero'), 'py3.12', metadata, config_file) + + +def successful_process(command, *, stdout=b'', stderr=None): + return subprocess.CompletedProcess(command, 0, stdout=stdout, stderr=stderr) + + +@pytest.fixture +def run_command(app, mocker): + def run(command, **kwargs): + if command[-2:] == ['config', 'current-context']: + return successful_process(command, stdout=b'kind-test\n') + if command[-4:] == ['get', 'nodes', '-o', 'json']: + nodes = {'items': [{'metadata': {'name': 'kind-control-plane'}, 'spec': {}}]} + return successful_process(command, stdout=json.dumps(nodes).encode()) + return successful_process(command) + + return mocker.patch.object(app.platform, 'run_command', side_effect=run) + + +def command_calls(run_command): + return [call.args[0] for call in run_command.call_args_list] + + +def option_value(command: list[str], option: str) -> str: + return command[command.index(option) + 1] + + +def find_kubectl_command(calls: list[list[str]], expected: list[str]) -> int: + matches = [] + for index, command in enumerate(calls): + kubectl_args = command[: command.index('--')] if '--' in command else command + contains_expected = any( + kubectl_args[start : start + len(expected)] == expected + for start in range(len(kubectl_args) - len(expected) + 1) + ) + if not contains_expected: + continue + + assert command[0] == 'kubectl' + assert option_value(kubectl_args, '--kubeconfig') == TEST_KUBECONFIG + matches.append(index) + + assert len(matches) == 1, f'Expected one kubectl command containing {expected!r}, found {len(matches)}' + return matches[0] + + +def exec_command_indices(calls: list[list[str]], expected: list[str]) -> list[int]: + return [ + index + for index, command in enumerate(calls) + if '--' in command and command[command.index('--') + 1 :] == expected + ] + + +def find_exec_command(calls: list[list[str]], expected: list[str], *, prefix: bool = False) -> int: + matches = [] + for index, command in enumerate(calls): + if 'exec' not in command or '--' not in command: + continue + + separator = command.index('--') + payload = command[separator + 1 :] + matches_expected = payload[: len(expected)] == expected if prefix else payload == expected + if not matches_expected: + continue + + kubectl_args = command[:separator] + assert command[0] == 'kubectl' + assert option_value(kubectl_args, '--kubeconfig') == TEST_KUBECONFIG + assert option_value(kubectl_args, '--namespace') == TEST_NAMESPACE + assert f'pod/{TEST_POD}' in kubectl_args + assert option_value(kubectl_args, '--container') == TEST_CONTAINER + matches.append(index) + + assert len(matches) == 1, f'Expected one kubectl exec payload matching {expected!r}, found {len(matches)}' + return matches[0] + + +def find_copy_command(calls: list[list[str]], source: str, destination: str) -> int: + matches = [] + for index, command in enumerate(calls): + if 'cp' not in command or source not in command or destination not in command: + continue + + assert command[0] == 'kubectl' + assert option_value(command, '--kubeconfig') == TEST_KUBECONFIG + assert option_value(command, '--container') == TEST_CONTAINER + assert command.index(source) < command.index(destination) + matches.append(index) + + assert len(matches) == 1, f'Expected one kubectl cp from {source!r} to {destination!r}, found {len(matches)}' + return matches[0] + + +def test_restart_waits_for_agent_before_and_after_restart(agent, mocker): + operations = mocker.Mock() + execute = mocker.patch.object(agent, '_exec') + wait_for_agent = mocker.patch.object(agent, '_wait_for_agent') + operations.attach_mock(execute, 'execute') + operations.attach_mock(wait_for_agent, 'wait_for_agent') + + agent._restart_agent_process() + + assert operations.mock_calls == [ + mocker.call.wait_for_agent(), + mocker.call.execute(mocker.ANY), + mocker.call.wait_for_agent(), + ] + + +def test_rejects_container_that_lost_prepared_state(agent, mocker): + execute = mocker.patch.object( + agent, + '_exec', + return_value=subprocess.CompletedProcess([], 1), + ) + + with pytest.raises(RuntimeError, match='may have restarted'): + agent._require_prepared() + + execute.assert_called_once_with(['test', '-f', TEST_PREPARED_MARKER], check=False) + + +def test_start_uses_selected_image_rbac_config_and_local_packages( + agent, metadata, config_file, auto_conf, temp_dir, run_command +): + local_base = temp_dir / 'datadog_checks_base' + local_base.ensure_dir_exists() + integration = temp_dir / 'velero' + integration.ensure_dir_exists() + metadata['start_commands'] = ['echo start'] + metadata['post_install_commands'] = ['echo post-install'] + + agent.start( + agent_build='registry.example.com/datadog-agent:test', + local_packages={local_base: '[kube]', integration: '[deps]'}, + env_vars={'DD_SITE': 'datadoghq.com'}, + ) + + calls = command_calls(run_command) + prefix = ['kubectl', '--kubeconfig', TEST_KUBECONFIG] + context_index = find_kubectl_command(calls, ['config', 'current-context']) + topology_index = find_kubectl_command(calls, ['get', 'nodes', '-o', 'json']) + create_calls = [call for call in run_command.call_args_list if call.args[0] == [*prefix, 'create', '-f', '-']] + assert len(create_calls) == 1 + create_index = run_command.call_args_list.index(create_calls[0]) + assert context_index < topology_index < create_index + manifest = json.loads(create_calls[0].kwargs['input']) + resources = {item['kind']: item for item in manifest['items']} + assert resources['Namespace']['metadata']['name'] == 'ddev-agent' + assert resources['ServiceAccount']['metadata']['namespace'] == 'ddev-agent' + assert resources['ClusterRole']['metadata']['name'] == 'ddev-agent' + assert resources['ClusterRoleBinding']['metadata']['name'] == 'ddev-agent' + assert resources['ClusterRoleBinding']['roleRef']['name'] == 'ddev-agent' + assert resources['ClusterRoleBinding']['subjects'][0]['namespace'] == 'ddev-agent' + assert resources['Pod']['metadata']['namespace'] == 'ddev-agent' + assert resources['Pod']['spec']['containers'][0]['name'] == 'agent' + assert resources['Pod']['spec']['containers'][0]['image'] == 'registry.example.com/datadog-agent:test' + assert resources['Pod']['spec']['containers'][0]['imagePullPolicy'] == 'Always' + assert resources['Pod']['spec']['serviceAccountName'] == 'ddev-agent' + resource_labels = {'app.kubernetes.io/managed-by': 'ddev'} + for kind in ('Namespace', 'ServiceAccount', 'ClusterRole', 'ClusterRoleBinding'): + assert resources[kind]['metadata']['labels'] == resource_labels + assert resources['Pod']['metadata']['labels'] == { + **resource_labels, + 'app.kubernetes.io/name': 'ddev-agent', + } + assert resources['ClusterRole']['rules'] == [ + {'apiGroups': [''], 'resources': ['nodes'], 'verbs': ['get', 'list', 'watch']}, + { + 'apiGroups': [''], + 'resources': ['nodes/metrics', 'nodes/spec', 'nodes/stats', 'nodes/proxy'], + 'verbs': ['get'], + }, + { + 'apiGroups': [''], + 'resources': ['pods', 'endpoints', 'services'], + 'verbs': ['get', 'list', 'watch'], + }, + ] + env = {item['name']: item.get('value') for item in resources['Pod']['spec']['containers'][0]['env']} + assert env['DD_API_KEY'] == 'a' * 32 + assert env['DD_SITE'] == 'datadoghq.com' + assert env['DD_AUTOCONFIG_FROM_ENVIRONMENT'] == 'true' + assert 'DD_KUBERNETES_KUBELET_HOST' in env + assert 'DD_KUBERNETES_KUBELET_NODENAME' in env + + start_index = find_exec_command(calls, ['echo', 'start']) + post_install_index = find_exec_command(calls, ['echo', 'post-install']) + local_base_copy_index = find_copy_command( + calls, + local_base.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/home/datadog_checks_base', + ) + find_exec_command( + calls, + [ + '/opt/datadog-agent/embedded/bin/python3', + '-m', + 'pip', + 'install', + '--disable-pip-version-check', + '-e', + '/home/datadog_checks_base[kube]', + ], + ) + find_copy_command( + calls, + config_file.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/etc/datadog-agent/conf.d/velero.d/conf.yaml', + ) + find_copy_command( + calls, + auto_conf.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/etc/datadog-agent/conf.d/velero.d/auto_conf.yaml', + ) + restart_index = find_exec_command(calls, ['sh', '-c'], prefix=True) + marker_index = find_exec_command(calls, ['touch', TEST_PREPARED_MARKER]) + prepared_indices = exec_command_indices(calls, ['test', '-f', TEST_PREPARED_MARKER]) + assert len(prepared_indices) == 1 + assert start_index < post_install_index < marker_index < restart_index < prepared_indices[0] + assert metadata['kubernetes']['local_packages'] == { + str(local_base): '[kube]', + str(integration): '[deps]', + } + local_base_copy = run_command.call_args_list[local_base_copy_index] + assert local_base_copy.kwargs['check'] is True + assert local_base_copy.kwargs['cwd'] == local_base.resolve().parent + + +@pytest.mark.parametrize('wait_timeout', [0, True]) +def test_rejects_invalid_wait_timeout_before_creating_resources(agent, metadata, run_command, wait_timeout): + metadata['kubernetes']['wait_timeout'] = wait_timeout + + with pytest.raises(ValueError, match='wait_timeout'): + agent.start(agent_build='', local_packages={}, env_vars={}) + + run_command.assert_not_called() + + +@pytest.mark.parametrize( + ('env_vars', 'expected'), + [ + ({AgentEnvVars.DOGSTATSD_PORT: '8125'}, 'DogStatsD exposure'), + ({AgentEnvVars.LOGS_ENABLED: 'true'}, 'shared-file log collection'), + ], +) +def test_rejects_unsupported_options_before_creating_resources(agent, run_command, env_vars, expected): + with pytest.raises(NotImplementedError, match=expected): + agent.start(agent_build='', local_packages={}, env_vars=env_vars) + + run_command.assert_not_called() + + +def test_start_detects_container_replacement_during_restart(agent, app, mocker): + def run(command, **kwargs): + if command[-2:] == ['config', 'current-context']: + return successful_process(command, stdout=b'kind-test\n') + if command[-4:] == ['get', 'nodes', '-o', 'json']: + nodes = {'items': [{'metadata': {'name': 'kind-control-plane'}, 'spec': {}}]} + return successful_process(command, stdout=json.dumps(nodes).encode()) + # A replaced container no longer carries the marker stamped before the restart. + if command[-3:] == ['test', '-f', TEST_PREPARED_MARKER]: + return subprocess.CompletedProcess(command, 1) + return successful_process(command) + + mocker.patch.object(app.platform, 'run_command', side_effect=run) + + with pytest.raises(RuntimeError, match='may have restarted'): + agent.start(agent_build='', local_packages={}, env_vars={}) + + +def test_rejects_non_kind_context_before_inspecting_cluster(agent, app, mocker): + run_command = mocker.patch.object( + app.platform, + 'run_command', + return_value=successful_process([], stdout=b'prod\n'), + ) + + with pytest.raises(RuntimeError, match='non-Kind Kubernetes context `prod`'): + agent.start(agent_build='', local_packages={}, env_vars={}) + + run_command.assert_called_once_with( + ['kubectl', '--kubeconfig', TEST_KUBECONFIG, 'config', 'current-context'], + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + ) + + +def test_rejects_multi_node_clusters(agent, app, mocker): + nodes = {'items': [{'metadata': {'name': 'one'}, 'spec': {}}, {'metadata': {'name': 'two'}, 'spec': {}}]} + + def run(command, **kwargs): + if command[-2:] == ['config', 'current-context']: + return successful_process(command, stdout=b'kind-test\n') + return successful_process(command, stdout=json.dumps(nodes).encode()) + + run_command = mocker.patch.object(app.platform, 'run_command', side_effect=run) + + with pytest.raises(NotImplementedError, match='exactly one schedulable node'): + agent.start(agent_build='', local_packages={}, env_vars={}) + + calls = command_calls(run_command) + context_index = find_kubectl_command(calls, ['config', 'current-context']) + topology_index = find_kubectl_command(calls, ['get', 'nodes', '-o', 'json']) + assert context_index < topology_index + assert not any('create' in command for command in calls) + for index in (context_index, topology_index): + assert run_command.call_args_list[index].kwargs['stdout'] is subprocess.PIPE + assert run_command.call_args_list[index].kwargs['stderr'] is subprocess.PIPE + + +def test_validation_ignores_kubectl_warnings(agent, app, mocker): + warning = b'Warning: deprecated API\n' + nodes = {'items': [{'metadata': {'name': 'kind-control-plane'}, 'spec': {}}]} + run_command = mocker.patch.object( + app.platform, + 'run_command', + side_effect=[ + successful_process([], stdout=b'kind-test\n', stderr=warning), + successful_process([], stdout=json.dumps(nodes).encode(), stderr=warning), + ], + ) + + agent._validate_context() + agent._validate_topology() + + assert all(call.kwargs['stderr'] is subprocess.PIPE for call in run_command.call_args_list) + + +def test_partial_manifest_creation_failure_is_propagated(agent, app, mocker): + nodes = {'items': [{'metadata': {'name': 'kind-control-plane'}, 'spec': {}}]} + + def run(command, **kwargs): + if command[-2:] == ['config', 'current-context']: + return successful_process(command, stdout=b'kind-test\n') + if command[-4:] == ['get', 'nodes', '-o', 'json']: + return successful_process(command, stdout=json.dumps(nodes).encode()) + if command[-3:] == ['create', '-f', '-']: + return subprocess.CompletedProcess(command, 1, stdout=b'', stderr=b'namespace already exists') + return successful_process(command) + + run_command = mocker.patch.object(app.platform, 'run_command', side_effect=run) + + with pytest.raises(RuntimeError, match='Unable to create Kubernetes Agent resources'): + agent.start(agent_build='', local_packages={}, env_vars={}) + + calls = command_calls(run_command) + assert sum(command[-3:] == ['create', '-f', '-'] for command in calls) == 1 + assert not any('apply' in command or 'wait' in command for command in calls) + + +def test_invoke_synchronizes_config_and_environment(agent, config_file, auto_conf, run_command): + config_file.write_text('instances:\n - openmetrics_endpoint: http://velero:8085/metrics\n') + + agent.invoke(['check', 'velero', '--json'], env_vars={'ZED': 'last', 'ALPHA': 'first'}) + + calls = command_calls(run_command) + config_copy_index = find_copy_command( + calls, + config_file.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/etc/datadog-agent/conf.d/velero.d/conf.yaml', + ) + auto_conf_copy_index = find_copy_command( + calls, + auto_conf.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/etc/datadog-agent/conf.d/velero.d/auto_conf.yaml', + ) + invoke_index = find_exec_command( + calls, + ['env', 'ALPHA=first', 'ZED=last', 'agent', 'check', 'velero', '--json'], + ) + prepared_indices = exec_command_indices(calls, ['test', '-f', TEST_PREPARED_MARKER]) + assert run_command.call_args_list[invoke_index].kwargs['check'] is True + assert len(prepared_indices) == 2 + assert prepared_indices[0] < config_copy_index < auto_conf_copy_index < invoke_index < prepared_indices[1] + + +def test_invoke_removes_pod_config_when_host_config_is_absent(agent, config_file, run_command): + config_file.remove() + + agent.invoke(['status']) + + calls = command_calls(run_command) + find_exec_command(calls, ['rm', '-f', '/etc/datadog-agent/conf.d/velero.d/conf.yaml']) + + +def test_stop_is_no_op(agent, run_command): + agent.stop() + + run_command.assert_not_called() + + +def test_restart_resynchronizes_config_auto_conf_and_editable_sources( + agent, metadata, config_file, temp_dir, run_command +): + local_package = temp_dir / 'velero-source' + local_package.ensure_dir_exists() + metadata['kubernetes']['local_packages'] = {str(local_package): '[deps]'} + + agent.restart() + + calls = command_calls(run_command) + package_copy_index = find_copy_command( + calls, + local_package.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/home/velero-source', + ) + config_copy_index = find_copy_command( + calls, + config_file.name, + f'{TEST_NAMESPACE}/{TEST_POD}:/etc/datadog-agent/conf.d/velero.d/conf.yaml', + ) + restart_index = find_exec_command(calls, ['sh', '-c'], prefix=True) + prepared_indices = exec_command_indices(calls, ['test', '-f', TEST_PREPARED_MARKER]) + assert len(prepared_indices) == 2 + assert package_copy_index < prepared_indices[0] < config_copy_index < restart_index < prepared_indices[1] + assert not any('pip' in command for command in calls) + + +def test_restart_rejects_unsafe_local_package_destination(agent, metadata, run_command): + metadata['kubernetes']['local_packages'] = {'..': '[deps]'} + + with pytest.raises(ValueError, match='local package path'): + agent.restart() + + run_command.assert_not_called() + + +def test_enter_shell_starts_interactive_bash(agent, run_command): + agent.enter_shell() + + calls = command_calls(run_command) + shell_index = find_exec_command(calls, ['bash']) + kubectl_args = calls[shell_index][: calls[shell_index].index('--')] + assert '-it' in kubectl_args + assert run_command.call_args_list[shell_index].kwargs['check'] is True + + +def test_show_logs_targets_agent_container(agent, run_command): + agent.show_logs() + + calls = command_calls(run_command) + logs_index = find_kubectl_command(calls, ['logs']) + command = calls[logs_index] + assert option_value(command, '--namespace') == TEST_NAMESPACE + assert f'pod/{TEST_POD}' in command + assert option_value(command, '--container') == TEST_CONTAINER + assert run_command.call_args_list[logs_index].kwargs['check'] is True + + +def test_kubeconfig_validation(app, get_integration, config_file): + agent = KubernetesAgent(app, get_integration('velero'), 'py3.12', {'kubernetes': {}}, config_file) + with pytest.raises(ValueError, match='non-empty `kubeconfig`'): + _ = agent._kubeconfig + + +@pytest.mark.parametrize('kubernetes_metadata', [None, 'kubeconfig'], ids=['missing', 'not_a_mapping']) +def test_rejects_metadata_without_kubernetes_mapping(app, get_integration, config_file, kubernetes_metadata): + metadata = {} if kubernetes_metadata is None else {'kubernetes': kubernetes_metadata} + agent = KubernetesAgent(app, get_integration('velero'), 'py3.12', metadata, config_file) + + with pytest.raises(ValueError, match='requires a `kubernetes` metadata mapping'): + _ = agent._kubeconfig diff --git a/docs/developer/ddev/plugins.md b/docs/developer/ddev/plugins.md index 470dc3ad8ec51..3673721508999 100644 --- a/docs/developer/ddev/plugins.md +++ b/docs/developer/ddev/plugins.md @@ -137,8 +137,9 @@ is desired. This fixture is responsible for starting and stopping environments a #### Metadata -- `env_type` - This is the type of interface that will be used to interact with the Agent. Currently, we support `docker` (default), `vagrant` and `local`. +- `agent_type` - The interface used to interact with the Agent. Supported values are `docker` (default), `vagrant`, and `kubernetes`. - `env_vars` - A `dict` of environment variables and their values that will be present when starting the Agent. -- `docker_volumes` - A `list` of `str` representing [Docker volume mounts][docker-volume-docs] if `env_type` is `docker` e.g. `/local/path:/agent/container/path:ro`. -- `docker_platform` - The container architecture to use if `env_type` is `docker`. Currently, we support `linux` (default) and `windows`. -- `logs_config` - A `list` of configs that will be used by the Logs Agent. You will never need to use this directly, but rather via [higher level abstractions](test.md#logs). +- `docker_volumes` - A `list` of `str` representing [Docker volume mounts][docker-volume-docs] if `agent_type` is `docker` e.g. `/local/path:/agent/container/path:ro`. +- `docker_platform` - The container architecture to use if `agent_type` is `docker`. Currently, we support `linux` (default) and `windows`. +- `kubernetes` - Configuration used when `agent_type` is `kubernetes`. It requires `kubeconfig` and optionally accepts `auto_conf`, `image_pull_policy`, and `wait_timeout`. The backend is supported only with a disposable, fixture-owned cluster provisioned by `kind_run`; it uses the fixed `ddev-agent` namespace, and fixture teardown deletes the whole cluster. +- `logs_config` - A `list` of configs used by the Logs Agent. The [higher-level shared-file log helper](test.md#logs) currently supports only `agent_type: docker`. diff --git a/docs/developer/ddev/test.md b/docs/developer/ddev/test.md index 34895ef89ad78..d5119efa949bd 100644 --- a/docs/developer/ddev/test.md +++ b/docs/developer/ddev/test.md @@ -91,6 +91,42 @@ You can use the `%HOST%` template variable in your configuration instead of hard Note: Vagrant environments are not supported in CI environments due to virtualization constraints. +### Kubernetes Agent + +The Kubernetes Agent backend runs the Datadog Agent only in a disposable Kubernetes cluster that is owned by the E2E +fixture. The currently supported provisioner is `kind_run`, which creates the cluster before the fixture yields and deletes +the whole cluster during fixture teardown. The backend is not designed for shared or pre-existing clusters. + +```python +@pytest.fixture(scope='session') +def dd_environment(): + with kind_run(conditions=[setup_workload]) as kubeconfig: + yield { + 'instances': [ + # Use endpoints reachable from inside the cluster. + {'openmetrics_endpoint': 'http://my-service.default.svc.cluster.local:8080/metrics'}, + ], + }, { + 'agent_type': 'kubernetes', + 'kubernetes': { + 'kubeconfig': kubeconfig, + # Optional: install an Autodiscovery template for this integration. + 'auto_conf': os.path.join(CHECK_ROOT, 'datadog_checks', CHECK_NAME, 'data', 'auto_conf.yaml'), + }, + } +``` + +The backend uses the Agent image selected by `ddev env start --agent` or `DDEV_E2E_AGENT`, installs and synchronizes +local packages requested by `--dev` or `--base`, and implements Agent commands through `kubectl exec`. Static and +discovery E2E tests therefore continue to use `dd_agent_check` and `dd_agent_check_discovery`. + +Agent images default to the `Always` pull policy so mutable release and development tags are refreshed. Environments +that import a local image into the cluster can set `image_pull_policy` to `IfNotPresent` or `Never`. The backend creates +one Agent in the fixed `ddev-agent` namespace with fixed cluster-scoped RBAC names. Resource cleanup happens when +`kind_run` deletes the whole disposable cluster. + +The current implementation supports exactly one schedulable Kubernetes node. + ### Terraform The `terraform_run` utility makes it easy to create services from a directory of [Terraform][terraform-home] files. diff --git a/keda/tests/conftest.py b/keda/tests/conftest.py index 44f8cb07eb0d4..e2251222309c8 100644 --- a/keda/tests/conftest.py +++ b/keda/tests/conftest.py @@ -3,12 +3,10 @@ # Licensed under a 3-clause BSD style license (see LICENSE) import copy import os -from contextlib import ExitStack import pytest from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward from datadog_checks.dev.subprocess import run_command from . import common @@ -27,13 +25,18 @@ def setup_ked(): @pytest.fixture(scope='session') def dd_environment(): - with kind_run(conditions=[setup_ked], sleep=30) as kubeconfig, ExitStack() as stack: - keda_host, keda_port = stack.enter_context( - port_forward(kubeconfig, 'keda', 8080, 'deployment', 'keda-operator-metrics-apiserver') + with kind_run(conditions=[setup_ked], sleep=30) as kubeconfig: + instances = [ + {'openmetrics_endpoint': ('http://keda-operator-metrics-apiserver.keda.svc.cluster.local:8080/metrics')} + ] + + yield ( + {'instances': instances}, + { + 'agent_type': 'kubernetes', + 'kubernetes': {'kubeconfig': kubeconfig}, + }, ) - instances = [{'openmetrics_endpoint': f'http://{keda_host}:{keda_port}/metrics'}] - - yield {'instances': instances} @pytest.fixture diff --git a/kuma/tests/conftest.py b/kuma/tests/conftest.py index e7f54ba7284fc..2311f425a0eed 100644 --- a/kuma/tests/conftest.py +++ b/kuma/tests/conftest.py @@ -2,82 +2,75 @@ # All rights reserved # Licensed under a 3-clause BSD style license (see LICENSE) import os -import time -from contextlib import ExitStack +import shlex import pytest -import requests from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward from datadog_checks.dev.subprocess import run_command +KUMA_NAMESPACE = 'kuma-system' +KUMA_SERVICE = 'kuma-control-plane' +KUMA_STARTUP_TIMEOUT = 600 +KUMA_METRICS_ENDPOINT = f'http://{KUMA_SERVICE}.{KUMA_NAMESPACE}.svc.cluster.local:5680/metrics' +KUMA_SOURCE_COMMAND = shlex.join( + [ + '/opt/datadog-agent/embedded/bin/python3', + '-c', + "import datadog_checks.kuma; print('Resolved Kuma module source:', datadog_checks.kuma.__file__)", + ] +) + def setup_kuma(): - kuma_version = os.environ.get("KUMA_VERSION", "2.10.6") - run_command(["kubectl", "create", "namespace", "kuma-system"]) - run_command(["helm", "repo", "add", "kuma", "https://kumahq.github.io/charts"]) - run_command(["helm", "repo", "update"]) + kuma_version = os.environ.get('KUMA_VERSION', '2.10.6') + run_command(['kubectl', 'create', 'namespace', KUMA_NAMESPACE], check=True) + run_command(['helm', 'repo', 'add', 'kuma', 'https://kumahq.github.io/charts'], check=True) + run_command(['helm', 'repo', 'update'], check=True) run_command( [ - "helm", - "upgrade", - "--install", - "kuma", - "kuma/kuma", - "--version", + 'helm', + 'upgrade', + '--install', + 'kuma', + 'kuma/kuma', + '--version', kuma_version, - "--create-namespace", - "-n", - "kuma-system", - ] + '--create-namespace', + '-n', + KUMA_NAMESPACE, + ], + check=True, ) run_command( - ["kubectl", "rollout", "status", "deployment/kuma-control-plane", "-n", "kuma-system", "--timeout=180s"] + [ + 'kubectl', + 'rollout', + 'status', + f'deployment/{KUMA_SERVICE}', + '-n', + KUMA_NAMESPACE, + # The Kuma deployment's readiness probe checks its in-cluster /ready endpoint. + f'--timeout={KUMA_STARTUP_TIMEOUT}s', + ], + check=True, ) -def wait_for_kuma_readiness(api_url, api_port, max_wait=600): - """Wait for Kuma control plane to be ready by querying the /config endpoint.""" - config_url = f'http://{api_url}:{api_port}/config' - start_time = time.monotonic() - - while time.monotonic() - start_time < max_wait: - try: - response = requests.get(config_url, timeout=5) - if response.ok: - print(f"Kuma control plane is ready at {config_url} (took {time.monotonic() - start_time} seconds)") - return - except (requests.exceptions.RequestException, requests.exceptions.Timeout): - pass - - print(f"Waiting for Kuma control plane to be ready at {config_url}...") - time.sleep(0.5) - - raise TimeoutError(f"Kuma control plane did not become ready within {max_wait} seconds") - - @pytest.fixture(scope='session') def dd_environment(dd_save_state): with kind_run(conditions=[setup_kuma]) as kubeconfig: - with ExitStack() as stack: - kuma_metrics_url, kuma_metrics_port = stack.enter_context( - port_forward(kubeconfig, 'kuma-system', 5680, 'service', 'kuma-control-plane') - ) - kuma_api_url, kuma_api_port = stack.enter_context( - port_forward(kubeconfig, 'kuma-system', 5681, 'service', 'kuma-control-plane') - ) - - # Wait for Kuma control plane to be ready - wait_for_kuma_readiness(kuma_api_url, kuma_api_port) - - metrics_endpoint = f'http://{kuma_metrics_url}:{kuma_metrics_port}/metrics' - - env_instance = {'openmetrics_endpoint': metrics_endpoint} - - dd_save_state("kuma_instance", env_instance) - - yield env_instance + instance = {'openmetrics_endpoint': KUMA_METRICS_ENDPOINT} + metadata = { + 'agent_type': 'kubernetes', + 'kubernetes': { + 'kubeconfig': kubeconfig, + }, + 'post_install_commands': [KUMA_SOURCE_COMMAND], + } + + dd_save_state('kuma_instance', instance) + yield instance, metadata @pytest.fixture(scope='session') diff --git a/kyverno/tests/conftest.py b/kyverno/tests/conftest.py index d75c270f964d7..66f1120e53e36 100644 --- a/kyverno/tests/conftest.py +++ b/kyverno/tests/conftest.py @@ -2,13 +2,11 @@ # All rights reserved # Licensed under a 3-clause BSD style license (see LICENSE) import os -from contextlib import ExitStack import pytest from datadog_checks.dev import get_here from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward from datadog_checks.dev.subprocess import run_command HERE = get_here() @@ -34,24 +32,30 @@ def setup_kyverno(): @pytest.fixture(scope='session') def dd_environment(): - with kind_run(conditions=[setup_kyverno], sleep=30) as kubeconfig, ExitStack() as stack: - kyverno_host1, kyverno_port1 = stack.enter_context( - port_forward(kubeconfig, 'kyverno', 8000, 'deployment', 'kyverno-admission-controller'), - ) - kyverno_host2, kyverno_port2 = stack.enter_context( - port_forward(kubeconfig, 'kyverno', 8000, 'deployment', 'kyverno-background-controller'), - ) - kyverno_host3, kyverno_port3 = stack.enter_context( - port_forward(kubeconfig, 'kyverno', 8000, 'deployment', 'kyverno-cleanup-controller'), - ) - kyverno_host4, kyverno_port4 = stack.enter_context( - port_forward(kubeconfig, 'kyverno', 8000, 'deployment', 'kyverno-reports-controller'), - ) + with kind_run(conditions=[setup_kyverno], sleep=30) as kubeconfig: instances = [ - {'openmetrics_endpoint': f'http://{kyverno_host1}:{kyverno_port1}/metrics'}, - {'openmetrics_endpoint': f'http://{kyverno_host2}:{kyverno_port2}/metrics'}, - {'openmetrics_endpoint': f'http://{kyverno_host3}:{kyverno_port3}/metrics'}, - {'openmetrics_endpoint': f'http://{kyverno_host4}:{kyverno_port4}/metrics'}, + {'openmetrics_endpoint': 'http://kyverno-svc-metrics.kyverno.svc.cluster.local:8000/metrics'}, + { + 'openmetrics_endpoint': ( + 'http://kyverno-background-controller-metrics.kyverno.svc.cluster.local:8000/metrics' + ) + }, + { + 'openmetrics_endpoint': ( + 'http://kyverno-cleanup-controller-metrics.kyverno.svc.cluster.local:8000/metrics' + ) + }, + { + 'openmetrics_endpoint': ( + 'http://kyverno-reports-controller-metrics.kyverno.svc.cluster.local:8000/metrics' + ) + }, ] - - yield {'instances': instances} + metadata = { + 'agent_type': 'kubernetes', + 'kubernetes': { + 'kubeconfig': kubeconfig, + }, + } + + yield {'instances': instances}, metadata diff --git a/tekton/tests/conftest.py b/tekton/tests/conftest.py index c84cffdf590cd..46723b60bc7a8 100644 --- a/tekton/tests/conftest.py +++ b/tekton/tests/conftest.py @@ -2,15 +2,12 @@ # All rights reserved # Licensed under a 3-clause BSD style license (see LICENSE) import os -from contextlib import ExitStack -from urllib.request import urlopen import pytest from datadog_checks.dev import get_here, run_command from datadog_checks.dev.conditions import WaitFor from datadog_checks.dev.kind import kind_run -from datadog_checks.dev.kube_port_forward import port_forward HERE = get_here() @@ -23,11 +20,10 @@ def _wait_for_resource(*, kind, name, condition, namespace=None, timeout="300s") run_command(command) -def wait_for_metric_families(endpoint: str, families: list[str]) -> None: - with urlopen(endpoint, timeout=5) as response: - metrics = response.read().decode("utf-8") +def wait_for_metric_families(service_proxy_path: str, families: list[str]) -> None: + result = run_command(["kubectl", "get", "--raw", service_proxy_path], capture="out", check=True) - missing = [family for family in families if family not in metrics] + missing = [family for family in families if family not in result.stdout] if missing: raise AssertionError("Tekton metric families are not available yet: {}".format(", ".join(missing))) @@ -72,39 +68,39 @@ def setup_tekton(): name = f"{action}-{definition}-run" _wait_for_resource(kind=kind, name=name, condition="Succeeded", namespace="tekton-pipelines") + WaitFor( + wait_for_metric_families, + attempts=60, + wait=2, + args=( + ('/api/v1/namespaces/tekton-pipelines/services/http:tekton-triggers-controller:9000/proxy/metrics'), + [ + 'controller_clusterinterceptor_count', + 'controller_clustertriggerbinding_count', + 'controller_eventlistener_count', + 'controller_triggerbinding_count', + 'controller_triggertemplate_count', + ], + ), + )() + @pytest.fixture(scope='session') def dd_environment(): - with kind_run(conditions=[setup_tekton], sleep=10) as kubeconfig, ExitStack() as stack: - instances = {} - - pipeline_host, pipeline_port = stack.enter_context( - port_forward(kubeconfig, 'tekton-pipelines', 9090, 'service', 'tekton-pipelines-controller') - ) - instances['pipelines_controller_endpoint'] = f'http://{pipeline_host}:{pipeline_port}/metrics' - - trigger_host, trigger_port = stack.enter_context( - port_forward(kubeconfig, 'tekton-pipelines', 9000, 'service', 'tekton-triggers-controller') - ) - instances['triggers_controller_endpoint'] = f'http://{trigger_host}:{trigger_port}/metrics' - - WaitFor( - wait_for_metric_families, - attempts=60, - wait=2, - args=( - instances['triggers_controller_endpoint'], - [ - 'controller_clusterinterceptor_count', - 'controller_clustertriggerbinding_count', - 'controller_eventlistener_count', - 'controller_triggerbinding_count', - 'controller_triggertemplate_count', - ], + with kind_run(conditions=[setup_tekton], sleep=10) as kubeconfig: + instances = { + 'pipelines_controller_endpoint': ( + 'http://tekton-pipelines-controller.tekton-pipelines.svc.cluster.local:9090/metrics' ), - )() + 'triggers_controller_endpoint': ( + 'http://tekton-triggers-controller.tekton-pipelines.svc.cluster.local:9000/metrics' + ), + } - yield {'instances': [instances]} + yield ( + {'instances': [instances]}, + {'agent_type': 'kubernetes', 'kubernetes': {'kubeconfig': kubeconfig}}, + ) @pytest.fixture