diff --git a/github-runner-manager/pyproject.toml b/github-runner-manager/pyproject.toml index 62c9b0d798..ca4a0fea8f 100644 --- a/github-runner-manager/pyproject.toml +++ b/github-runner-manager/pyproject.toml @@ -79,9 +79,15 @@ select = ["E", "W", "F", "C", "N", "R", "D", "H"] # Ignore W503,E203 because using black creates errors with this # Ignore D107 Missing docstring in __init__ ignore = ["W503", "D107", "E203"] -# D100, D101, D102, D103, D104: Ignore docstring style issues in tests -# DCO020, DCO030, DCO050: Ignore docstring argument,returns,raises sections in tests -per-file-ignores = ["tests/*:D100,D101,D102,D103,D104,D205,D212, DCO020, DCO030, DCO050"] +per-file-ignores = [ + # Ignore factory methods attributes docstring + "tests/unit/factories/*:DCO060", + # Ignore no return values (DCO031) in docstring for abstract methods + "src/github_runner_manager/manager/cloud_runner_manager.py:DCO031", + # DCO020, DCO030, DCO050, DCO060: Ignore docstring argument, returns, raises, attribute + # sections in tests + "tests/*:D100,D101,D102,D103,D104,D205,D212,DCO020,DCO030,DCO050,DCO060", +] docstring-convention = "google" # Check for properly formatted copyright header in each file copyright-check = "True" diff --git a/github-runner-manager/src/github_runner_manager/github_client.py b/github-runner-manager/src/github_runner_manager/github_client.py index d3ee48c83f..61e42532b8 100644 --- a/github-runner-manager/src/github_runner_manager/github_client.py +++ b/github-runner-manager/src/github_runner_manager/github_client.py @@ -23,11 +23,7 @@ from requests import RequestException from typing_extensions import assert_never -from github_runner_manager.configuration.github import ( - GitHubOrg, - GitHubPath, - GitHubRepo, -) +from github_runner_manager.configuration.github import GitHubOrg, GitHubPath, GitHubRepo from github_runner_manager.manager.models import InstanceID from github_runner_manager.platform.platform_provider import ( DeleteRunnerBusyError, diff --git a/github-runner-manager/src/github_runner_manager/manager/cloud_runner_manager.py b/github-runner-manager/src/github_runner_manager/manager/cloud_runner_manager.py index 7cd7641110..06532e2c27 100644 --- a/github-runner-manager/src/github_runner_manager/manager/cloud_runner_manager.py +++ b/github-runner-manager/src/github_runner_manager/manager/cloud_runner_manager.py @@ -35,20 +35,6 @@ class HealthState(Enum): UNHEALTHY = auto() UNKNOWN = auto() - @staticmethod - def from_value(health: bool | None) -> "HealthState": - """Create from a health value. - - Args: - health: The health value as boolean or None. - - Returns: - The health state. - """ - if health is None: - return HealthState.UNKNOWN - return HealthState.HEALTHY if health else HealthState.UNHEALTHY - class CloudRunnerState(str, Enum): """Represent state of the instance hosting the runner. @@ -275,11 +261,25 @@ def get_runners(self) -> Sequence[CloudRunnerInstance]: """Get cloud self-hosted runners.""" @abc.abstractmethod - def delete_runner(self, instance_id: InstanceID) -> RunnerMetrics | None: - """Delete self-hosted runner. + def delete_vms(self, instance_ids: Sequence[InstanceID]) -> list[InstanceID]: + """Delete cloud VM instances. + + Args: + instance_ids: The ID of the VMs to request deletion. + + Returns: + The deleted instance IDs. + """ + + @abc.abstractmethod + def extract_metrics(self, instance_ids: Sequence[InstanceID]) -> list[RunnerMetrics]: + """Extract metrics from cloud VMs. Args: - instance_id: The instance id of the runner to delete. + instance_ids: The VM instance IDs to fetch the metrics from. + + Returns: + The fetched runner metrics. """ @abc.abstractmethod diff --git a/github-runner-manager/src/github_runner_manager/manager/models.py b/github-runner-manager/src/github_runner_manager/manager/models.py index ef81f38175..fe4b591aee 100644 --- a/github-runner-manager/src/github_runner_manager/manager/models.py +++ b/github-runner-manager/src/github_runner_manager/manager/models.py @@ -32,8 +32,8 @@ class InstanceID: """ prefix: str - reactive: bool | None suffix: str + reactive: bool | None = None @property def name(self) -> str: diff --git a/github-runner-manager/src/github_runner_manager/manager/runner_manager.py b/github-runner-manager/src/github_runner_manager/manager/runner_manager.py index aecd28ba39..e6509d51f3 100644 --- a/github-runner-manager/src/github_runner_manager/manager/runner_manager.py +++ b/github-runner-manager/src/github_runner_manager/manager/runner_manager.py @@ -23,11 +23,9 @@ from github_runner_manager.metrics import events as metric_events from github_runner_manager.metrics import github as github_metrics from github_runner_manager.metrics import runner as runner_metrics -from github_runner_manager.metrics.reconcile import CLEANED_RUNNERS_TOTAL from github_runner_manager.metrics.runner import RunnerMetrics from github_runner_manager.openstack_cloud.constants import CREATE_SERVER_TIMEOUT from github_runner_manager.platform.platform_provider import ( - DeleteRunnerBusyError, PlatformApiError, PlatformProvider, PlatformRunnerHealth, @@ -82,27 +80,33 @@ class RunnerInstance: platform_state: PlatformRunnerState | None cloud_state: CloudRunnerState - def __init__( - self, + @classmethod + def from_cloud_and_platform_health( + cls, cloud_instance: CloudRunnerInstance, platform_health_state: PlatformRunnerHealth | None, - ): + ) -> "RunnerInstance": """Construct an instance. Args: cloud_instance: Information on the cloud instance. platform_health_state: Health state in the platform provider. + + Returns: + The RunnerInstance instantiated from cloud instance and platform state. """ - self.name = cloud_instance.name - self.instance_id = cloud_instance.instance_id - self.metadata = cloud_instance.metadata - self.health = cloud_instance.health - self.platform_state = ( - PlatformRunnerState.from_platform_health(platform_health_state) - if platform_health_state is not None - else None + return cls( + name=cloud_instance.name, + instance_id=cloud_instance.instance_id, + metadata=cloud_instance.metadata, + health=cloud_instance.health, + platform_state=( + PlatformRunnerState.from_platform_health(platform_health_state) + if platform_health_state is not None + else None + ), + cloud_state=cloud_instance.state, ) - self.cloud_state = cloud_instance.state class RunnerManager: @@ -183,7 +187,7 @@ def get_runners(self) -> tuple[RunnerInstance, ...]: health_runners_map = {runner.identity.instance_id: runner for runner in runners_health} for cloud_runner in cloud_runners: if cloud_runner.instance_id not in health_runners_map: - runner_instance = RunnerInstance(cloud_runner, None) + runner_instance = RunnerInstance.from_cloud_and_platform_health(cloud_runner, None) runner_instance.health = HealthState.UNKNOWN runner_instances.append(runner_instance) continue @@ -194,7 +198,9 @@ def get_runners(self) -> tuple[RunnerInstance, ...]: cloud_runner.health = HealthState.HEALTHY else: cloud_runner.health = HealthState.UNHEALTHY - runner_instance = RunnerInstance(cloud_runner, health_runner) + runner_instance = RunnerInstance.from_cloud_and_platform_health( + cloud_runner, health_runner + ) runner_instances.append(runner_instance) return cast(tuple[RunnerInstance], tuple(runner_instances)) @@ -268,7 +274,9 @@ def _cleanup_resources( ) cloud_runners = self._cloud.get_runners() logger.info("cleanup cloud_runners %s", cloud_runners) - runners_health_response = self._platform.get_runners_health(cloud_runners) + runners_health_response = self._platform.get_runners_health( + requested_runners=cloud_runners + ) logger.info("cleanup health_response %s", runners_health_response) # Clean dangling resources in the cloud @@ -295,7 +303,7 @@ def _cleanup_resources( ) ) - if maximum_runners_to_delete: + if maximum_runners_to_delete is not None: cloud_runners_to_delete.sort( key=partial(_runner_deletion_sort_key, health_runners_map) ) @@ -313,58 +321,70 @@ def _delete_cloud_runners( runners_health: Sequence[PlatformRunnerHealth], delete_busy_runners: bool = False, ) -> Iterable[runner_metrics.RunnerMetrics]: - """Delete runners in the platform ant the cloud. + """Delete runners in the platform and the cloud. If delete_busy_runners is False, when the platform provider fails in deleting the runner because it can be busy, will mean that that runner should not be deleted. + + Runners without health information should not be deleted. """ - extracted_runner_metrics = [] - health_runners_map = {health.identity.instance_id: health for health in runners_health} - for cloud_runner in cloud_runners: - logging.info("Trying to delete cloud_runner %s", cloud_runner) - runner_health = health_runners_map.get(cloud_runner.instance_id) - if runner_health and runner_health.runner_in_platform: - try: - self._platform.delete_runner(runner_health.identity) - except DeleteRunnerBusyError: - if not delete_busy_runners: - logger.warning( - "Skipping deletion as the runner is busy. %s", cloud_runner.instance_id - ) - continue - logger.info("Deleting busy runner: %s", cloud_runner.instance_id) - except PlatformApiError as exc: - if not delete_busy_runners: - logger.warning( - "Failed to delete platform runner %s. %s. Skipping.", - cloud_runner.instance_id, - exc, - ) - continue - logger.warning( - "Deleting runner: %s after platform failure %s.", - cloud_runner.instance_id, - exc, - ) + if not cloud_runners: + return [] - logging.info("Delete runner in cloud: %s", cloud_runner.instance_id) - runner_metric = self._cloud.delete_runner(cloud_runner.instance_id) - CLEANED_RUNNERS_TOTAL.labels(self.manager_name).inc(1) - if not runner_metric: - logger.error("No metrics returned after deleting %s", cloud_runner.instance_id) - else: - extracted_runner_metrics.append(runner_metric) - return extracted_runner_metrics + runner_identity_map = { + health_info.identity.instance_id: health_info.identity + for health_info in runners_health + } + platform_runner_ids_to_delete = [ + # The runner_id cannot be None due to the if condition. the type system + # isn't able to catch that. + cast(str, runner_identity_map[runner.instance_id].metadata.runner_id) + for runner in cloud_runners + if runner.instance_id in runner_identity_map + and runner_identity_map[runner.instance_id].metadata.runner_id + ] + logger.info("Deleting runners from platform: %s", platform_runner_ids_to_delete) + deleted_runner_ids = self._platform.delete_runners( + runner_ids=platform_runner_ids_to_delete + ) + logger.info( + "Deleted runners from platform: %s (diff: %s)", + deleted_runner_ids, + set(platform_runner_ids_to_delete) - set(deleted_runner_ids), + ) + + logger.info("Cloud runners: %s", cloud_runners) + cloud_vm_ids_to_delete = [ + runner.instance_id + for runner in cloud_runners + # We can delete all VMs if delete_busy_runners is True + if delete_busy_runners + # We can delete the VM if no runner is associated with it + or not runner.metadata.runner_id + # We can delete the VM if it has been deleted from the Platform provider. + or runner.metadata.runner_id in deleted_runner_ids + ] + logger.info("Extracting metrics from cloud VMs: %s", cloud_vm_ids_to_delete) + extracted_metrics = self._cloud.extract_metrics(instance_ids=cloud_vm_ids_to_delete) + logger.info("Extracted metrics from cloud VMs: %s", extracted_metrics) + logger.info("Deleting VMs %s", cloud_vm_ids_to_delete) + deleted_vm_ids = self._cloud.delete_vms(instance_ids=cloud_vm_ids_to_delete) + logger.info( + "Deleted VMs: %s, (diff: %s)", + deleted_vm_ids, + set(cloud_vm_ids_to_delete) - set(deleted_vm_ids), + ) + return tuple(extracted_metrics) def _clean_platform_runners(self, runners: list[RunnerIdentity]) -> None: """Clean the specified runners in the platform.""" - for runner in runners: - try: - self._platform.delete_runner(runner) - except DeleteRunnerBusyError: - logger.warning("Tried to delete busy runner in cleanup %s", runner) - except PlatformApiError: - logger.warning("Failed to delete platform runner %s", runner) + if not runners: + return + + runner_ids_to_delete = [ + runner.metadata.runner_id for runner in runners if runner.metadata.runner_id + ] + self._platform.delete_runners(runner_ids=runner_ids_to_delete) @staticmethod def _spawn_runners( @@ -525,7 +545,7 @@ def _create_runner(args: _CreateRunnerArgs) -> InstanceID: ) except RunnerError: logger.warning("Deleting runner %s from platform after creation failed", instance_id) - args.platform_provider.delete_runner(runner_info.identity) + args.platform_provider.delete_runners(runner_ids=[args.metadata.runner_id]) raise return instance_id diff --git a/github-runner-manager/src/github_runner_manager/manager/runner_scaler.py b/github-runner-manager/src/github_runner_manager/manager/runner_scaler.py index b0f8e44ef4..e451c06aa6 100644 --- a/github-runner-manager/src/github_runner_manager/manager/runner_scaler.py +++ b/github-runner-manager/src/github_runner_manager/manager/runner_scaler.py @@ -94,7 +94,7 @@ class _ReconcileMetricData: start_timestamp: float end_timestamp: float metric_stats: IssuedMetricEventsStats - runner_list: tuple[RunnerInstance] + runner_list: tuple[RunnerInstance, ...] flavor: str expected_runner_quantity: int @@ -375,7 +375,7 @@ def _reconcile_non_reactive(self, expected_quantity: int) -> _ReconcileResult: return _ReconcileResult(runner_diff=runner_diff, metric_stats=metric_stats) @staticmethod - def _log_runners(runner_list: tuple[RunnerInstance]) -> None: + def _log_runners(runner_list: tuple[RunnerInstance, ...]) -> None: """Log information about the runners found. Args: @@ -446,7 +446,6 @@ def _issue_reconciliation_metric( IDLE_RUNNERS_COUNT.labels(manager_name).set(len(idle_runners)) try: - metric_events.issue_event( metric_events.Reconciliation( timestamp=time.time(), diff --git a/github-runner-manager/src/github_runner_manager/metrics/runner.py b/github-runner-manager/src/github_runner_manager/metrics/runner.py index 6426819f1a..983ca1dfd0 100644 --- a/github-runner-manager/src/github_runner_manager/metrics/runner.py +++ b/github-runner-manager/src/github_runner_manager/metrics/runner.py @@ -3,13 +3,13 @@ """Classes and function to extract the metrics from storage and issue runner metrics events.""" +import concurrent.futures import io import json import logging from dataclasses import dataclass -from datetime import datetime from json import JSONDecodeError -from typing import Optional, Type +from typing import Optional, Sequence, Type import paramiko import paramiko.ssh_exception @@ -18,7 +18,6 @@ from github_runner_manager.errors import IssueMetricEventError, RunnerMetricsError, SSHError from github_runner_manager.manager.cloud_runner_manager import ( - CloudRunnerInstance, PostJobMetrics, PreJobMetrics, RunnerMetrics, @@ -31,6 +30,7 @@ PRE_JOB_METRICS_FILE_NAME, RUNNER_INSTALLED_TS_FILE_NAME, ) +from github_runner_manager.openstack_cloud.openstack_cloud import OpenstackCloud, OpenstackInstance logger = logging.getLogger(__name__) @@ -41,45 +41,130 @@ class PullFileError(Exception): """Represents an error while pulling a file from the runner instance.""" -def pull_runner_metrics(instance_id: InstanceID, ssh_conn: SSHConnection) -> "PulledMetrics": +@dataclass +class _PullRunnerMetricsConfig: + """Configurations for pulling runner metrics from a VM. + + Attributes: + cloud_service: The OpenStack cloud service. + instance_id: The instance ID to fetch the runner metric from. + """ + + cloud_service: OpenstackCloud + instance_id: InstanceID + + +def pull_runner_metrics( + cloud_service: OpenstackCloud, instance_ids: Sequence[InstanceID] +) -> "list[PulledMetrics]": """Pull metrics from runner. + This function uses multiprocessing to fetch metrics in parallel. + Args: - instance_id: The name of the runner. - ssh_conn: The SSH connection to the runner. + cloud_service: The OpenStack cloud service. + instance_ids: The instance IDs to fetch the metrics from. Returns: Metrics pulled from the instance. """ - logger.debug("Pulling metrics for %s", instance_id) - pulled_metrics = PulledMetrics() + if not instance_ids: + return [] + pull_metrics_configs = [ + _PullRunnerMetricsConfig(cloud_service=cloud_service, instance_id=instance_id) + for instance_id in instance_ids + ] + pulled_metrics: list[PulledMetrics] = [] + with concurrent.futures.ThreadPoolExecutor(max_workers=min(len(instance_ids), 30)) as executor: + future_to_pull_metrics_config = { + executor.submit(_pull_runner_metrics, config): config + for config in pull_metrics_configs + } + for future in concurrent.futures.as_completed(future_to_pull_metrics_config): + pull_config = future_to_pull_metrics_config[future] + metric = future.result() + if not metric: + logger.warning("No metrics pulled for %s", pull_config.instance_id) + else: + pulled_metrics.append(metric) - try: - pulled_metrics.runner_installed = ssh_pull_file( - ssh_conn=ssh_conn, - remote_path=str(RUNNER_INSTALLED_TS_FILE_NAME), - max_size=MAX_METRICS_FILE_SIZE, - ) - pulled_metrics.pre_job_metrics = ssh_pull_file( - ssh_conn=ssh_conn, - remote_path=str(PRE_JOB_METRICS_FILE_NAME), - max_size=MAX_METRICS_FILE_SIZE, - ) - pulled_metrics.post_job_metrics = ssh_pull_file( - ssh_conn=ssh_conn, - remote_path=str(POST_JOB_METRICS_FILE_NAME), - max_size=MAX_METRICS_FILE_SIZE, + return pulled_metrics + + +def _pull_runner_metrics(pull_config: _PullRunnerMetricsConfig) -> "PulledMetrics | None": + """Pull metrics from a single runner via SSH file pull. + + Args: + pull_config: Configurations for pulling the runner metrics. + + Returns: + PulledMetrics if metrics were available. None otherwise. + """ + instance = pull_config.cloud_service.get_instance(instance_id=pull_config.instance_id) + if not instance: + logger.warning( + "Skipping fetching metrics, instance not found: %s", pull_config.instance_id ) - except PullFileError as exc: + return None + + runner_installed, pre_job_metrics, post_job_metrics = "", "", "" + try: + with pull_config.cloud_service.get_ssh_connection(instance=instance) as ssh_conn: + try: + runner_installed = _ssh_pull_file( + ssh_conn=ssh_conn, + remote_path=str(RUNNER_INSTALLED_TS_FILE_NAME), + max_size=MAX_METRICS_FILE_SIZE, + ) + except PullFileError as exc: + logger.warning( + "Failed to pull runner_installed metrics for %s: %s.", + pull_config.instance_id, + exc, + ) + try: + pre_job_metrics = _ssh_pull_file( + ssh_conn=ssh_conn, + remote_path=str(PRE_JOB_METRICS_FILE_NAME), + max_size=MAX_METRICS_FILE_SIZE, + ) + except PullFileError as exc: + logger.warning( + "Failed to pull pre_job metrics for %s: %s.", + pull_config.instance_id, + exc, + ) + try: + post_job_metrics = _ssh_pull_file( + ssh_conn=ssh_conn, + remote_path=str(POST_JOB_METRICS_FILE_NAME), + max_size=MAX_METRICS_FILE_SIZE, + ) + except PullFileError as exc: + logger.warning( + "Failed to pull post_job metrics for %s: %s.", + pull_config.instance_id, + exc, + ) + except SSHError: logger.warning( - "Failed to pull metrics for %s: %s . Will not be able to issue all metrics", - instance_id, - exc, + "Failed to create SSH connection for pulling metrics: %s", instance.instance_id ) - return pulled_metrics + return None + + return ( + PulledMetrics( + instance=instance, + runner_installed=runner_installed, + pre_job_metrics=pre_job_metrics, + post_job_metrics=post_job_metrics, + ) + if (runner_installed or pre_job_metrics or post_job_metrics) + else None + ) -def ssh_pull_file(ssh_conn: SSHConnection, remote_path: str, max_size: int) -> str: +def _ssh_pull_file(ssh_conn: SSHConnection, remote_path: str, max_size: int) -> str: """Pull file from the runner instance. Args: @@ -140,33 +225,29 @@ def ssh_pull_file(ssh_conn: SSHConnection, remote_path: str, max_size: int) -> s return value -@dataclass +@dataclass(frozen=True) class PulledMetrics: """Metrics pulled from a runner. Attributes: + instance: The instance in which the metrics were pulled from. runner_installed: String with the runner-installed file. pre_job_metrics: String with the pre-job-metrics file. post_job_metrics: String with the post-job-metrics file. """ + instance: OpenstackInstance runner_installed: str | None = None pre_job_metrics: str | None = None post_job_metrics: str | None = None - def to_runner_metrics( - self, instance: CloudRunnerInstance, installation_start: datetime - ) -> RunnerMetrics | None: - """. - - Args: - instance: Cloud runner instance. - installation_start: Creation time of the runner. + def to_runner_metrics(self) -> RunnerMetrics | None: + """Convert PulledMetrics to RunnerMetrics instance. Returns: The RunnerMetrics object for the runner or None if it can not be built. """ - instance_id = instance.instance_id + instance_id = self.instance.instance_id if self.runner_installed is None: logger.error( "Invalid pulled metrics. No runner_installed information for %s.", instance_id @@ -203,7 +284,7 @@ def to_runner_metrics( try: return RunnerMetrics( - installation_start_timestamp=installation_start.timestamp(), + installation_start_timestamp=self.instance.created_at.timestamp(), installed_timestamp=float(self.runner_installed), pre_job=( # pylint: disable=not-a-mapping PreJobMetrics(**pre_job_metrics) if pre_job_metrics else None @@ -212,11 +293,14 @@ def to_runner_metrics( PostJobMetrics(**post_job_metrics) if post_job_metrics else None ), instance_id=instance_id, - metadata=instance.metadata, + metadata=self.instance.metadata, ) except ValueError: logger.exception( - "Error creating RunnerMetrics %s, %s, %s", instance_id, installation_start, self + "Error creating RunnerMetrics %s, %s, %s", + instance_id, + self.instance.created_at, + self, ) return None diff --git a/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_cloud.py b/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_cloud.py index d81cc71caa..083eb0a549 100644 --- a/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_cloud.py +++ b/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_cloud.py @@ -2,6 +2,7 @@ # See LICENSE file for licensing details. """Class for accessing OpenStack API for managing servers.""" +import concurrent.futures import contextlib import copy import functools @@ -12,7 +13,7 @@ from datetime import datetime, timezone from functools import reduce from pathlib import Path -from typing import Any, Callable, Iterable, Iterator, ParamSpec, TypeVar, cast +from typing import Any, Callable, Iterable, Iterator, ParamSpec, Sequence, TypeVar, cast import keystoneauth1.exceptions import openstack @@ -75,6 +76,29 @@ _MIN_KEYPAIR_AGE_IN_SECONDS_BEFORE_DELETION = 60 +class DeleteVMError(openstack.exceptions.SDKException): + """Represents an error while deleting a VM instance. + + Attributes: + instance_id: The instance ID that was failed to delete. + """ + + instance_id: InstanceID + + def __init__( + self, instance_id: InstanceID, message: str | None = None, extra_data: Any = None + ): + """Initialize the OpenstackVMDeleteError. + + Args: + instance_id: The instance ID of the failed delete VM. + message: The delete error message for parent SDKException. + extra_data: Extra data for parent SDKException if any. + """ + self.instance_id = instance_id + super().__init__(message, extra_data) + + @dataclass class OpenstackInstance: """Represents an OpenStack instance. @@ -158,6 +182,42 @@ def exception_handling_wrapper(*args: P.args, **kwargs: P.kwargs) -> T: return exception_handling_wrapper +@dataclass +class _DeleteVMConfig: + """Configurations for deleting a VM. + + Attributes: + instance_id: The ID of the VM to request deletion. + credentials: The OpenStack connection credentials. + max_api_version: The OpenStack maximum compute API version. + keys_dir: The path to the directory in which the SSH key files are stored. + wait: Whether to wait for the VM delete to complete. + timeout: Timeout in seconds for VM deletion to complete. + """ + + instance_id: InstanceID + credentials: OpenStackCredentials + max_api_version: str + keys_dir: Path + wait: bool = False + timeout: int = 10 * 60 + + +@dataclass +class _DeleteKeypairConfig: + """Configurations for deleting an OpenStack keypair. + + Attributes: + keys_dir: The path to the directory in which the SSH key files are stored. + instance_id: The instance ID of the key owner. + conn: The OpenStack connection instance. + """ + + keys_dir: Path + instance_id: InstanceID + conn: OpenstackConnection + + class OpenstackCloud: """Client to interact with OpenStack cloud. @@ -240,11 +300,22 @@ def launch_instance( "Attempting clean up of openstack server %s that timeout during creation", instance_id, ) - self._delete_instance(conn, instance_id) + OpenstackCloud._delete_instance( + _DeleteVMConfig( + instance_id=instance_id, + credentials=self._credentials, + max_api_version=self._max_compute_api_version, + keys_dir=self._ssh_key_dir, + ) + ) raise OpenStackError(f"Timeout creating openstack server {instance_id}") from err except openstack.exceptions.SDKException as err: logger.exception("Failed to create openstack server %s", instance_id) - self._delete_keypair(conn, instance_id) + OpenstackCloud._delete_keypair( + _DeleteKeypairConfig( + keys_dir=self._ssh_key_dir, instance_id=instance_id, conn=conn + ) + ) raise OpenStackError(f"Failed to create openstack server {instance_id}") from err return OpenstackInstance(server, self.prefix) @@ -267,37 +338,103 @@ def get_instance(self, instance_id: InstanceID) -> OpenstackInstance | None: return OpenstackInstance(server, self.prefix) return None - @_catch_openstack_errors - def delete_instance(self, instance_id: InstanceID) -> None: + @staticmethod + def _delete_instance(delete_config: _DeleteVMConfig) -> bool: """Delete a openstack instance. Args: - instance_id: The instance ID of the instance to delete. - """ - logger.info("Deleting openstack server with %s", instance_id) + delete_config: The configuration used to delete a cloud VM instance. - with self._get_openstack_connection() as conn: - self._delete_instance(conn, instance_id) + Raises: + DeleteVMError: If there was an error deleting the VM instance. + """ + with openstack.connect( + auth_url=delete_config.credentials.auth_url, + project_name=delete_config.credentials.project_name, + username=delete_config.credentials.username, + password=delete_config.credentials.password, + region_name=delete_config.credentials.region_name, + user_domain_name=delete_config.credentials.user_domain_name, + project_domain_name=delete_config.credentials.project_domain_name, + compute_api_version=delete_config.max_api_version, + ) as conn: + try: + logger.info("Deleting server %s", delete_config.instance_id.name) + deleted = conn.delete_server( + name_or_id=delete_config.instance_id.name, + wait=delete_config.wait, + timeout=delete_config.timeout, + ) + logger.info( + "Deleted server %s (true delete: %s)", delete_config.instance_id.name, deleted + ) + except ( + openstack.exceptions.SDKException, + openstack.exceptions.ResourceTimeout, + ) as exc: + raise DeleteVMError( + instance_id=delete_config.instance_id, + message=f"Failed to delete server {delete_config.instance_id.name}", + ) from exc + + OpenstackCloud._delete_keypair( + _DeleteKeypairConfig( + keys_dir=delete_config.keys_dir, + instance_id=delete_config.instance_id, + conn=conn, + ) + ) - def _delete_instance(self, conn: OpenstackConnection, instance_id: InstanceID) -> None: - """Delete a openstack instance. + return deleted - Raises: - OpenStackError: Unable to delete OpenStack server. + def delete_instances( + self, instance_ids: Sequence[InstanceID], wait: bool = False, timeout: int = 60 * 10 + ) -> list[InstanceID]: + """Delete Openstack VM instances. Args: - conn: The openstack connection to use. - instance_id: The full name of the server. + instance_ids: The VM instance IDs to requeest deletion. + wait: Whether to wait for VM deletion to complete. + timeout: Timeout in seconds to wait for VM deletion to complete. + + Returns: + The deleted VM instance IDs if wait is True, deleted requested VM instance IDs + otherwise. """ - try: - res = conn.delete_server(name_or_id=instance_id.name) - logger.info("openstack delete result for %s: %s", instance_id, res) - self._delete_keypair(conn, instance_id) - except ( - openstack.exceptions.SDKException, - openstack.exceptions.ResourceTimeout, - ) as err: - raise OpenStackError(f"Failed to remove openstack runner {instance_id}") from err + deleted_instance_ids: list[InstanceID] = [] + + # Guard no instance IDs since multiprocessing Pool may raise an exception. + if not instance_ids: + return deleted_instance_ids + + delete_configs = [ + _DeleteVMConfig( + instance_id=instance_id, + credentials=self._credentials, + max_api_version=self._max_compute_api_version, + keys_dir=self._ssh_key_dir, + wait=wait, + timeout=timeout, + ) + for instance_id in instance_ids + ] + with concurrent.futures.ThreadPoolExecutor( + max_workers=min(len(instance_ids), 30) + ) as executor: + future_to_delete_instance_config = { + executor.submit(OpenstackCloud._delete_instance, config): config + for config in delete_configs + } + for future in concurrent.futures.as_completed(future_to_delete_instance_config): + delete_config = future_to_delete_instance_config[future] + try: + if not future.result(): + continue + deleted_instance_ids.append(delete_config.instance_id) + except DeleteVMError as exc: + logger.error("Failed to delete OpenStack VM instance: %s", exc.instance_id) + + return deleted_instance_ids @_catch_openstack_errors @contextlib.contextmanager @@ -314,7 +451,7 @@ def get_ssh_connection(self, instance: OpenstackInstance) -> Iterator[SSHConnect Yields: SSH connection object. """ - key_path = self._get_key_path(instance.instance_id.name) + key_path = self._get_key_path(instance.instance_id) if not key_path.exists(): raise KeyfileError( @@ -463,7 +600,13 @@ def _cleanup_openstack_keypairs( if str(key.name) in exclude_keys: continue try: - self._delete_keypair(conn, InstanceID.build_from_name(self.prefix, key.name)) + OpenstackCloud._delete_keypair( + _DeleteKeypairConfig( + keys_dir=self._ssh_key_dir, + instance_id=InstanceID.build_from_name(self.prefix, key.name), + conn=conn, + ) + ) except openstack.exceptions.SDKException: logger.warning( "Unable to delete OpenStack keypair associated with deleted key file %s ", @@ -570,23 +713,32 @@ def _setup_keypair( key_path.chmod(0o400) return keypair - def _delete_keypair(self, conn: OpenstackConnection, instance_id: InstanceID) -> None: + @staticmethod + def _delete_keypair(delete_keypair_config: _DeleteKeypairConfig) -> None: """Delete OpenStack keypair. Args: - conn: The connection object to access OpenStack cloud. - instance_id: The name of the keypair. + delete_keypair_config: Configurations for deleting the KeyPair. """ - logger.debug("Deleting keypair for %s", instance_id) + logger.info("Deleting key: %s", delete_keypair_config.instance_id) try: # Keypair have unique names, access by ID is not needed. - if not conn.delete_keypair(instance_id.name): - logger.warning("Unable to delete keypair for %s", instance_id) + if not delete_keypair_config.conn.delete_keypair( + delete_keypair_config.instance_id.name + ): + logger.warning("Failed to delete key: %s", delete_keypair_config.instance_id.name) + return except (openstack.exceptions.SDKException, openstack.exceptions.ResourceTimeout): - logger.warning("Unable to delete keypair for %s", instance_id, stack_info=True) + logger.warning( + "Error attempting to delete key: %s", + delete_keypair_config.instance_id.name, + stack_info=True, + ) + return - key_path = self._get_key_path(instance_id.name) + key_path = delete_keypair_config.keys_dir / f"{delete_keypair_config.instance_id}.key" key_path.unlink(missing_ok=True) + logger.info("Deleted key: %s", delete_keypair_config.instance_id) @staticmethod def _ensure_security_group( diff --git a/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_runner_manager.py b/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_runner_manager.py index 3a690c64c0..0ceed7c1d6 100644 --- a/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_runner_manager.py +++ b/github-runner-manager/src/github_runner_manager/openstack_cloud/openstack_runner_manager.py @@ -16,18 +16,14 @@ MissingServerConfigError, OpenStackError, RunnerCreateError, - SSHError, ) from github_runner_manager.manager.cloud_runner_manager import ( CloudRunnerInstance, CloudRunnerManager, CloudRunnerState, + RunnerMetrics, ) -from github_runner_manager.manager.models import ( - InstanceID, - RunnerContext, - RunnerIdentity, -) +from github_runner_manager.manager.models import InstanceID, RunnerContext, RunnerIdentity from github_runner_manager.manager.runner_manager import HealthState from github_runner_manager.metrics import runner as runner_metrics from github_runner_manager.openstack_cloud.constants import ( @@ -164,73 +160,18 @@ def cleanup(self) -> None: """Cleanup runner and resource on the cloud.""" self._openstack_cloud.cleanup() - def _build_cloud_runner_instance( - self, instance: OpenstackInstance, healthy: bool | None = None - ) -> CloudRunnerInstance: + def _build_cloud_runner_instance(self, instance: OpenstackInstance) -> CloudRunnerInstance: """Build a new cloud runner instance from an openstack instance.""" metadata = instance.metadata return CloudRunnerInstance( name=instance.instance_id.name, metadata=metadata, instance_id=instance.instance_id, - health=HealthState.from_value(healthy), + health=HealthState.UNKNOWN, state=CloudRunnerState.from_openstack_server_status(instance.status), created_at=instance.created_at, ) - def delete_runner(self, instance_id: InstanceID) -> runner_metrics.RunnerMetrics | None: - """Delete self-hosted runners. - - Args: - instance_id: The instance id of the runner to delete. - - Returns: - Any metrics collected during the deletion of the runner. - """ - logger.debug("Delete instance %s", instance_id) - instance = self._openstack_cloud.get_instance(instance_id) - if instance is None: - logger.warning( - "Unable to delete instance %s as it is not found", - instance_id, - ) - return None - - pulled_metrics = self._delete_runner(instance) - logger.debug( - "Metrics extracted, deleting instance %s %s", instance_id, instance.instance_id - ) - logger.debug("Instance deleted successfully %s %s", instance_id, instance.instance_id) - logger.debug("Extract metrics for runner %s %s", instance_id, instance.instance_id) - cloud_instance = self._build_cloud_runner_instance(instance) - return pulled_metrics.to_runner_metrics(cloud_instance, instance.created_at) - - def _delete_runner(self, instance: OpenstackInstance) -> runner_metrics.PulledMetrics: - """Delete self-hosted runners by openstack instance. - - Args: - instance: The OpenStack instance. - """ - pulled_metrics = runner_metrics.PulledMetrics() - try: - with self._openstack_cloud.get_ssh_connection(instance) as ssh_conn: - pulled_metrics = runner_metrics.pull_runner_metrics(instance.instance_id, ssh_conn) - except SSHError: - logger.exception( - "Failed to get SSH connection while removing %s", instance.instance_id - ) - logger.warning( - "Skipping runner remove script for %s due to SSH issues", instance.instance_id - ) - - try: - self._openstack_cloud.delete_instance(instance.instance_id) - except OpenStackError: - logger.exception( - "Unable to delete openstack instance for runner %s", instance.instance_id - ) - return pulled_metrics - def _generate_cloud_init(self, runner_context: RunnerContext) -> str: """Generate cloud init userdata. @@ -317,3 +258,37 @@ def _get_repo_policy_compliance_client(self) -> RepoPolicyComplianceClient | Non service_config.repo_policy_compliance.token, ) return None + + def delete_vms( + self, instance_ids: Sequence[InstanceID], wait: bool = False, timeout: int = 60 * 10 + ) -> list[InstanceID]: + """Delete VMs. + + Args: + instance_ids: The ID of the VMs to request deletion. + wait: Whether to wait for the delete to be complete. + timeout: Timeout in seconds to wait for the deletion to complete. + + Returns: + The instance IDs requested for deletion. + """ + return self._openstack_cloud.delete_instances( + instance_ids=instance_ids, wait=wait, timeout=timeout + ) + + def extract_metrics(self, instance_ids: Sequence[InstanceID]) -> list[RunnerMetrics]: + """Extract metrics from cloud VMs. + + Args: + instance_ids: The ID of the VMs to fetch metrics from. + + Returns: + Metrics from VMs. + """ + return [ + converted_metrics + for pulled_metrics in runner_metrics.pull_runner_metrics( + cloud_service=self._openstack_cloud, instance_ids=instance_ids + ) + if (converted_metrics := pulled_metrics.to_runner_metrics()) + ] diff --git a/github-runner-manager/src/github_runner_manager/platform/github_provider.py b/github-runner-manager/src/github_runner_manager/platform/github_provider.py index 247979df5d..11eabf49ba 100644 --- a/github-runner-manager/src/github_runner_manager/platform/github_provider.py +++ b/github-runner-manager/src/github_runner_manager/platform/github_provider.py @@ -3,13 +3,19 @@ """Client for managing self-hosted runner on GitHub side.""" +import concurrent.futures import logging +from dataclasses import dataclass from enum import Enum from pydantic import HttpUrl -from github_runner_manager.configuration.github import GitHubConfiguration, GitHubRepo -from github_runner_manager.github_client import GithubClient, GithubRunnerNotFoundError +from github_runner_manager.configuration.github import GitHubConfiguration, GitHubPath, GitHubRepo +from github_runner_manager.github_client import ( + DeleteRunnerBusyError, + GithubClient, + GithubRunnerNotFoundError, +) from github_runner_manager.manager.models import ( InstanceID, RunnerContext, @@ -28,10 +34,25 @@ logger = logging.getLogger(__name__) +@dataclass +class _DeleteRunnerConfig: + """Configurations for deleting a runner. + + Attributes: + runner_id: The ID of the runner to delete. + path: The path (repository/org) in which the the runner was registered to. + github_client: The GitHub client to use to call delete runner. + """ + + runner_id: str + path: GitHubPath + github_client: GithubClient + + class GitHubRunnerPlatform(PlatformProvider): """Manage self-hosted runner on GitHub side.""" - def __init__(self, prefix: str, path: str, github_client: GithubClient): + def __init__(self, prefix: str, path: GitHubPath, github_client: GithubClient): """Construct the object. Args: @@ -143,17 +164,59 @@ def get_runners_health(self, requested_runners: list[RunnerIdentity]) -> Runners non_requested_runners=non_requested_runners, ) - def delete_runner(self, runner_identity: RunnerIdentity) -> None: - """Delete a runner from GitHub. + def delete_runners(self, runner_ids: list[str]) -> list[str]: + """Delete runners from GitHub. - This method will raise DeleteRunnerBusyError if the runner is not deletable, that is, - if it is busy. If the runner does not exist it will not fail. + This method will ignore DeleteRunnerBusyErrors and print a warning log. Args: - runner_identity: Identity of the runner to delete. + runner_ids: The GitHub runner IDs to delete. + + Returns: + The runner IDs that were deleted successfully. """ - logger.info("Delete runner in GitHub: %s", runner_identity) - self._client.delete_runner(self._path, int(runner_identity.metadata.runner_id)) + logger.info("Delete runners from GitHub provider: %s", runner_ids) + # Guard multiprocessing.Pool from having 0 processes which will raise an error. + if not runner_ids: + return [] + + delete_configs = [ + _DeleteRunnerConfig(runner_id=runner_id, path=self._path, github_client=self._client) + for runner_id in runner_ids + ] + deleted_runner_ids: list[str] = [] + with concurrent.futures.ThreadPoolExecutor( + max_workers=min(len(runner_ids), 30) + ) as executor: + future_to_delete_runner_config = { + executor.submit(GitHubRunnerPlatform._delete_runner, config): config + for config in delete_configs + } + for future in concurrent.futures.as_completed(future_to_delete_runner_config): + delete_config = future_to_delete_runner_config[future] + try: + future.result() + deleted_runner_ids.append(delete_config.runner_id) + except DeleteRunnerBusyError: + logger.warning( + "Delete runner attempt failed, busy runner: %s", + delete_config.runner_id, + ) + + return deleted_runner_ids + + @staticmethod + def _delete_runner(delete_runner_config: _DeleteRunnerConfig) -> None: + """Delete a single runner from GitHub. + + This method is a wrapper to be called via multiprocessing pool for parallel deletion. + + Args: + delete_runner_config: The configuration to use for deleting the runner. + """ + delete_runner_config.github_client.delete_runner( + path=delete_runner_config.path, runner_id=int(delete_runner_config.runner_id) + ) def get_runner_context( self, metadata: RunnerMetadata, instance_id: InstanceID, labels: list[str] diff --git a/github-runner-manager/src/github_runner_manager/platform/jobmanager_provider.py b/github-runner-manager/src/github_runner_manager/platform/jobmanager_provider.py index ccace7aca4..91c4f4c5fa 100644 --- a/github-runner-manager/src/github_runner_manager/platform/jobmanager_provider.py +++ b/github-runner-manager/src/github_runner_manager/platform/jobmanager_provider.py @@ -28,10 +28,7 @@ PlatformRunnerHealth, RunnersHealthResponse, ) -from github_runner_manager.types_.github import ( - GitHubRunnerStatus, - SelfHostedRunner, -) +from github_runner_manager.types_.github import GitHubRunnerStatus, SelfHostedRunner logger = logging.getLogger(__name__) @@ -136,15 +133,19 @@ def get_runners_health(self, requested_runners: list[RunnerIdentity]) -> Runners failed_requested_runners=failed_runners, ) - def delete_runner(self, runner_identity: RunnerIdentity) -> None: + def delete_runners(self, runner_ids: list[str]) -> list[str]: """Delete a runner from jobmanager. This method does nothing, as the jobmanager does not implement it. Args: - runner_identity: The identity of the runner to delete. + runner_ids: The runner IDs to delete. + + Returns: + The runner IDs requested for deletion. """ logger.debug("No need to delete runners in the jobmanager.") + return runner_ids def get_runner_context( self, metadata: RunnerMetadata, instance_id: InstanceID, labels: list[str] diff --git a/github-runner-manager/src/github_runner_manager/platform/platform_provider.py b/github-runner-manager/src/github_runner_manager/platform/platform_provider.py index 3f41b78535..8b5c5fa4bd 100644 --- a/github-runner-manager/src/github_runner_manager/platform/platform_provider.py +++ b/github-runner-manager/src/github_runner_manager/platform/platform_provider.py @@ -86,13 +86,11 @@ def get_runners_health( """ @abc.abstractmethod - def delete_runner(self, runner_identity: RunnerIdentity) -> None: - """Delete a runner. - - Can raise DeleteRunnerBusyError + def delete_runners(self, runner_ids: list[str]) -> list[str]: + """Delete runners. Args: - runner_identity: Runner to delete. + runner_ids: Runner IDs to delete. """ @abc.abstractmethod diff --git a/github-runner-manager/tests/unit/factories/metrics_factory.py b/github-runner-manager/tests/unit/factories/metrics_factory.py new file mode 100644 index 0000000000..43b649a6c4 --- /dev/null +++ b/github-runner-manager/tests/unit/factories/metrics_factory.py @@ -0,0 +1,78 @@ +# Copyright 2025 Canonical Ltd. +# See LICENSE file for licensing details. + +"""Factories for Metrics objects.""" + +import factory + +from github_runner_manager.manager.cloud_runner_manager import CodeInformation +from github_runner_manager.metrics.events import Event, RunnerInstalled, RunnerStop + + +class EventFactory(factory.Factory): + """Factory for creating Event instances.""" + + class Meta: + """Meta class for Event. + + Attributes: + model: The metadata reference model. + """ + + model = Event + + timestamp = factory.Faker("unix_time", end_datetime="now") + event = factory.LazyAttribute(lambda obj: obj.__class__.__name__.lower()) + + +class RunnerInstalledFactory(EventFactory): + """Factory for creating RunnerInstalled instances.""" + + class Meta: + """Meta class for RunnerInstalled. + + Attributes: + model: The metadata reference model. + """ + + model = RunnerInstalled + + flavor = factory.Faker("word", ext_word_list=["large", "xlarge"]) + duration = factory.Faker("random_int", min=1, max=3600) + + +class CodeInformationFactory(factory.Factory): + """Factory for creating CodeInformation instances.""" + + class Meta: + """Meta class for CodeInformation. + + Attributes: + model: The metadata reference model. + """ + + model = CodeInformation + + code = factory.Faker("random_int", min=100, max=599) + + +class RunnerStopFactory(EventFactory): + """Factory for creating RunnerStop instances.""" + + class Meta: + """Meta class for RunnerStop. + + Attributes: + model: The metadata reference model. + """ + + model = RunnerStop + + flavor = factory.Faker("word", ext_word_list=["large", "xlarge"]) + workflow = factory.Faker("sentence", nb_words=3) + repo = factory.Faker("word") + github_event = factory.Faker("word") + status = factory.Faker("sentence", nb_words=5) + status_info = factory.SubFactory(CodeInformationFactory) + job_duration = factory.Faker("random_int", min=1, max=3600) + job_conclusion = factory.Faker("word", ext_word_list=["success", "failure", "cancelled"]) diff --git a/github-runner-manager/tests/unit/factories/runner_instance_factory.py b/github-runner-manager/tests/unit/factories/runner_instance_factory.py new file mode 100644 index 0000000000..cf594127e2 --- /dev/null +++ b/github-runner-manager/tests/unit/factories/runner_instance_factory.py @@ -0,0 +1,199 @@ +# Copyright 2025 Canonical Ltd. +# See LICENSE file for licensing details. + +"""Factories for Runner instance objects.""" + +import secrets +from datetime import datetime, timezone + +import factory + +from github_runner_manager.manager.cloud_runner_manager import ( + CloudRunnerInstance, + CloudRunnerState, +) +from github_runner_manager.manager.cloud_runner_manager import HealthState as CloudHelathState +from github_runner_manager.manager.models import InstanceID, RunnerIdentity, RunnerMetadata +from github_runner_manager.manager.runner_manager import RunnerInstance +from github_runner_manager.platform.platform_provider import ( + PlatformRunnerHealth, + PlatformRunnerState, +) +from github_runner_manager.types_.github import SelfHostedRunner + + +class InstanceIDFactory(factory.Factory): + """Factory class for creating InstanceID.""" + + class Meta: + """Meta class for InstanceID. + + Attributes: + model: The metadata reference model. + """ + + model = InstanceID + + prefix = factory.Faker("word") + reactive = factory.Iterator([True, False, None]) + suffix = factory.LazyAttribute(lambda _: secrets.token_hex(6)) + + +class RunnerMetadataFactory(factory.Factory): + """Factory for creating RunnerMetadata instances.""" + + class Meta: + """Meta class for RunnerMetadata. + + Attributes: + model: The metadata reference model. + """ + + model = RunnerMetadata + + platform_name = factory.Faker("word", ext_word_list=["github", "jobmanager"]) + runner_id = str(factory.Faker("random_int", min=1, max=10000)) + url = factory.Faker("url") + + +class CloudRunnerInstanceFactory(factory.Factory): + """Factory for creating CloudRunnerInstance instances.""" + + class Meta: + """Meta class for CloudRunnerInstance. + + Attributes: + model: The metadata reference model. + """ + + model = CloudRunnerInstance + + name = factory.Faker("word") + instance_id = factory.SubFactory(InstanceIDFactory) + metadata = factory.SubFactory(RunnerMetadataFactory) + health = CloudHelathState.HEALTHY + state = CloudRunnerState.ACTIVE + created_at = factory.LazyFunction(lambda: datetime.now(tz=timezone.utc)) + + @classmethod + def from_self_hosted_runner(cls, self_hosted_runner: SelfHostedRunner) -> CloudRunnerInstance: + """Construct CloudRunnerInstance associated to self hosted runner. + + Args: + self_hosted_runner: The target self hosted runner to associate. + + Returns: + The Instantiated CloudRunnerInstance. + """ + return CloudRunnerInstanceFactory( + instance_id=self_hosted_runner.identity.instance_id, + metadata=RunnerMetadataFactory(runner_id=str(self_hosted_runner.id)), + ) + + +class RunnerIdentityFactory(factory.Factory): + """Factory for creating RunnerIdentity instances.""" + + class Meta: + """Meta class for RunnerIdentity. + + Attributes: + model: The metadata reference model. + """ + + model = RunnerIdentity + + instance_id = factory.SubFactory(InstanceIDFactory) + metadata = factory.SubFactory(RunnerMetadataFactory) + + +class PlatformRunnerHealthFactory(factory.Factory): + """Factory for creating PlatformRunnerHealth instances.""" + + class Meta: + """Meta class for PlatformRunnerHealth. + + Attributes: + model: The metadata reference model. + """ + + model = PlatformRunnerHealth + + identity = factory.SubFactory(RunnerIdentityFactory) + online = factory.Faker("boolean") + busy = factory.Faker("boolean") + deletable = factory.Faker("boolean") + runner_in_platform = factory.Faker("boolean") + + +class RunnerInstanceFactory(factory.Factory): + """Factory for creating RunnerInstance instances.""" + + class Meta: + """Meta class for RunnerInstance. + + Attributes: + model: The metadata reference model. + """ + + model = RunnerInstance + + name = factory.Faker("name") + instance_id = factory.SubFactory(InstanceIDFactory) + metadata = factory.LazyAttribute(lambda _: {"key": "value"}) + health = CloudHelathState.HEALTHY + platform_state = platform_state = factory.LazyFunction( + lambda: secrets.choice(list(PlatformRunnerState)) + ) + cloud_state = CloudRunnerState.ACTIVE + + @classmethod + def from_state( + cls, cloud_runner: CloudRunnerInstance, platform_health: PlatformRunnerHealth | None = None + ) -> RunnerInstance: + """Generate RunnerInstance from cloud runner and platform runner states. + + Args: + cloud_runner: The cloud runner to generate state from. + platform_health: The platform runner to generate state from. + + Returns: + The generated RunnerInstance. + """ + return RunnerInstance( + name=cloud_runner.name, + instance_id=cloud_runner.instance_id, + metadata=cloud_runner.metadata, + health=cloud_runner.health, + platform_state=( + PlatformRunnerState.from_platform_health(platform_health) + if platform_health is not None + else None + ), + cloud_state=cloud_runner.state, + ) + + +class SelfHostedRunnerFactory(factory.Factory): + """Factory for creating SelfHostedRunner instances.""" + + class Meta: + """Meta class for SelfHostedRunner. + + Attributes: + model: The metadata reference model. + """ + + model = SelfHostedRunner + + busy = factory.Faker("boolean") + id = factory.Faker("random_int", min=1, max=10000) + labels = factory.List([factory.Faker("word") for _ in range(3)]) + status = factory.Faker("word", ext_word_list=["online", "offline"]) + deletable = factory.Faker("boolean") + # identity.metadata.runner_id should be equal to the id attribute. + identity = factory.LazyAttribute( + lambda obj: RunnerIdentityFactory( + metadata=RunnerMetadata(platform_name="github", runner_id=obj.id), + ) + ) diff --git a/github-runner-manager/tests/unit/fake_runner_managers.py b/github-runner-manager/tests/unit/fake_runner_managers.py new file mode 100644 index 0000000000..d879946882 --- /dev/null +++ b/github-runner-manager/tests/unit/fake_runner_managers.py @@ -0,0 +1,307 @@ +# Copyright 2025 Canonical Ltd. +# See LICENSE file for licensing details. +import logging +from typing import Sequence +from unittest.mock import MagicMock + +from pydantic import HttpUrl + +from github_runner_manager.manager.cloud_runner_manager import ( + CloudRunnerInstance, + CloudRunnerManager, +) +from github_runner_manager.manager.models import ( + InstanceID, + RunnerContext, + RunnerIdentity, + RunnerMetadata, +) +from github_runner_manager.metrics.runner import RunnerMetrics +from github_runner_manager.openstack_cloud.openstack_cloud import _MAX_NOVA_COMPUTE_API_VERSION +from github_runner_manager.platform.platform_provider import ( + JobInfo, + PlatformProvider, + PlatformRunnerHealth, + RunnersHealthResponse, +) +from github_runner_manager.types_.github import GitHubRunnerStatus, SelfHostedRunner +from tests.unit.factories.runner_instance_factory import CloudRunnerInstanceFactory + +logger = logging.getLogger(__name__) + + +class FakeOpenstackCloud: + """Fake implementation of OpenstackCloud.""" + + _MOCK_COMPUTE_ENDPOINT = "mock-compute-endpoint" + _MOCK_COMPUTE_ENDPOINT_RESPONSE = {"version": {"version": _MAX_NOVA_COMPUTE_API_VERSION}} + + def __init__( + self, + initial_servers: list[InstanceID], + server_to_errors: dict[InstanceID, Exception] | None = None, + ) -> None: + """Initialize the OpenstackCloud mock object.""" + self.servers = {instance.name: instance for instance in initial_servers} + self._injected_errors = { + instance.name: exc for instance, exc in (server_to_errors or {}).items() + } + + def __enter__(self) -> "FakeOpenstackCloud": + """Fake enter method for context entering.""" + return self + + def __exit__(self, *args, **kwargs) -> None: + """Fake exit method for context exiting.""" + return + + def connect(self) -> "FakeOpenstackCloud": + """Fake OpenStack lib's connect function.""" + return self + + @property + def compute(self) -> "FakeOpenstackCloud": + """Fake the compute API attribute.""" + return self + + def get_endpoint(self) -> str: + """Fake endpoint string for compute endpoint.""" + return self._MOCK_COMPUTE_ENDPOINT + + @property + def session(self) -> dict: + """Fake the connection session attribute.""" + compute_endpoint_mock = MagicMock() + compute_endpoint_mock.json.return_value = self._MOCK_COMPUTE_ENDPOINT_RESPONSE + return {self._MOCK_COMPUTE_ENDPOINT: compute_endpoint_mock} + + def delete_server( + self, + name_or_id: str, + wait: bool = False, + timeout: int = 180, + delete_ips: bool = False, + delete_ip_retry: int = 1, + ) -> bool: + """Fake method for deleting server.""" + injected_test_error = self._injected_errors.pop(name_or_id, None) + if injected_test_error: + raise injected_test_error + + if self.servers.pop(name_or_id, None): + return True + return False + + def delete_keypair(self, *args, **kwargs): + """Fake delete keypair method.""" + pass + + +class FakeCloudRunnerManager(CloudRunnerManager): + """Fake of CloudRunnerManager. + + Metrics is not supported in this fake. + + Attributes: + name_prefix: The naming prefix for runners managed. + """ + + @property + def name_prefix(self) -> str: + """The naming prefix for runners managed.""" + return "fake_cloud_runner_manager" + + def __init__(self, initial_cloud_runners: list[CloudRunnerInstance]) -> None: + """Initialize the Cloud Runner Manager. + + Args: + initial_cloud_runners: A list of initial Cloud Runner Instances. + """ + self._cloud_runners = {runner.instance_id: runner for runner in initial_cloud_runners} + + def create_runner( + self, runner_identity: RunnerIdentity, runner_context: RunnerContext + ) -> CloudRunnerInstance: + """Create a runner instance for the given runner identity and context. + + Args: + runner_identity: The runner identity to create a runner for. + runner_context: The context for the runner to create a runner for. + + Returns: + The created runner instance. + """ + created_runner = CloudRunnerInstanceFactory(instance_id=runner_identity.instance_id) + self._cloud_runners[runner_identity.instance_id] = created_runner + return created_runner + + def get_runners(self) -> Sequence[CloudRunnerInstance]: + """Get all the cloud runner instances managed by the manager. + + Returns: + A list of cloud runner instances. + """ + return list(self._cloud_runners.values()) + + def delete_vms(self, instance_ids: Sequence[InstanceID]) -> list[InstanceID]: + """Delete VMs with given instance ids. + + Args: + instance_ids: A list of instance ids to delete. + + Returns: + A list of instance ids that were deleted. + """ + deleted_instance_ids: list[InstanceID] = [] + for instance_id in instance_ids: + cloud_runner = self._cloud_runners.pop(instance_id, None) + if not cloud_runner: + continue + deleted_instance_ids.append(cloud_runner.instance_id) + return deleted_instance_ids + + def extract_metrics(self, instance_ids: Sequence[InstanceID]) -> list[RunnerMetrics]: + """Extract metrics from VMs with given instance ids. + + The fake runner manager does not implement this. + + Args: + instance_ids: A list of instance ids to extract metrics from. + + Returns: + A list of metrics extracted from VMs with given instance ids. + """ + return [] + + def cleanup(self) -> None: + """Cleanup cloud resources. + + The fake runner manager does not implement this. + """ + pass + + +class FakeGitHubRunnerPlatform(PlatformProvider): + """Fake GitHub platform provider.""" + + def __init__(self, initial_runners: Sequence[SelfHostedRunner]) -> None: + """Initialize the fake platform. + + Args: + initial_runners: Runners to instantiate the platform with. + """ + self._runners = {runner.identity.instance_id: runner for runner in initial_runners} + + def get_runner_health(self, runner_identity: RunnerIdentity) -> PlatformRunnerHealth: + """Get runner health of a runner with given runner identity. + + Args: + runner_identity: The identity of the runner to query health status. + + Returns: + The PlatformRunnerHealth status of the runner. + """ + runner = self._runners.get(runner_identity.instance_id, None) + if not runner: + return PlatformRunnerHealth( + identity=runner_identity, + online=False, + busy=False, + deletable=True, + runner_in_platform=False, + ) + return PlatformRunnerHealth( + identity=runner_identity, + online=runner.status == GitHubRunnerStatus.ONLINE, + busy=runner.busy, + deletable=False, + ) + + def get_runners_health(self, requested_runners: list[RunnerIdentity]) -> RunnersHealthResponse: + """Batch get runners health. + + Args: + requested_runners: The runners to get. the health information for. + + Returns: + The requested runners health info. + """ + response = RunnersHealthResponse() + + for requested_runner in requested_runners: + runner = self._runners.get(requested_runner.instance_id, None) + if runner: + response.requested_runners.append( + self.get_runner_health(runner_identity=runner.identity) + ) + continue + response.failed_requested_runners.append(requested_runner) + + requested_runner_ids = set(runner.instance_id for runner in requested_runners) + for instance_id, runner in self._runners.items(): + if instance_id in requested_runner_ids: + continue + response.non_requested_runners.append(runner.identity) + return response + + def delete_runners(self, runner_ids: list[str]) -> list[str]: + """Delete runners from platform. + + Args: + runner_ids: The runner IDs to delete. + + Returns: + The successfully deleted runners. + """ + deleted_runner_ids: list[str] = [] + runner_id_map = {str(runner.id): runner for runner in self._runners.values()} + for runner_id in runner_ids: + runner = runner_id_map.get(runner_id, None) + if not runner: + continue + self._runners.pop(runner.identity.instance_id) + deleted_runner_ids.append(runner_id) + return deleted_runner_ids + + def get_runner_context( + self, metadata: RunnerMetadata, instance_id: InstanceID, labels: list[str] + ) -> tuple[RunnerContext, SelfHostedRunner]: + """Get a context for a runner. + + Args: + metadata: The runner's metadata. + instance_id: The ID of the instance. + labels: The labels of the instance. + + Raises: + NotImplementedError: This method is not tested with this mock. + """ + raise NotImplementedError + + def check_job_been_picked_up(self, metadata: RunnerMetadata, job_url: HttpUrl) -> bool: + """Check if a job has been picked up by the runner. + + Args: + metadata: The metadata of the runner. + job_url: The URL of the job. + + Raises: + NotImplementedError: This method is not tested with this mock. + """ + raise NotImplementedError + + def get_job_info( + self, metadata: RunnerMetadata, repository: str, workflow_run_id: str, runner: InstanceID + ) -> JobInfo: + """Get information about a job. + + Args: + metadata: The metadata of the runner. + repository: The name of the repository. + workflow_run_id: The ID of the workflow run. + runner: The ID of the runner. + + Raises: + NotImplementedError: This method is not tested with this mock. + """ + raise NotImplementedError diff --git a/github-runner-manager/tests/unit/manager/test_runner_manager.py b/github-runner-manager/tests/unit/manager/test_runner_manager.py index 7cc6cd1dae..090522e2a3 100644 --- a/github-runner-manager/tests/unit/manager/test_runner_manager.py +++ b/github-runner-manager/tests/unit/manager/test_runner_manager.py @@ -7,102 +7,192 @@ import pytest -from github_runner_manager.errors import RunnerCreateError from github_runner_manager.manager.cloud_runner_manager import ( CloudRunnerInstance, CloudRunnerManager, ) -from github_runner_manager.manager.models import ( - InstanceID, - RunnerContext, - RunnerIdentity, - RunnerMetadata, +from github_runner_manager.manager.models import RunnerMetadata +from github_runner_manager.manager.runner_manager import FlushMode, RunnerInstance, RunnerManager +from github_runner_manager.platform.platform_provider import PlatformProvider +from github_runner_manager.types_.github import SelfHostedRunner +from tests.unit.factories.runner_instance_factory import ( + CloudRunnerInstanceFactory, + RunnerInstanceFactory, + SelfHostedRunnerFactory, ) -from github_runner_manager.manager.runner_manager import RunnerManager -from github_runner_manager.platform.platform_provider import ( - PlatformProvider, - RunnersHealthResponse, -) -from github_runner_manager.types_.github import GitHubRunnerStatus, SelfHostedRunner +from tests.unit.fake_runner_managers import FakeCloudRunnerManager, FakeGitHubRunnerPlatform -def test_cleanup_removes_runners_in_platform_not_in_cloud(monkeypatch: pytest.MonkeyPatch): +@pytest.mark.parametrize( + "initial_runners, initial_cloud_runners, expected_runners, expected_cloud_runners, flush_mode", + [ + pytest.param( + [SelfHostedRunnerFactory()], + [], + [], + [], + FlushMode.FLUSH_IDLE, + id="one platform runner not in cloud is cleaned up", + ), + pytest.param( + [ + ( + idle_runner_with_cloud := SelfHostedRunnerFactory( + busy=False, + status="online", + ) + ), + ], + [ + runner_with_platform := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=idle_runner_with_cloud + ) + ], + [], + [], + FlushMode.FLUSH_IDLE, + id="one idle platform runner, matching cloud runner in cloud flushed", + ), + pytest.param( + [ + ( + busy_runner_with_cloud := SelfHostedRunnerFactory( + busy=True, + status="online", + deletable=False, + ) + ), + ], + [ + runner_with_platform := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=busy_runner_with_cloud + ) + ], + [busy_runner_with_cloud], + [runner_with_platform], + FlushMode.FLUSH_IDLE, + id="one busy platform runner, matching cloud runner in cloud is not flushed", + ), + pytest.param( + [ + ( + busy_runner_with_cloud := SelfHostedRunnerFactory( + busy=True, + status="online", + deletable=False, + ) + ), + ], + [ + runner_with_platform := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=busy_runner_with_cloud + ) + ], + [], + [], + FlushMode.FLUSH_BUSY, + id="one busy platform runner, matching cloud runner in cloud is flushed in flush busy", + ), + ], +) +def test_flush_runners( + initial_runners: list[SelfHostedRunner], + initial_cloud_runners: list[CloudRunnerInstance], + expected_runners: list[SelfHostedRunner], + expected_cloud_runners: list[CloudRunnerInstance], + flush_mode: FlushMode, +): """ - arrange: Given a runner in GitHub that is not in the cloud provider. - act: Call cleanup in the RunnerManager instance. - assert: The github runner should be deleted. + arrange: Given GitHub runners and Cloud runners. + act: Call flush in the RunnerManager instance. + assert: Expected github runners and cloud runners are flushed. """ - instance_id = InstanceID.build("prefix-0") - github_runner_identity = RunnerIdentity( - instance_id=instance_id, metadata=RunnerMetadata(platform_name="github", runner_id="1") - ) - - cloud_runner_manager = MagicMock() - cloud_runner_manager.get_runners.return_value = [] - github_provider = MagicMock() - runner_manager = RunnerManager( - "managername", - platform_provider=github_provider, - cloud_runner_manager=cloud_runner_manager, - labels=["label1", "label2"], + mock_platform = FakeGitHubRunnerPlatform(initial_runners=initial_runners) + mock_cloud = FakeCloudRunnerManager(initial_cloud_runners=initial_cloud_runners) + manager = RunnerManager( + "test-manager", platform_provider=mock_platform, cloud_runner_manager=mock_cloud, labels=[] ) - github_provider.get_runners_health.return_value = RunnersHealthResponse( - non_requested_runners=[github_runner_identity] - ) - - runner_manager.cleanup() + manager.flush_runners(flush_mode=flush_mode) - github_provider.delete_runner.assert_called_with(github_runner_identity) - cloud_runner_manager.delete_runner.assert_not_called() + assert list(mock_platform._runners.values()) == expected_runners + assert list(mock_cloud._cloud_runners.values()) == expected_cloud_runners -def test_failed_runner_in_openstack_cleans_github(monkeypatch: pytest.MonkeyPatch): +@pytest.mark.parametrize( + "initial_runners, initial_cloud_runners, expected_runners, expected_cloud_runners", + [ + pytest.param( + [SelfHostedRunnerFactory()], [], [], [], id="one platform runner not in cloud" + ), + pytest.param( + [ + (runner_with_cloud := SelfHostedRunnerFactory()), + SelfHostedRunnerFactory(), + ], + [ + runner_with_platform := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=runner_with_cloud + ) + ], + [runner_with_cloud], + [runner_with_platform], + id="one platform runner not in cloud, one in cloud", + ), + pytest.param( + [], + [runner_without_platform := CloudRunnerInstanceFactory()], + [], + [runner_without_platform], + id="cloud runner only in cloud", + ), + pytest.param( + [], + [runner_with_platform], + [], + [runner_with_platform], + id="cloud runner with platform runner only in cloud", + ), + pytest.param( + [SelfHostedRunnerFactory(), SelfHostedRunnerFactory(), SelfHostedRunnerFactory()], + [], + [], + [], + id="multiple runners not in cloud, none in cloud", + ), + pytest.param( + [runner_with_cloud, SelfHostedRunnerFactory()], + [runner_with_platform, runner_without_platform := CloudRunnerInstanceFactory()], + [runner_with_cloud], + [runner_with_platform, runner_without_platform], + id="some in cloud, some not in cloud", + ), + ], +) +def test_runner_maanger_cleanup( + initial_runners: list[SelfHostedRunner], + initial_cloud_runners: list[CloudRunnerInstance], + expected_runners: list[SelfHostedRunner], + expected_cloud_runners: list[CloudRunnerInstance], +): """ - arrange: Prepare a RunnerManager with a cloud manager that will fail when creating a runner. - act: Create a Runner. - assert: When there was a failure to create a runner in the cloud manager, - only that github runner will be deleted in GitHub. + arrange: Given GitHub runners and Cloud runners. + act: Call cleanup in the RunnerManager instance. + assert: Expected github runners and cloud runners cleanup is run. """ - cloud_instances: tuple[CloudRunnerInstance, ...] = () - cloud_runner_manager = MagicMock() - cloud_runner_manager.get_runners.return_value = cloud_instances - cloud_runner_manager.name_prefix = "unit-0" - github_provider = MagicMock(spec=PlatformProvider) - - runner_manager = RunnerManager( - "managername", - platform_provider=github_provider, - cloud_runner_manager=cloud_runner_manager, - labels=["label1", "label2"], - ) - - identity = RunnerIdentity( - instance_id=InstanceID.build("invalid"), - metadata=RunnerMetadata(platform_name="github", runner_id="1"), + mock_platform = FakeGitHubRunnerPlatform(initial_runners=initial_runners) + mock_cloud = FakeCloudRunnerManager(initial_cloud_runners=initial_cloud_runners) + manager = RunnerManager( + "test-manager", platform_provider=mock_platform, cloud_runner_manager=mock_cloud, labels=[] ) - github_runner = SelfHostedRunner( - identity=identity, - id=1, - labels=[], - status=GitHubRunnerStatus.OFFLINE, - busy=True, - ) - - def _get_runner_context(instance_id, metadata, labels): - """Return the runner context.""" - nonlocal github_runner - github_runner.identity.instance_id = instance_id - return RunnerContext(shell_run_script="agent"), github_runner - github_provider.get_runner_context.side_effect = _get_runner_context - cloud_runner_manager.create_runner.side_effect = RunnerCreateError("") + manager.cleanup() - _ = runner_manager.create_runners(1, RunnerMetadata(), True) - github_provider.delete_runner.assert_called_once_with(github_runner.identity) + assert list(mock_platform._runners.values()) == expected_runners + assert list(mock_cloud._cloud_runners.values()) == expected_cloud_runners -def test_create_runner() -> None: +def test_runner_manager_create_runners() -> None: """ arrange: None. act: call runner_manager.create_runners. @@ -127,3 +217,147 @@ def test_create_runner() -> None: assert instance_id cloud_runner_manager.create_runner.assert_called_once() + + +@pytest.mark.parametrize( + "initial_runners, initial_cloud_runners, expected_runner_instances", + [ + pytest.param([], [], (), id="no runners"), + pytest.param( + [SelfHostedRunnerFactory()], [], (), id="platform runner without cloud runner" + ), + pytest.param( + [], + [cloud_runner := CloudRunnerInstanceFactory()], + (RunnerInstanceFactory(name=cloud_runner.name),), + id="cloud runner without platform runner", + ), + pytest.param( + [runner_with_cloud := SelfHostedRunnerFactory()], + [ + cloud_runner := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=runner_with_cloud + ) + ], + (RunnerInstanceFactory(name=cloud_runner.name),), + id="platform runner with cloud runner", + ), + ], +) +def test_runner_manager_get_runners( + initial_runners: list[SelfHostedRunner], + initial_cloud_runners: list[CloudRunnerInstance], + expected_runner_instances: tuple[RunnerInstance], +): + """ + arrange: Given GitHub runners and Cloud runners. + act: when RunnerManager.get_runners is called. + assert: expected RunnerInstances are returned. + """ + mock_platform = FakeGitHubRunnerPlatform(initial_runners=initial_runners) + mock_cloud = FakeCloudRunnerManager(initial_cloud_runners=initial_cloud_runners) + manager = RunnerManager( + "test-manager", platform_provider=mock_platform, cloud_runner_manager=mock_cloud, labels=[] + ) + + # Test for number of runners matching and that the runner belongs to the cloud instance + # provided. The instance contents cannot be tested due to coupled logic in how the + # RunnerInstance is generated - which would require business logic in the test of setting up + # RunnerInstances from cloud runners and self hosted runners. + result = manager.get_runners() + assert len(result) == len(expected_runner_instances) + runner_names = {runner.name for runner in result} + assert all(runner.name in runner_names for runner in expected_runner_instances) + + +@pytest.mark.parametrize( + "initial_runners, initial_cloud_runners, num_delete, expected_runners, expected_cloud_runners", + [ + pytest.param([], [], 1, [], [], id="no runners to delete"), + pytest.param( + [runner_with_cloud := SelfHostedRunnerFactory()], + [ + runner_with_platform := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=runner_with_cloud + ) + ], + 0, + [runner_with_cloud], + [runner_with_platform], + id="num delete runners 0", + ), + pytest.param( + [runner_with_cloud], + [runner_with_platform], + 1, + [], + [], + id="delete 1 runner", + ), + ], +) +def test_runner_manager_deterministic_delete_runners( + initial_runners: list[SelfHostedRunner], + initial_cloud_runners: list[CloudRunnerInstance], + num_delete: int, + expected_runners: list[SelfHostedRunner], + expected_cloud_runners: list[CloudRunnerInstance], +): + """ + arrange: given initial runners (platform and cloud). + act: when RunnerManager.delete_runners is called. + assert: expected cloud & platform runners remain. + """ + mock_platform = FakeGitHubRunnerPlatform(initial_runners=initial_runners) + mock_cloud = FakeCloudRunnerManager(initial_cloud_runners=initial_cloud_runners) + manager = RunnerManager( + "test-manager", platform_provider=mock_platform, cloud_runner_manager=mock_cloud, labels=[] + ) + + manager.delete_runners(num_delete) + + assert list(mock_platform._runners.values()) == expected_runners + assert list(mock_cloud._cloud_runners.values()) == expected_cloud_runners + + +@pytest.mark.parametrize( + "initial_runners, initial_cloud_runners, num_delete," + "expected_runners_count, expected_cloud_runners_count", + [ + pytest.param( + [runner_with_cloud, runner_with_cloud_two := SelfHostedRunnerFactory()], + [ + runner_with_platform, + runner_with_platform_two := CloudRunnerInstanceFactory.from_self_hosted_runner( + self_hosted_runner=runner_with_cloud_two + ), + ], + 1, + 1, + 1, + id="delete 1 runner out of two runners", + ), + ], +) +def test_runner_manager_non_deterministic_delete_runners( + initial_runners: list[SelfHostedRunner], + initial_cloud_runners: list[CloudRunnerInstance], + num_delete: int, + expected_runners_count: int, + expected_cloud_runners_count: int, +): + """ + arrange: given initial runners (platform and cloud). + act: when RunnerManager.delete_runners is called. + assert: expected cloud & platform runners remain. + """ + mock_platform = FakeGitHubRunnerPlatform(initial_runners=initial_runners) + mock_cloud = FakeCloudRunnerManager(initial_cloud_runners=initial_cloud_runners) + manager = RunnerManager( + "test-manager", platform_provider=mock_platform, cloud_runner_manager=mock_cloud, labels=[] + ) + + manager.delete_runners(num_delete) + + assert len(mock_platform._runners.values()) == expected_runners_count + assert len(mock_platform._runners.values()) == expected_cloud_runners_count diff --git a/github-runner-manager/tests/unit/metrics/test_runner.py b/github-runner-manager/tests/unit/metrics/test_runner.py index dfcadd80ee..f26ec3974a 100644 --- a/github-runner-manager/tests/unit/metrics/test_runner.py +++ b/github-runner-manager/tests/unit/metrics/test_runner.py @@ -21,7 +21,13 @@ from github_runner_manager.metrics import runner as runner_metrics from github_runner_manager.metrics import type as metrics_type from github_runner_manager.metrics.events import RunnerInstalled, RunnerStart, RunnerStop -from github_runner_manager.metrics.runner import PullFileError, ssh_pull_file +from github_runner_manager.metrics.runner import ( + PulledMetrics, + PullFileError, + SSHError, + _ssh_pull_file, + pull_runner_metrics, +) from github_runner_manager.types_.github import JobConclusion @@ -41,28 +47,106 @@ def runner_fs_base_fixture(tmp_path: Path) -> Path: return runner_fs_base -def _create_metrics_data(instance_id: InstanceID) -> RunnerMetrics: - """Create a RunnerMetrics object that is suitable for most tests. +def test_pull_runner_metrics_errors(caplog: pytest.LogCaptureFixture): + """ + arrange: given a mocked cloud service that raises exceptions are different points. + act: when pull_runner_metrics function is called. + assert: no metrics are pulled and errors are logged. + """ + test_instances = [] + get_instance_side_effects: list[None | Exception] = [] + get_ssh_connection_side_effects = [] + # Setup for instance not exists + test_instances.append( + ( + not_exists_instance := InstanceID( + prefix="instance-not-exists", reactive=False, suffix="1" + ) + ) + ) + get_instance_side_effects.append(None) + # Setup for instance fail ssh connection + test_instances.append( + (fail_ssh_conn_instance := InstanceID(prefix="fail-ssh-conn", reactive=False, suffix="2")) + ) + fail_ssh_conn_instance_mock = MagicMock() + fail_ssh_conn_instance_mock.instance_id = fail_ssh_conn_instance + get_instance_side_effects.append(fail_ssh_conn_instance_mock) + get_ssh_connection_side_effects.append(SSHError()) + # Setup for instance fail pull file + test_instances.append( + ( + fail_pull_file_instance := InstanceID( + prefix="fail-pull-file", reactive=False, suffix="3" + ) + ) + ) + fail_pull_file_instance_mock = MagicMock() + fail_pull_file_instance_mock.instance_id = fail_pull_file_instance + get_instance_side_effects.append(fail_pull_file_instance_mock) + ssh_connection_mock = MagicMock() + ssh_connection_mock.return_value = ssh_connection_mock + ssh_connection_mock.__enter__ = ssh_connection_mock + ssh_connection_mock.run.side_effect = [TimeoutError] + get_ssh_connection_side_effects.append(ssh_connection_mock) + # Mock cloud service setup + mock_cloud_service = MagicMock() + mock_cloud_service.get_instance = MagicMock(side_effect=get_instance_side_effects) + mock_cloud_service.get_ssh_connection = MagicMock(side_effect=get_ssh_connection_side_effects) + + assert pull_runner_metrics(cloud_service=mock_cloud_service, instance_ids=test_instances) == [] + assert ( + f"Skipping fetching metrics, instance not found: {not_exists_instance}" in caplog.messages + ) + assert ( + f"Failed to create SSH connection for pulling metrics: {fail_ssh_conn_instance}" + in caplog.messages + ) + assert ( + f"Failed to create SSH connection for pulling metrics: {fail_pull_file_instance}" + in caplog.messages + ) - Args: - instance_id: The test runner name. - Returns: - Test metrics data. +def test_pull_runner_metrics(): """ - return RunnerMetrics( - installation_start_timestamp=1, - installed_timestamp=2, - pre_job=PreJobMetrics( - timestamp=3, - workflow="workflow1", - workflow_run_id="workflow_run_id1", - repository="org1/repository1", - event="push", - ), - post_job=PostJobMetrics(timestamp=3, status=PostJobStatus.NORMAL), - instance_id=instance_id, - metadata=RunnerMetadata(), + arrange: given a mock cloud service get_instance method and get_ssh_connection method. + act: when pull_runner_metrics function is called. + assert: metrics are pulled from corresponding instances correctly. + """ + mock_cloud_service = MagicMock() + mock_ssh_conn = MagicMock() + mock_ssh_conn.return_value = mock_ssh_conn + mock_ssh_conn.__enter__ = mock_ssh_conn + test_remote_file_contents = "test-contents" + mock_ssh_conn.get = lambda remote, local: local.write( + bytes(test_remote_file_contents, encoding="utf-8") + ) + mock_cloud_service.get_ssh_connection = mock_ssh_conn + mock_instance_one, mock_instance_two = (MagicMock(), MagicMock()) + mock_cloud_service.get_instance.side_effect = [mock_instance_one, mock_instance_two] + + # Compare the set as the order is not guaranteed but it does not matter. + assert set( + pull_runner_metrics( + cloud_service=mock_cloud_service, + instance_ids=[mock_instance_one.instance_id, mock_instance_two.instance_id], + ) + ) == set( + [ + PulledMetrics( + instance=mock_instance_one, + runner_installed=test_remote_file_contents, + pre_job_metrics=test_remote_file_contents, + post_job_metrics=test_remote_file_contents, + ), + PulledMetrics( + instance=mock_instance_two, + runner_installed=test_remote_file_contents, + pre_job_metrics=test_remote_file_contents, + post_job_metrics=test_remote_file_contents, + ), + ] ) @@ -126,6 +210,31 @@ def test_issue_events(issue_event_mock: MagicMock): ) +def _create_metrics_data(instance_id: InstanceID) -> RunnerMetrics: + """Create a RunnerMetrics object that is suitable for most tests. + + Args: + instance_id: The test runner name. + + Returns: + Test metrics data. + """ + return RunnerMetrics( + installation_start_timestamp=1, + installed_timestamp=2, + pre_job=PreJobMetrics( + timestamp=3, + workflow="workflow1", + workflow_run_id="workflow_run_id1", + repository="org1/repository1", + event="push", + ), + post_job=PostJobMetrics(timestamp=3, status=PostJobStatus.NORMAL), + instance_id=instance_id, + metadata=RunnerMetadata(), + ) + + def test_issue_events_pre_job_before_runner_installed(issue_event_mock: MagicMock): """ arrange: A runner with pre-job timestamp smaller than installed timestamp. @@ -363,7 +472,7 @@ def _ssh_get(remote, local) -> None: ssh_conn.get.side_effect = _ssh_get - response = ssh_pull_file(ssh_conn, remote_path, max_size) + response = _ssh_pull_file(ssh_conn, remote_path, max_size) assert response == "content from the file" @@ -400,7 +509,7 @@ def _ssh_get(remote, local) -> None: ssh_conn.get.side_effect = _ssh_get with pytest.raises(PullFileError) as exc: - _ = ssh_pull_file(ssh_conn, remote_path, max_size) + _ = _ssh_pull_file(ssh_conn, remote_path, max_size) assert "max" in str(exc) @@ -426,5 +535,5 @@ def _ssh_run(command, **kwargs) -> Optional[Result]: ssh_conn.run.side_effect = _ssh_run with pytest.raises(PullFileError) as exc: - _ = ssh_pull_file(ssh_conn, remote_path, max_size) + _ = _ssh_pull_file(ssh_conn, remote_path, max_size) assert "too large" in str(exc) diff --git a/github-runner-manager/tests/unit/mock_runner_managers.py b/github-runner-manager/tests/unit/mock_runner_managers.py deleted file mode 100644 index 051b2b2efe..0000000000 --- a/github-runner-manager/tests/unit/mock_runner_managers.py +++ /dev/null @@ -1,491 +0,0 @@ -# Copyright 2025 Canonical Ltd. -# See LICENSE file for licensing details. -import hashlib -import logging -import random -import secrets -from dataclasses import dataclass -from datetime import datetime, timezone -from typing import Iterable, Sequence -from unittest.mock import MagicMock - -from pydantic import HttpUrl - -from github_runner_manager.configuration.github import GitHubPath -from github_runner_manager.github_client import GithubClient -from github_runner_manager.manager.cloud_runner_manager import ( - CloudRunnerInstance, - CloudRunnerManager, - CloudRunnerState, -) -from github_runner_manager.manager.models import ( - InstanceID, - RunnerContext, - RunnerIdentity, - RunnerMetadata, -) -from github_runner_manager.metrics.runner import RunnerMetrics -from github_runner_manager.platform.github_provider import ( - PlatformRunnerState, -) -from github_runner_manager.platform.platform_provider import ( - JobInfo, - PlatformProvider, - PlatformRunnerHealth, - RunnersHealthResponse, -) -from github_runner_manager.types_.github import ( - GitHubRunnerStatus, - JITConfig, - RunnerApplication, - SelfHostedRunner, -) - -logger = logging.getLogger(__name__) - -# Compressed tar file for testing. -# Python `tarfile` module works on only files. -# Hardcoding a sample tar file is simpler. -TEST_BINARY = ( - b"\x1f\x8b\x08\x00\x00\x00\x00\x00\x00\x03\xed\xd1\xb1\t\xc30\x14\x04P\xd5\x99B\x13\x04\xc9" - b"\xb6\xacyRx\x01[\x86\x8c\x1f\x05\x12HeHaB\xe0\xbd\xe6\x8a\x7f\xc5\xc1o\xcb\xd6\xae\xed\xde" - b"\xc2\x89R7\xcf\xd33s-\xe93_J\xc8\xd3X{\xa9\x96\xa1\xf7r\x1e\x87\x1ab:s\xd4\xdb\xbe\xb5\xdb" - b"\x1ac\xcfe=\xee\x1d\xdf\xffT\xeb\xff\xbf\xfcz\x04\x00\x00\x00\x00\x00\x00\x00\x00\x00_{\x00" - b"\xc4\x07\x85\xe8\x00(\x00\x00" -) - - -class MockGhapiClient: - """Mock for Ghapi client.""" - - def __init__(self, token: str): - """Initialization method for GhapiClient fake. - - Args: - token: The client token value. - """ - self.token = token - self.actions = MockGhapiActions() - - def last_page(self) -> int: - """Last page number stub. - - Returns: - Always zero. - """ - return 0 - - -class MockGhapiActions: - """Mock for actions in Ghapi client.""" - - def __init__(self): - """A placeholder method for test stub/fakes initialization.""" - hash = hashlib.sha256() - hash.update(TEST_BINARY) - self.test_hash = hash.hexdigest() - self.registration_token_repo = secrets.token_hex() - self.registration_token_org = secrets.token_hex() - - def _list_runner_applications(self): - """A placeholder method for test fake. - - Returns: - A fake runner applications list. - """ - runners = [] - runners.append( - RunnerApplication( - os="linux", - architecture="x64", - download_url="https://www.example.com", - filename="test_runner_binary", - sha256_checksum=self.test_hash, - ) - ) - return runners - - def list_runner_applications_for_repo(self, owner: str, repo: str): - """A placeholder method for test stub. - - Args: - owner: Placeholder for repository owner. - repo: Placeholder for repository name. - - Returns: - A fake runner applications list. - """ - return self._list_runner_applications() - - def list_runner_applications_for_org(self, org: str): - """A placeholder method for test stub. - - Args: - org: Placeholder for repository owner. - - Returns: - A fake runner applications list. - """ - return self._list_runner_applications() - - def create_registration_token_for_repo(self, owner: str, repo: str): - """A placeholder method for test stub. - - Args: - owner: Placeholder for repository owner. - repo: Placeholder for repository name. - - Returns: - Registration token stub. - """ - return JITConfig( - {"token": self.registration_token_repo, "expires_at": "2020-01-22T12:13:35.123-08:00"} - ) - - def list_self_hosted_runners_for_repo( - self, owner: str, repo: str, per_page: int, page: int = 0 - ): - """A placeholder method for test stub. - - Args: - owner: Placeholder for repository owner. - repo: Placeholder for repository name. - per_page: Placeholder for responses per page. - page: Placeholder for response page number. - - Returns: - Empty runners stub. - """ - return {"runners": []} - - def list_self_hosted_runners_for_org(self, org: str, per_page: int, page: int = 0): - """A placeholder method for test stub. - - Args: - org: Placeholder for repository owner. - per_page: Placeholder for responses per page. - page: Placeholder for response page number. - - Returns: - Empty runners stub. - """ - return {"runners": []} - - def delete_self_hosted_runner_from_repo(self, owner: str, repo: str, runner_id: str): - """A placeholder method for test stub. - - Args: - owner: Placeholder for repository owner. - repo: Placeholder for repository name. - runner_id: Placeholder for runenr_id. - """ - pass - - def delete_self_hosted_runner_from_org(self, org: str, runner_id: str): - """A placeholder method for test stub. - - Args: - org: Placeholder for organisation. - runner_id: Placeholder for runner id. - """ - pass - - -@dataclass -class MockRunner: - """Mock of a runner. - - Attributes: - name: The name of the runner. - instance_id: The instance id of the runner. - metadata: Metadata of the server. - cloud_state: The cloud state of the runner. - platform_state: The github state of the runner. - health: The health state of the runner. - created_at: The cloud creation time of the runner. - deletable: If the runner is deletable. - """ - - name: str - instance_id: InstanceID - metadata: RunnerMetadata - cloud_state: CloudRunnerState - platform_state: PlatformRunnerState - health: bool - created_at: datetime - deletable: bool = False - - def __init__(self, instance_id: InstanceID): - """Construct the object. - - Args: - instance_id: InstanceID of the runner. - """ - self.name = instance_id.name - self.instance_id = instance_id - self.metadata = RunnerMetadata() - self.cloud_state = CloudRunnerState.ACTIVE - self.platform_state = PlatformRunnerState.IDLE - self.health = True - # By default a runner that has just being created. - self.created_at = datetime.now(timezone.utc) - - def to_cloud_runner(self) -> CloudRunnerInstance: - """Construct CloudRunnerInstance from this object. - - Returns: - The CloudRunnerInstance instance. - """ - return CloudRunnerInstance( - name=self.name, - metadata=self.metadata, - instance_id=self.instance_id, - health=self.health, - state=self.cloud_state, - created_at=self.created_at, - ) - - -@dataclass -class SharedMockRunnerManagerState: - """State shared by mock runner managers. - - For sharing the mock runner states between MockCloudRunnerManager and MockGitHubRunnerPlatform. - - Attributes: - runners: The runners. - """ - - runners: dict[InstanceID, MockRunner] - - def __init__(self): - """Construct the object.""" - self.runners = {} - - -class MockCloudRunnerManager(CloudRunnerManager): - """Mock of CloudRunnerManager. - - Metrics is not supported in this mock. - - Attributes: - name_prefix: The naming prefix for runners managed. - prefix: The naming prefix for runners managed. - state: The shared state between mocks runner managers. - """ - - def __init__(self, state: SharedMockRunnerManagerState): - """Construct the object. - - Args: - state: The shared state between cloud and github runner managers. - """ - self.prefix = f"mock_{secrets.token_hex(4)}" - self.state = state - - @property - def name_prefix(self) -> str: - """Get the name prefix of the self-hosted runners.""" - return self.prefix - - def create_runner( - self, - runner_identity: RunnerIdentity, - runner_context: RunnerContext, - ) -> None: - """Create a self-hosted runner. - - Args: - runner_identity: Identity of the runner to create. - runner_context: Context for the runner. - - Returns: - The CloudRunnerInstance for the runner - """ - runner = MockRunner(runner_identity.instance_id) - self.state.runners[runner_identity.instance_id] = runner - return runner.to_cloud_runner() - - def get_runners(self) -> Sequence[CloudRunnerInstance]: - """Get cloud self-hosted runners. - - Returns: - Information on the runner instances. - """ - return [runner.to_cloud_runner() for runner in self.state.runners.values()] - - def delete_runner(self, instance_id: InstanceID) -> RunnerMetrics | None: - """Delete self-hosted runner. - - Args: - instance_id: The instance id of the runner to delete. - - Returns: - Any runner metrics produced during deletion. - """ - runner = self.state.runners.pop(instance_id, None) - if runner is not None: - return MagicMock() - return [] - - def cleanup(self) -> None: - """Cleanup runner dangling resources on the cloud.""" - - -class MockGitHubRunnerPlatform(PlatformProvider): - """Mock of GitHubRunnerPlatform. - - Attributes: - github: The GitHub client. - name_prefix: The naming prefix for runner managed. - state: The shared state between mock runner managers. - path: The GitHub path to register the runners under. - """ - - def __init__(self, name_prefix: str, path: GitHubPath, state: SharedMockRunnerManagerState): - """Construct the object. - - Args: - name_prefix: The naming prefix for runner managed. - path: The GitHub path to register the runners under. - state: The shared state between mock runner managers. - """ - self.github = GithubClient("mock_token") - self.github._client = MockGhapiClient("mock_token") - self.name_prefix = name_prefix - self.state = state - self.path = path - - def get_runner_health( - self, - runner_identity: RunnerIdentity, - ) -> PlatformRunnerHealth: - """Get info on self-hosted runner. - - Args: - runner_identity: Identity of the runner. - - Returns: - Information about the health of the runner - """ - if runner_identity.instance_id in self.state.runners: - runner = self.state.runners[runner_identity.instance_id] - return PlatformRunnerHealth( - identity=runner_identity, - online=runner.platform_state != PlatformRunnerState.OFFLINE, - busy=runner.platform_state == PlatformRunnerState.BUSY, - deletable=runner.deletable, - ) - return PlatformRunnerHealth( - identity=runner_identity, - online=False, - busy=False, - deletable=True, - ) - - def get_runners_health( - self, requested_runners: list[RunnerIdentity] - ) -> "list[PlatformRunnerHealth]": - """Get information from the requested runners health. - - Args: - requested_runners: List of runners to get health information for. - - Returns: - Health information for the runners. - """ - found_identities = [] - for identity in requested_runners: - if identity.instance_id in self.state.runners: - runner = self.state.runners[identity.instance_id] - if runner.health: - found_identities.append(identity) - requested_runners = [self.get_runner_health(identity) for identity in found_identities] - return RunnersHealthResponse(requested_runners=requested_runners) - - def get_runner_context( - self, metadata: RunnerMetadata, instance_id: str, labels: list[str] - ) -> tuple[RunnerContext, SelfHostedRunner]: - """Get the registration JIT token for registering runners on GitHub. - - Args: - metadata: Metadata of the server. - instance_id: Instance ID of the runner. - labels: Labels for the runner. - - Returns: - The registration token and the SelfHostedRunner - """ - runner = MagicMock(spec=list(SelfHostedRunner.__fields__.keys())) - runner.id = 5 - return RunnerContext(shell_run_script="fake-agent"), runner - - def get_runners( - self, states: Iterable[PlatformRunnerState] | None = None - ) -> tuple[SelfHostedRunner, ...]: - """Get the runners. - - Args: - states: The states to filter for. - - Returns: - List of runners. - """ - if states is None: - states = [member.value for member in PlatformRunnerState] - - platform_state_set = set(states) - runner_id = random.randint(1, 1000000) - return tuple( - SelfHostedRunner( - busy=runner.platform_state == PlatformRunnerState.BUSY, - id=runner_id, - labels=[], - instance_id=InstanceID.build_from_name(self.name_prefix, runner.name), - status=( - GitHubRunnerStatus.OFFLINE - if runner.platform_state == PlatformRunnerState.OFFLINE - else GitHubRunnerStatus.ONLINE - ), - metadata=RunnerMetadata(platform_name="github", runner_id=str(runner_id)), - ) - for runner in self.state.runners.values() - if runner.platform_state in platform_state_set - ) - - def delete_runner(self, runner_identity: RunnerIdentity) -> None: - """Delete a runner. - - Args: - runner_identity: Runner to delete. - """ - if runner_identity.instance_id in self.state.runners: - del self.state.runners[runner_identity.instance_id] - - def check_job_been_picked_up(self, metadata: RunnerMetadata, job_url: HttpUrl) -> bool: - """Check if the job has already been picked up. - - Args: - metadata: Metadata of the instance. - job_url: The URL of the job. - - Raises: - NotImplementedError: Work in progress. - """ - raise NotImplementedError - - def get_job_info( - self, metadata: RunnerMetadata, repository: str, workflow_run_id: str, runner: InstanceID - ) -> JobInfo: - """Get the Job info from the provider. - - Args: - metadata: Metadata of the runner. - repository: repository to get the job from. - workflow_run_id: workflow run id of the job. - runner: runner to get the job from. - - Raises: - NotImplementedError: Work in progress. - """ - raise NotImplementedError diff --git a/github-runner-manager/tests/unit/openstack_cloud/test_openstack_cloud.py b/github-runner-manager/tests/unit/openstack_cloud/test_openstack_cloud.py index 04f776985e..9f169e6411 100644 --- a/github-runner-manager/tests/unit/openstack_cloud/test_openstack_cloud.py +++ b/github-runner-manager/tests/unit/openstack_cloud/test_openstack_cloud.py @@ -11,22 +11,28 @@ import keystoneauth1.exceptions import openstack +import openstack.exceptions import pytest from openstack.compute.v2.keypair import Keypair from openstack.connection import Connection from openstack.network.v2.security_group import SecurityGroup as OpenstackSecurityGroup from openstack.network.v2.security_group_rule import SecurityGroupRule +from pytest import LogCaptureFixture +import github_runner_manager.openstack_cloud.openstack_cloud from github_runner_manager.errors import OpenStackError, SSHError from github_runner_manager.openstack_cloud.openstack_cloud import ( _MAX_NOVA_COMPUTE_API_VERSION, _MIN_KEYPAIR_AGE_IN_SECONDS_BEFORE_DELETION, _TEST_STRING, DEFAULT_SECURITY_RULES, + InstanceID, OpenstackCloud, OpenStackCredentials, + _DeleteKeypairConfig, get_missing_security_rules, ) +from tests.unit.fake_runner_managers import FakeOpenstackCloud FAKE_ARG = "fake" FAKE_PREFIX = "fake_prefix" @@ -52,6 +58,19 @@ def openstack_cloud_fixture(monkeypatch): return OpenstackCloud(creds, FAKE_PREFIX, FAKE_ARG) +@pytest.fixture(name="mock_openstack_conn", scope="function") +def mock_openstack_conn_fixture(monkeypatch: pytest.MonkeyPatch): + """Patch OpenStack connection.""" + connection_mock = MagicMock() + connection_mock.__enter__.return_value = connection_mock + monkeypatch.setattr( + github_runner_manager.openstack_cloud.openstack_cloud.openstack, + "connect", + MagicMock(return_value=connection_mock), + ) + return connection_mock + + @pytest.mark.parametrize( "public_method, args", [ @@ -65,7 +84,6 @@ def openstack_cloud_fixture(monkeypatch): id="launch_instance", ), pytest.param("get_instance", {"instance_id": FAKE_ARG}, id="get_instance"), - pytest.param("delete_instance", {"instance_id": FAKE_ARG}, id="delete_instance"), pytest.param("get_instances", {}, id="get_instances"), pytest.param("cleanup", {}, id="cleanup"), ], @@ -305,6 +323,105 @@ def mock_ssh_connection(*_args, **_kwargs) -> MagicMock: assert "No connectable SSH addresses found" in str(err.value) +# We test this internal method because this fails silently without bubbling up exceptions due to +# it's non-critical nature. +def test__delete_keypair_fail( + openstack_cloud: OpenstackCloud, mock_openstack_conn: MagicMock, caplog: LogCaptureFixture +): + """ + arrange: given a mocked openstack delete_keypair method that returns False. + act: when _delete_keypair method is called. + assert: None is returned and the failure is logged. + """ + mock_openstack_conn.delete_keypair = MagicMock(return_value=False) + test_key_instance_id = InstanceID(prefix="test-key-delete", reactive=False, suffix="fail") + + assert ( + openstack_cloud._delete_keypair( + _DeleteKeypairConfig( + keys_dir=MagicMock(), instance_id=test_key_instance_id, conn=mock_openstack_conn + ) + ) + is None + ) + assert f"Failed to delete key: {test_key_instance_id.name}" in caplog.messages + + +def test__delete_keypair_error( + openstack_cloud: OpenstackCloud, mock_openstack_conn: MagicMock, caplog: LogCaptureFixture +): + """ + arrange: given a mocked openstack delete_keypair method that returns False. + act: when _delete_keypair method is called. + assert: None is returned and the failure is logged. + """ + mock_openstack_conn.delete_keypair = MagicMock( + side_effect=[openstack.exceptions.ResourceTimeout()] + ) + test_key_instance_id = InstanceID(prefix="test-key-delete", reactive=False, suffix="fail") + + assert ( + openstack_cloud._delete_keypair( + _DeleteKeypairConfig( + keys_dir=MagicMock(), instance_id=test_key_instance_id, conn=mock_openstack_conn + ) + ) + is None + ) + assert f"Error attempting to delete key: {test_key_instance_id.name}" in caplog.messages + + +def test_delete_instances_partial_server_delete_failure( + monkeypatch: pytest.MonkeyPatch, openstack_cloud: OpenstackCloud, caplog: LogCaptureFixture +): + """ + arrange: given a mocked openstack connection that errors on few failed requests. + act: when delete_instances method is called. + assert: successfully deleted instance IDs are returned and failed instances are logged. + """ + successful_delete_id = InstanceID(prefix="success", reactive=False, suffix="") + already_deleted_id = InstanceID(prefix="already_deleted", reactive=False, suffix="") + timeout_id = InstanceID(prefix="timeout error", reactive=False, suffix="") + mock_cloud = FakeOpenstackCloud( + initial_servers=[successful_delete_id, timeout_id], + server_to_errors={timeout_id: openstack.exceptions.ResourceTimeout()}, + ) + monkeypatch.setattr( + github_runner_manager.openstack_cloud.openstack_cloud.openstack, + "connect", + MagicMock(return_value=mock_cloud), + ) + + deleted_instance_ids = openstack_cloud.delete_instances( + instance_ids=[successful_delete_id, already_deleted_id, timeout_id] + ) + + assert successful_delete_id in deleted_instance_ids + assert already_deleted_id not in deleted_instance_ids + assert timeout_id not in deleted_instance_ids + assert f"Failed to delete OpenStack VM instance: {timeout_id}" in caplog.messages + + +def test_delete_instances( + openstack_cloud: OpenstackCloud, + mock_openstack_conn: MagicMock, +): + """ + arrange: given a mocked openstack connection. + act: when delete_instances method is called. + assert: deleted instance IDs are returned. + """ + mock_openstack_conn.delete_server = MagicMock(side_effect=[True, False]) + successful_delete_id = InstanceID(prefix="success", reactive=False, suffix="") + already_deleted_id = InstanceID(prefix="already_deleted", reactive=False, suffix="") + + deleted_instance_ids = openstack_cloud.delete_instances( + instance_ids=[successful_delete_id, already_deleted_id] + ) + + assert deleted_instance_ids == [successful_delete_id] + + @pytest.mark.parametrize( "max_compute_api_version, expected_version", [ diff --git a/github-runner-manager/tests/unit/openstack_cloud/test_openstack_runner_manager.py b/github-runner-manager/tests/unit/openstack_cloud/test_openstack_runner_manager.py index 444645c684..07a406cbe3 100644 --- a/github-runner-manager/tests/unit/openstack_cloud/test_openstack_runner_manager.py +++ b/github-runner-manager/tests/unit/openstack_cloud/test_openstack_runner_manager.py @@ -4,20 +4,11 @@ """Module for unit-testing OpenStack runner manager.""" import logging import textwrap -from datetime import datetime, timezone -from typing import Iterable from unittest.mock import MagicMock import pytest from github_runner_manager.configuration import ProxyConfig, SupportServiceConfig, UserInfo -from github_runner_manager.manager.cloud_runner_manager import ( - CodeInformation, - PostJobMetrics, - PostJobStatus, - PreJobMetrics, - RunnerMetrics, -) from github_runner_manager.manager.models import ( InstanceID, RunnerContext, @@ -25,19 +16,12 @@ RunnerMetadata, ) from github_runner_manager.metrics import runner -from github_runner_manager.metrics.runner import PullFileError -from github_runner_manager.openstack_cloud import openstack_cloud -from github_runner_manager.openstack_cloud.constants import ( - POST_JOB_METRICS_FILE_NAME, - PRE_JOB_METRICS_FILE_NAME, - RUNNER_INSTALLED_TS_FILE_NAME, -) from github_runner_manager.openstack_cloud.openstack_cloud import OpenstackCloud from github_runner_manager.openstack_cloud.openstack_runner_manager import ( OpenStackRunnerManager, OpenStackRunnerManagerConfig, + runner_metrics, ) -from tests.unit.factories import openstack_factory logger = logging.getLogger(__name__) @@ -219,158 +203,31 @@ def test_create_runner_without_aproxy( assert "aproxy" not in openstack_cloud.launch_instance.call_args.kwargs["cloud_init"] -def _params_test_delete_extract_metrics(): - """Builds parametrized input for the test_delete_extract_metrics. - - The following values are returned: - runner_installed_metrics,pre_job_metrics,post_job_metrics,result +def test_delete_vms(runner_manager: OpenStackRunnerManager): """ - openstack_created_at = ( - datetime.strptime(openstack_factory.SERVER_CREATED_AT, "%Y-%m-%dT%H:%M:%SZ") - .replace(tzinfo=timezone.utc) - .timestamp() - ) - openstack_installed_at = openstack_created_at + 20 - pre_job_timestamp = openstack_installed_at + 20 - post_job_timestamp = openstack_installed_at + 20 - pre_job_metrics_str = f"""{{ - "timestamp": {pre_job_timestamp}, - "workflow": "Workflow Dispatch Tests", - "workflow_run_id": "13831611664", - "repository": "canonical/github-runner-operator", - "event": "workflow_dispatch" - }}""" - post_job_metrics_str = f"""{{ - "timestamp": {post_job_timestamp}, "status": "normal", "status_info": {{"code" : "200"}} - }}""" + arrange: given a mocked cloud service. + act: when delete_vms method is called. + assert: the mocked service call is made and the deleted instance IDs are returned. + """ + test_instance_ids = [InstanceID(prefix="test-prefix", reactive=None, suffix="test-suffix")] + mock_cloud = MagicMock() + mock_cloud.delete_instances = MagicMock(return_value=test_instance_ids) + runner_manager._openstack_cloud = mock_cloud - return [ - pytest.param(None, None, None, None, id="All None. No metrics returned."), - pytest.param( - "", None, None, None, id="Invalid runner-installed metrics. No metrics returned." - ), - pytest.param( - str(openstack_installed_at), - None, - None, - RunnerMetrics( - instance_id=InstanceID( - prefix=OPENSTACK_INSTANCE_PREFIX, reactive=False, suffix="unhealthy" - ), - installation_start_timestamp=openstack_created_at, - installed_timestamp=openstack_installed_at, - metadata=RunnerMetadata(), - ), - id="Only installed_timestamp. Metric returned.", - ), - pytest.param( - str(openstack_installed_at), - pre_job_metrics_str, - None, - RunnerMetrics( - instance_id=InstanceID( - prefix=OPENSTACK_INSTANCE_PREFIX, reactive=False, suffix="unhealthy" - ), - installation_start_timestamp=openstack_created_at, - installed_timestamp=openstack_installed_at, - pre_job=PreJobMetrics( - timestamp=pre_job_timestamp, - workflow="Workflow Dispatch Tests", - workflow_run_id="13831611664", - repository="canonical/github-runner-operator", - event="workflow_dispatch", - ), - metadata=RunnerMetadata(), - ), - id="installed_timestamp and pre_job_metrics. Metric returned.", - ), - pytest.param( - str(openstack_installed_at), - pre_job_metrics_str, - post_job_metrics_str, - RunnerMetrics( - metadata=RunnerMetadata(), - instance_id=InstanceID( - prefix=OPENSTACK_INSTANCE_PREFIX, reactive=False, suffix="unhealthy" - ), - installation_start_timestamp=openstack_created_at, - installed_timestamp=openstack_installed_at, - pre_job=PreJobMetrics( - timestamp=pre_job_timestamp, - workflow="Workflow Dispatch Tests", - workflow_run_id="13831611664", - repository="canonical/github-runner-operator", - event="workflow_dispatch", - ), - post_job=PostJobMetrics( - timestamp=post_job_timestamp, - status=PostJobStatus.NORMAL, - status_info=CodeInformation(code=200), - ), - ), - id="installed_timestamp, pre_job_metrics and post_job_metrics. Metric returned", - ), - ] + assert test_instance_ids == runner_manager.delete_vms(instance_ids=test_instance_ids) + mock_cloud.delete_instances.assert_called_once() -@pytest.mark.parametrize( - "runner_installed_metrics,pre_job_metrics,post_job_metrics,result", - _params_test_delete_extract_metrics(), -) -def test_delete_extract_metrics( - runner_manager: OpenStackRunnerManager, - runner_installed_metrics: str | None, - pre_job_metrics: str | None, - post_job_metrics: str | None, - result: Iterable[RunnerMetrics], - monkeypatch: pytest.MonkeyPatch, -): +def test_extract_metrics(runner_manager: OpenStackRunnerManager, monkeypatch: pytest.MonkeyPatch): """ - arrange: Given different values for values of metrics for a runner. - act: Delete the runner for those metrics. - assert: The expected RunnerMetrics object is obtained, or None if there should not be one. + arrange: given a mocked metrics service. + act: when extract_metrics method is called. + assert: converted metrics are returned. """ - ssh_pull_file_mock = MagicMock() - monkeypatch.setattr( - "github_runner_manager.metrics.runner.ssh_pull_file", - ssh_pull_file_mock, + pull_metrics_mock = MagicMock( + return_value=[(test_metric_one := MagicMock()), (test_metric_two := MagicMock())] ) + monkeypatch.setattr(runner_metrics, "pull_runner_metrics", pull_metrics_mock) - def _ssh_pull_file(remote_path, *args, **kwargs): - """Get a file from the runner.""" - logger.info("ssh_pull_file: remote_path %s", remote_path) - res = None - if remote_path == str(PRE_JOB_METRICS_FILE_NAME): - res = pre_job_metrics - elif remote_path == str(POST_JOB_METRICS_FILE_NAME): - res = post_job_metrics - elif remote_path == str(RUNNER_INSTALLED_TS_FILE_NAME): - res = runner_installed_metrics - if res is None: - raise PullFileError("Nothing found or invalid file.") - return res - - ssh_pull_file_mock.side_effect = _ssh_pull_file - - instance_id = InstanceID( - prefix=OPENSTACK_INSTANCE_PREFIX, reactive=False, suffix="unhealthy" - ).name - openstack_cloud_mock = _create_openstack_cloud_mock(instance_id) - runner_manager._openstack_cloud = openstack_cloud_mock - - runner_metrics = runner_manager.delete_runner(instance_id) - - assert runner_metrics == result - - -def _create_openstack_cloud_mock(instance_name: str) -> MagicMock: - """Create an OpenstackCloud mock which returns servers with a given list of server names.""" - openstack_cloud_mock = MagicMock(spec=OpenstackCloud) - openstack_cloud_mock.get_instance.return_value = openstack_cloud.OpenstackInstance( - server=openstack_factory.ServerFactory( - status="ACTIVE", - name=instance_name, - ), - prefix=OPENSTACK_INSTANCE_PREFIX, - ) - return openstack_cloud_mock + metrics = runner_manager.extract_metrics(instance_ids=MagicMock()) + assert metrics == [test_metric_one.to_runner_metrics(), test_metric_two.to_runner_metrics()] diff --git a/github-runner-manager/tests/unit/platform/test_github_provider.py b/github-runner-manager/tests/unit/platform/test_github_provider.py index 3175938b4e..23588792d6 100644 --- a/github-runner-manager/tests/unit/platform/test_github_provider.py +++ b/github-runner-manager/tests/unit/platform/test_github_provider.py @@ -15,6 +15,7 @@ GitHubRunnerPlatform, ) from github_runner_manager.platform.platform_provider import ( + DeleteRunnerBusyError, PlatformRunnerHealth, RunnersHealthResponse, ) @@ -239,3 +240,34 @@ def test_get_runners_health( runners_health_response = platform.get_runners_health(requested_runners) assert runners_health_response == expected_health_response + + +def test_github_provider_delete_busy_runner_error(): + """ + arrange: given a mocked GitHub client that raises DeleteRunnerBusyError. + act: when GitHubRunnerPlatform.delete_runners is called. + assert: act: no ids are returned. + """ + mock_github_client = MagicMock() + mock_github_client.delete_runner.side_effect = DeleteRunnerBusyError + github_provider = GitHubRunnerPlatform( + prefix="test", path="test", github_client=mock_github_client + ) + test_delete_ids = ["1", "2", "3"] + + assert github_provider.delete_runners(test_delete_ids) == [] + + +def test_github_provider_delete_runners(): + """ + arrange: given a mocked GitHub client. + act: when GitHubRunnerPlatform.delete_runners is called. + assert: act: the deleted runner IDs are returned. + """ + mock_github_client = MagicMock() + github_provider = GitHubRunnerPlatform( + prefix="test", path="test", github_client=mock_github_client + ) + test_delete_ids = ["1", "2", "3"] + + assert sorted(github_provider.delete_runners(test_delete_ids)) == sorted(test_delete_ids) diff --git a/github-runner-manager/tests/unit/test_runner_scaler.py b/github-runner-manager/tests/unit/test_runner_scaler.py index d8e2739003..89329bc833 100644 --- a/github-runner-manager/tests/unit/test_runner_scaler.py +++ b/github-runner-manager/tests/unit/test_runner_scaler.py @@ -30,13 +30,16 @@ GitHubPath, GitHubRepo, ) -from github_runner_manager.errors import CloudError, ReconcileError from github_runner_manager.manager import runner_manager as runner_manager_module from github_runner_manager.manager.cloud_runner_manager import CloudRunnerState from github_runner_manager.manager.models import InstanceID -from github_runner_manager.manager.runner_manager import FlushMode, RunnerManager -from github_runner_manager.manager.runner_scaler import RunnerScaler -from github_runner_manager.metrics.events import Reconciliation +from github_runner_manager.manager.runner_manager import ( + IssuedMetricEventsStats, + RunnerInstance, + RunnerManager, +) +from github_runner_manager.manager.runner_scaler import FlushMode, RunnerInfo, RunnerScaler +from github_runner_manager.metrics.events import RunnerStart, RunnerStop from github_runner_manager.openstack_cloud.configuration import ( OpenStackConfiguration, OpenStackCredentials, @@ -47,11 +50,7 @@ ) from github_runner_manager.platform.github_provider import PlatformRunnerState from github_runner_manager.reactive.types_ import ReactiveProcessConfig -from tests.unit.mock_runner_managers import ( - MockCloudRunnerManager, - MockGitHubRunnerPlatform, - SharedMockRunnerManagerState, -) +from tests.unit.factories.runner_instance_factory import RunnerInstanceFactory logger = logging.getLogger(__name__) @@ -79,16 +78,6 @@ def github_path_fixture() -> GitHubPath: return GitHubRepo(owner="mock_owner", repo="mock_repo") -@pytest.fixture(scope="function", name="mock_runner_managers") -def mock_runner_managers_fixture( - github_path: GitHubPath, -) -> tuple[MockCloudRunnerManager, MockGitHubRunnerPlatform]: - state = SharedMockRunnerManagerState() - mock_cloud = MockCloudRunnerManager(state) - mock_github = MockGitHubRunnerPlatform(mock_cloud.name_prefix, github_path, state) - return (mock_cloud, mock_github) - - @pytest.fixture(scope="function", name="issue_events_mock") def issue_events_mock_fixture(monkeypatch: pytest.MonkeyPatch): issue_events_mock = MagicMock() @@ -388,328 +377,188 @@ def test_build_runner_scaler( ) -def test_get_no_runner(runner_manager: RunnerManager, user_info: UserInfo): - """ - Arrange: A RunnerScaler with no runners. - Act: Get runner information. - Assert: Information should contain no runners. - """ - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=0, max_quantity=0) - assert_runner_info(runner_scaler, online=0) - - -def test_flush_no_runner(runner_manager: RunnerManager, user_info: UserInfo): - """ - Arrange: A RunnerScaler with no runners. - Act: - 1. Flush idle runners. - 2. Flush busy runners. - Assert: - 1. No change in number of runners. Runner info should contain no runners. - 2. No change in number of runners. - """ - # 1. - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=0, max_quantity=0) - diff = runner_scaler.flush(flush_mode=FlushMode.FLUSH_IDLE) - assert diff == 0 - assert_runner_info(runner_scaler, online=0) - - # 2. - diff = runner_scaler.flush(flush_mode=FlushMode.FLUSH_BUSY) - assert diff == 0 - assert_runner_info(runner_scaler, online=0) - - -def test_reconcile_runner_create_one(runner_manager: RunnerManager, user_info: UserInfo): - """ - Arrange: A RunnerScaler with no runners. - Act: Reconcile to no runners. - Assert: No changes. Runner info should contain no runners. - """ - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=0, max_quantity=0) - diff = runner_scaler.reconcile() - assert diff == 0 - assert_runner_info(runner_scaler, online=0) - - -def test_reconcile_runner_create_one_reactive( - monkeypatch: pytest.MonkeyPatch, runner_manager: RunnerManager, user_info: UserInfo +@pytest.mark.parametrize( + "runners, expected_runner_info", + [ + pytest.param( + [], + RunnerInfo(online=0, busy=0, offline=0, unknown=0, runners=(), busy_runners=()), + id="No runners", + ), + pytest.param( + [busy_runner := RunnerInstanceFactory(platform_state=PlatformRunnerState.BUSY)], + RunnerInfo( + online=1, + busy=1, + offline=0, + unknown=0, + runners=(busy_runner.name,), + busy_runners=(busy_runner.name,), + ), + id="One busy runner", + ), + pytest.param( + [idle_runner := RunnerInstanceFactory(platform_state=PlatformRunnerState.IDLE)], + RunnerInfo( + online=1, + busy=0, + offline=0, + unknown=0, + runners=(idle_runner.name,), + busy_runners=(), + ), + id="One idle runner", + ), + pytest.param( + [offline_runner := RunnerInstanceFactory(platform_state=PlatformRunnerState.OFFLINE)], + RunnerInfo( + online=0, + busy=0, + offline=1, + unknown=0, + runners=(), + busy_runners=(), + ), + id="One offline runner", + ), + pytest.param( + [unknown_runner := RunnerInstanceFactory(platform_state=None)], + RunnerInfo( + online=0, + busy=0, + offline=0, + unknown=1, + runners=(), + busy_runners=(), + ), + id="One unknown runner", + ), + pytest.param( + [busy_runner, idle_runner, offline_runner, unknown_runner], + RunnerInfo( + online=2, + busy=1, + offline=1, + unknown=1, + runners=(busy_runner.name, idle_runner.name), + busy_runners=(busy_runner.name,), + ), + id="One runner of each type", + ), + ], +) +def test_runner_scaler_get_runner_info( + runners: list[RunnerInstance], expected_runner_info: RunnerInfo ): """ - Arrange: Prepare one RunnerScaler in reactive mode. - Fake the reconcile function in reactive to return its input. - Act: Call reconcile with base quantity 0 and max quantity 5. - Assert: 5 processes should be returned in the result of the reconcile. + arrange: given a mock runner manager. + act: when RunnerScaler.get_runner_info is called. + assert the expected runner info is extracted. """ - reactive_process_config = MagicMock() + runner_manager = MagicMock() + runner_manager.get_runners.return_value = runners runner_scaler = RunnerScaler( - runner_manager, reactive_process_config, user_info, base_quantity=0, max_quantity=5 - ) - - from github_runner_manager.reactive.runner_manager import ReconcileResult - - def _fake_reactive_reconcile( - expected_quantity: int, runner_manager, reactive_process_config, user, python_path - ): - """Reactive reconcile fake.""" - return ReconcileResult(processes_diff=expected_quantity, metric_stats={"event": ""}) - - monkeypatch.setattr( - "github_runner_manager.reactive.runner_manager.reconcile", - MagicMock(side_effect=_fake_reactive_reconcile), + runner_manager=runner_manager, + reactive_process_config=None, + user=MagicMock(), + base_quantity=0, + max_quantity=0, ) - diff = runner_scaler.reconcile() - assert diff == 5 - assert_runner_info(runner_scaler, online=0) - -def test_reconcile_error_still_issue_metrics( - runner_manager: RunnerManager, - monkeypatch: pytest.MonkeyPatch, - issue_events_mock: MagicMock, - user_info: UserInfo, + assert runner_scaler.get_runner_info() == expected_runner_info + + +@pytest.mark.parametrize( + "cleanup_metrics, flush_metrics, expected_flushed", + [ + pytest.param({}, {}, 0, id="No changes"), + pytest.param({RunnerStart: 1}, {}, 0, id="No runner stop metrics"), + pytest.param({RunnerStop: 1}, {}, 1, id="Runner stop metric from cleanup"), + pytest.param({}, {RunnerStop: 1}, 1, id="Runner stop metric from flush"), + pytest.param( + {RunnerStop: 1}, + {RunnerStop: 1}, + 2, + id="Runner stop metrics from cleanup and flush", + ), + ], +) +def test_runner_scaler_flush_extract_metrics( + cleanup_metrics: IssuedMetricEventsStats, + flush_metrics: IssuedMetricEventsStats, + expected_flushed: int, ): """ - Arrange: A RunnerScaler with no runners which raises an error on reconcile. - Act: Reconcile to one runner. - Assert: ReconciliationEvent should be issued. + arrange: given a mocked runner manager with that returns the given metrics. + act: when RunnerScaler.flush is called. + assert: the expected number of flushed runners from metrics is returned. """ - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=1, max_quantity=0) - monkeypatch.setattr( - runner_scaler._manager, "cleanup", MagicMock(side_effect=Exception("Mock error")) - ) - with pytest.raises(Exception): - runner_scaler.reconcile() - issue_events_mock.assert_called_once() - issued_event = issue_events_mock.call_args[0][0] - assert isinstance(issued_event, Reconciliation) + runner_manager = MagicMock() + runner_manager.cleanup.return_value = cleanup_metrics + runner_manager.flush_runners.return_value = flush_metrics - -def test_reconcile_raises_reconcile_error( - runner_manager: RunnerManager, - monkeypatch: pytest.MonkeyPatch, - issue_events_mock: MagicMock, - user_info: UserInfo, -): - """ - Arrange: A RunnerScaler with no runners which raises a Cloud error on reconcile. - Act: Reconcile to one runner. - Assert: ReconcileError should be raised. - """ - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=1, max_quantity=0) - monkeypatch.setattr( - runner_scaler._manager, "cleanup", MagicMock(side_effect=CloudError("Mock error")) + runner_scaler = RunnerScaler( + runner_manager=runner_manager, + reactive_process_config=None, + user=MagicMock(), + base_quantity=0, + max_quantity=0, ) - with pytest.raises(ReconcileError) as exc: - runner_scaler.reconcile() - assert "Failed to reconcile runners." in str(exc.value) + assert runner_scaler.flush() == expected_flushed -def test_one_runner(runner_manager: RunnerManager, user_info: UserInfo): - """ - Arrange: A RunnerScaler with no runners. - Act: - 1. Reconcile to one runner. - 2. Reconcile to one runner. - 3. Flush idle runners. - 4. Reconcile to one runner. - Assert: - 1. Runner info has one runner. - 2. No changes to number of runner. - 3. Runner info has one runner. - """ - # 1. - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=1, max_quantity=0) - diff = runner_scaler.reconcile() - assert diff == 1 - assert_runner_info(runner_scaler, online=1) - - # 2. - diff = runner_scaler.reconcile() - assert diff == 0 - assert_runner_info(runner_scaler, online=1) - # 3. - runner_scaler.flush(flush_mode=FlushMode.FLUSH_IDLE) - assert_runner_info(runner_scaler, online=0) - - # 3. - diff = runner_scaler.reconcile() - assert diff == 1 - assert_runner_info(runner_scaler, online=1) - - -def test_flush_busy_on_idle_runner(runner_scaler_one_runner: RunnerScaler): - """ - Arrange: A RunnerScaler with one idle runner. - Act: Run flush busy runner. - Assert: No runners. - """ - runner_scaler = runner_scaler_one_runner - - runner_scaler.flush(flush_mode=FlushMode.FLUSH_BUSY) - assert_runner_info(runner_scaler, online=0) - - -def test_flush_busy_on_busy_runner( - runner_scaler_one_runner: RunnerScaler, -): - """ - Arrange: A RunnerScaler with one busy runner. - Act: Run flush busy runner. - Assert: No runners. - """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.BUSY) - - runner_scaler.flush(flush_mode=FlushMode.FLUSH_BUSY) - assert_runner_info(runner_scaler, online=0) - - -def test_get_runner_one_busy_runner( - runner_scaler_one_runner: RunnerScaler, -): - """ - Arrange: A RunnerScaler with one busy runner. - Act: Run get runners. - Assert: One busy runner. - """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.BUSY) - - assert_runner_info(runner_scaler=runner_scaler, online=1, busy=1) - - -def test_get_runner_offline_runner(runner_scaler_one_runner: RunnerScaler): - """ - Arrange: A RunnerScaler with one offline runner. - Act: Run get runners. - Assert: One offline runner. - """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.OFFLINE) - - assert_runner_info(runner_scaler=runner_scaler, offline=1) - - -def test_get_runner_unknown_runner(runner_scaler_one_runner: RunnerScaler): - """ - Arrange: A RunnerScaler with one offline runner. - Act: Run get runners. - Assert: One offline runner. - """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, health=False) - assert_runner_info(runner_scaler=runner_scaler, unknown=1) - - -def test_flush_idle_on_starting_offline_runner( - runner_scaler_one_runner: RunnerScaler, - runner_manager: RunnerManager, -): - """ - Arrange: A RunnerScaler with one offline runner that could be starting. - Act: Run flush idle runner. - Assert: The runner should not be deleted. - """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.OFFLINE) - runners_before = runner_manager.get_runners() - - runner_scaler.flush(flush_mode=FlushMode.FLUSH_IDLE) - assert_runner_info(runner_scaler, offline=1) - runners_after = runner_manager.get_runners() - assert runners_before[0].instance_id == runners_after[0].instance_id - - -def test_flush_idle_on_old_offline_runner( - runner_scaler_one_runner: RunnerScaler, - runner_manager: RunnerManager, -): +@pytest.mark.parametrize( + "flush_mode, expected_flush_mode", + [ + pytest.param(FlushMode.FLUSH_IDLE, FlushMode.FLUSH_IDLE, id="flush_idle"), + pytest.param(FlushMode.FLUSH_BUSY, FlushMode.FLUSH_BUSY, id="flush_busy"), + ], +) +def test_runner_scaler_flush_mode(flush_mode: FlushMode, expected_flush_mode: FlushMode): """ - Arrange: A RunnerScaler with one offline runner that had enough time to start. - Act: Run flush idle runner. - Assert: The runner should be deleted. No one should be created. + arrange: given a mocked runner manager. + act: when RunnerScaler.flush is called with the given flush mode. + assert: flush_runners is called with expected mode. """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.OFFLINE, old_runner=True) - - runner_scaler.flush(flush_mode=FlushMode.FLUSH_IDLE) - assert_runner_info(runner_scaler, offline=0) + runner_manager = MagicMock() + RunnerScaler( + runner_manager=runner_manager, + reactive_process_config=None, + user=MagicMock(), + base_quantity=0, + max_quantity=0, + ).flush(flush_mode=flush_mode) -def test_reconcile_on_starting_offline_runner( - runner_scaler_one_runner: RunnerScaler, - runner_manager: RunnerManager, -): - """ - Arrange: A RunnerScaler with one offline runner that could be starting. - Act: Run reconcile. - Assert: The runner should not be deleted. - """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.OFFLINE) - runners_before = runner_manager.get_runners() - - runner_scaler.reconcile() - assert_runner_info(runner_scaler, offline=1) - runners_after = runner_manager.get_runners() - assert runners_before[0].instance_id == runners_after[0].instance_id + runner_manager.flush_runners.assert_called_with(flush_mode=expected_flush_mode) -def test_reconcile_on_old_offline_runner( - runner_scaler_one_runner: RunnerScaler, - runner_manager: RunnerManager, +@pytest.mark.parametrize( + "runners, quantity, expected_diff", + [ + pytest.param([], 0, 0, id="no difference"), + pytest.param([], 1, 1, id="scale up one runner"), + pytest.param([RunnerInstanceFactory()], 0, -1, id="scale down one runner"), + ], +) +def test_runner_scaler__reconcile_non_reactive( + runners: list[RunnerInstance], quantity: int, expected_diff: int ): """ - Arrange: A RunnerScaler with one offline not busy runner that could have failed to start. - Act: Run reconcile. - Assert: The runner should be deleted. Another one will be created. + arrange: given a mocked runner manager. + act: when RunnerScaler._reconcile_non_reactive is called. + assert: expected runner diff is returned. """ - runner_scaler = runner_scaler_one_runner - set_one_runner_state(runner_scaler, PlatformRunnerState.OFFLINE, old_runner=True) - runners_before = runner_manager.get_runners() - - runner_scaler.reconcile() - assert_runner_info(runner_scaler, online=1, offline=0) - runners_after = runner_manager.get_runners() - assert runners_before[0].instance_id != runners_after[0].instance_id + runner_manager = MagicMock() + runner_manager.get_runners.return_value = runners + result = RunnerScaler( + runner_manager=runner_manager, + reactive_process_config=None, + user=MagicMock(), + base_quantity=0, + max_quantity=0, + )._reconcile_non_reactive(expected_quantity=quantity) -def test_delete_some_runners_in_reconcile(runner_manager: RunnerManager, user_info: UserInfo): - """ - Arrange: Run a reconcile to get 5 runners online. - Act: In a different runner_scaler, reconcile with 2 runners. - Assert: 3 runners should be delete and 2 runners should be online. The busy runner and the - runner without health information should be retained based on the desired ordering. - """ - runner_scaler = RunnerScaler(runner_manager, None, user_info, base_quantity=5, max_quantity=0) - diff = runner_scaler.reconcile() - assert diff == 5 - assert_runner_info(runner_scaler, online=5) - - # Update the runner_dict, so we can check that runners are deleted in order of "inconvenience". - # This test depends on the preservation of insertion order. - # See https://docs.python.org/3.7/library/stdtypes.html#dict.values - runner_dict = runner_scaler._manager._platform.state.runners - initial_mock_runners = list(runner_dict.values()) - initial_mock_runners[0].platform_state = PlatformRunnerState.IDLE - initial_mock_runners[1].platform_state = PlatformRunnerState.OFFLINE - initial_mock_runners[2].platform_state = PlatformRunnerState.BUSY - initial_mock_runners[3].deletable = True - initial_mock_runners[4].health = False # Runner without health information - - second_runner_scaler = RunnerScaler( - runner_manager, None, user_info, base_quantity=2, max_quantity=0 - ) - diff = second_runner_scaler.reconcile() - # Even as 3 runners were deleted, the deletable one was deleted in the cleanup, so - # the runner_scaler returns -2. - assert diff == -2 - assert_runner_info(second_runner_scaler, online=1, busy=1, unknown=1) - - assert len(runner_dict) == 2 - # The busy runner should not be deleted. - assert initial_mock_runners[2].instance_id in runner_dict - # The runner without health information should not be deleted - assert initial_mock_runners[4].instance_id in runner_dict + assert result.runner_diff == expected_diff