diff --git a/kraken/application_outage/actions.py b/kraken/application_outage/actions.py index bc66fb37..5a8f4476 100644 --- a/kraken/application_outage/actions.py +++ b/kraken/application_outage/actions.py @@ -4,13 +4,13 @@ import time import kraken.cerberus.setup as cerberus from jinja2 import Template import kraken.invoke.command as runcommand -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry from krkn_lib.utils.functions import get_yaml_item_value # Reads the scenario config, applies and deletes a network policy to # block the traffic for the specified duration -def run(scenarios_list, config, wait_duration, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list, config, wait_duration, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): failed_post_scenarios = "" scenario_telemetries: list[ScenarioTelemetry] = [] failed_scenarios = [] diff --git a/kraken/arcaflow_plugin/arcaflow_plugin.py b/kraken/arcaflow_plugin/arcaflow_plugin.py index 18799613..250006f1 100644 --- a/kraken/arcaflow_plugin/arcaflow_plugin.py +++ b/kraken/arcaflow_plugin/arcaflow_plugin.py @@ -6,11 +6,11 @@ import logging from pathlib import Path from typing import List from .context_auth import ContextAuth -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry -def run(scenarios_list: List[str], kubeconfig_path: str, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list: List[str], kubeconfig_path: str, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): scenario_telemetries: list[ScenarioTelemetry] = [] failed_post_scenarios = [] for scenario in scenarios_list: diff --git a/kraken/network_chaos/actions.py b/kraken/network_chaos/actions.py index 831bde0e..069dd0fb 100644 --- a/kraken/network_chaos/actions.py +++ b/kraken/network_chaos/actions.py @@ -7,14 +7,14 @@ import kraken.cerberus.setup as cerberus import kraken.node_actions.common_node_functions as common_node_functions from jinja2 import Environment, FileSystemLoader from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry from krkn_lib.utils.functions import get_yaml_item_value # krkn_lib # Reads the scenario config and introduces traffic variations in Node's host network interface. -def run(scenarios_list, config, wait_duration, kubecli: KrknKubernetes, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list, config, wait_duration, kubecli: KrknKubernetes, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): failed_post_scenarios = "" logging.info("Runing the Network Chaos tests") failed_post_scenarios = "" diff --git a/kraken/node_actions/run.py b/kraken/node_actions/run.py index c08ae3e2..ad382133 100644 --- a/kraken/node_actions/run.py +++ b/kraken/node_actions/run.py @@ -13,9 +13,8 @@ from kraken.node_actions.docker_node_scenarios import docker_node_scenarios import kraken.node_actions.common_node_functions as common_node_functions import kraken.cerberus.setup as cerberus from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry, ScenarioTelemetry -from krkn_lib.utils.functions import get_yaml_item_value - +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes +from krkn_lib.models.telemetry import ScenarioTelemetry node_general = False @@ -54,7 +53,7 @@ def get_node_scenario_object(node_scenario, kubecli: KrknKubernetes): # Run defined scenarios # krkn_lib -def run(scenarios_list, config, wait_duration, kubecli: KrknKubernetes, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list, config, wait_duration, kubecli: KrknKubernetes, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): scenario_telemetries: list[ScenarioTelemetry] = [] failed_scenarios = [] for node_scenario_config in scenarios_list: diff --git a/kraken/plugins/__init__.py b/kraken/plugins/__init__.py index e0109641..8f455966 100644 --- a/kraken/plugins/__init__.py +++ b/kraken/plugins/__init__.py @@ -13,7 +13,7 @@ from kraken.plugins.run_python_plugin import run_python_file from kraken.plugins.network.ingress_shaping import network_chaos from kraken.plugins.pod_network_outage.pod_network_outage_plugin import pod_outage from kraken.plugins.pod_network_outage.pod_network_outage_plugin import pod_egress_shaping -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry @@ -228,7 +228,7 @@ PLUGINS = Plugins( ) -def run(scenarios: List[str], kubeconfig_path: str, kraken_config: str, failed_post_scenarios: List[str], wait_duration: int, telemetry: KrknTelemetry) -> (List[str], list[ScenarioTelemetry]): +def run(scenarios: List[str], kubeconfig_path: str, kraken_config: str, failed_post_scenarios: List[str], wait_duration: int, telemetry: KrknTelemetryKubernetes) -> (List[str], list[ScenarioTelemetry]): scenario_telemetries: list[ScenarioTelemetry] = [] for scenario in scenarios: scenario_telemetry = ScenarioTelemetry() diff --git a/kraken/pod_scenarios/setup.py b/kraken/pod_scenarios/setup.py index 49315a04..f1ca1edc 100644 --- a/kraken/pod_scenarios/setup.py +++ b/kraken/pod_scenarios/setup.py @@ -7,7 +7,7 @@ import arcaflow_plugin_kill_pod import kraken.cerberus.setup as cerberus import kraken.post_actions.actions as post_actions from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry from arcaflow_plugin_sdk import serialization from krkn_lib.utils.functions import get_yaml_item_value @@ -75,7 +75,7 @@ def container_run(kubeconfig_path, failed_post_scenarios, wait_duration, kubecli: KrknKubernetes, - telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): + telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): failed_scenarios = [] scenario_telemetries: list[ScenarioTelemetry] = [] diff --git a/kraken/pvc/pvc_scenario.py b/kraken/pvc/pvc_scenario.py index c07929dc..ee92b682 100644 --- a/kraken/pvc/pvc_scenario.py +++ b/kraken/pvc/pvc_scenario.py @@ -5,13 +5,13 @@ import time import yaml from ..cerberus import setup as cerberus from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry from krkn_lib.utils.functions import get_yaml_item_value # krkn_lib -def run(scenarios_list, config, kubecli: KrknKubernetes, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list, config, kubecli: KrknKubernetes, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): """ Reads the scenario config and creates a temp file to fill up the PVC """ diff --git a/kraken/service_disruption/common_service_disruption_functions.py b/kraken/service_disruption/common_service_disruption_functions.py index e72dfa9d..84522982 100644 --- a/kraken/service_disruption/common_service_disruption_functions.py +++ b/kraken/service_disruption/common_service_disruption_functions.py @@ -5,7 +5,7 @@ import kraken.cerberus.setup as cerberus import kraken.post_actions.actions as post_actions import yaml from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry from krkn_lib.utils.functions import get_yaml_item_value @@ -158,7 +158,7 @@ def run( failed_post_scenarios, kubeconfig_path, kubecli: KrknKubernetes, - telemetry: KrknTelemetry + telemetry: KrknTelemetryKubernetes ) -> (list[str], list[ScenarioTelemetry]): scenario_telemetries: list[ScenarioTelemetry] = [] failed_scenarios = [] diff --git a/kraken/shut_down/common_shut_down_func.py b/kraken/shut_down/common_shut_down_func.py index 72e78609..d63df3e9 100644 --- a/kraken/shut_down/common_shut_down_func.py +++ b/kraken/shut_down/common_shut_down_func.py @@ -10,7 +10,7 @@ from ..node_actions.openstack_node_scenarios import OPENSTACKCLOUD from ..node_actions.az_node_scenarios import Azure from ..node_actions.gcp_node_scenarios import GCP from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry def multiprocess_nodes(cloud_object_function, nodes): @@ -128,7 +128,7 @@ def cluster_shut_down(shut_down_config, kubecli: KrknKubernetes): # krkn_lib -def run(scenarios_list, config, wait_duration, kubecli: KrknKubernetes, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list, config, wait_duration, kubecli: KrknKubernetes, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): failed_post_scenarios = [] failed_scenarios = [] scenario_telemetries: list[ScenarioTelemetry] = [] diff --git a/kraken/time_actions/common_time_functions.py b/kraken/time_actions/common_time_functions.py index a2688cae..480894c0 100644 --- a/kraken/time_actions/common_time_functions.py +++ b/kraken/time_actions/common_time_functions.py @@ -7,7 +7,7 @@ import random from ..cerberus import setup as cerberus from ..invoke import command as runcommand from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry from krkn_lib.utils.functions import get_yaml_item_value @@ -309,7 +309,7 @@ def check_date_time(object_type, names, kubecli:KrknKubernetes): # krkn_lib -def run(scenarios_list, config, wait_duration, kubecli:KrknKubernetes, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]): +def run(scenarios_list, config, wait_duration, kubecli:KrknKubernetes, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]): failed_scenarios = [] scenario_telemetries: list[ScenarioTelemetry] = [] for time_scenario_config in scenarios_list: diff --git a/kraken/zone_outage/actions.py b/kraken/zone_outage/actions.py index 6ed3de49..afbd194a 100644 --- a/kraken/zone_outage/actions.py +++ b/kraken/zone_outage/actions.py @@ -3,11 +3,11 @@ import logging import time from ..node_actions.aws_node_scenarios import AWS from ..cerberus import setup as cerberus -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes from krkn_lib.models.telemetry import ScenarioTelemetry -def run(scenarios_list, config, wait_duration, telemetry: KrknTelemetry) -> (list[str], list[ScenarioTelemetry]) : +def run(scenarios_list, config, wait_duration, telemetry: KrknTelemetryKubernetes) -> (list[str], list[ScenarioTelemetry]) : """ filters the subnet of interest and applies the network acl to create zone outage diff --git a/requirements.txt b/requirements.txt index a8f39a5d..f91b0509 100644 --- a/requirements.txt +++ b/requirements.txt @@ -19,7 +19,7 @@ ibm_cloud_sdk_core ibm_vpc itsdangerous==2.0.1 jinja2==3.0.3 -krkn-lib == 1.3.2 +krkn-lib>=1.4.0 kubernetes lxml >= 4.3.0 oauth2client>=4.1.3 diff --git a/run_kraken.py b/run_kraken.py index 512b796f..6bff7a73 100644 --- a/run_kraken.py +++ b/run_kraken.py @@ -27,7 +27,9 @@ import kraken.prometheus.client as promcli from kraken import plugins from krkn_lib.k8s import KrknKubernetes -from krkn_lib.telemetry import KrknTelemetry +from krkn_lib.ocp import KrknOpenshift +from krkn_lib.telemetry.k8s import KrknTelemetryKubernetes +from krkn_lib.telemetry.ocp import KrknTelemetryOpenshift from krkn_lib.models.telemetry import ChaosRunTelemetry from krkn_lib.utils import SafeLogger from krkn_lib.utils.functions import get_yaml_item_value @@ -153,11 +155,14 @@ def main(cfg): os.environ["KUBECONFIG"] = str(kubeconfig_path) # krkn-lib-kubernetes init kubecli = KrknKubernetes(kubeconfig_path=kubeconfig_path) + ocpcli = KrknOpenshift(kubeconfig_path=kubeconfig_path) except: kubecli.initialize_clients(None) # KrknTelemetry init - telemetry = KrknTelemetry(safe_logger, kubecli) + telemetry_k8s = KrknTelemetryKubernetes(safe_logger, kubecli) + telemetry_ocp = KrknTelemetryOpenshift(safe_logger, ocpcli) + # find node kraken might be running on kubecli.find_kraken_node() @@ -183,7 +188,9 @@ def main(cfg): # Cluster info logging.info("Fetching cluster info") - cv = kubecli.get_clusterversion_string() + cv = "" + if config["kraken"]["distribution"] == "openshift": + cv = ocpcli.get_clusterversion_string() if cv != "": logging.info(cv) else: @@ -252,7 +259,7 @@ def main(cfg): sys.exit(1) elif scenario_type == "arcaflow_scenarios": failed_post_scenarios, scenario_telemetries = arcaflow_plugin.run( - scenarios_list, kubeconfig_path, telemetry + scenarios_list, kubeconfig_path, telemetry_k8s ) chaos_telemetry.scenarios.extend(scenario_telemetries) @@ -263,7 +270,7 @@ def main(cfg): kraken_config, failed_post_scenarios, wait_duration, - telemetry + telemetry_k8s ) chaos_telemetry.scenarios.extend(scenario_telemetries) # krkn_lib @@ -276,7 +283,7 @@ def main(cfg): failed_post_scenarios, wait_duration, kubecli, - telemetry + telemetry_k8s ) chaos_telemetry.scenarios.extend(scenario_telemetries) @@ -284,7 +291,7 @@ def main(cfg): # krkn_lib elif scenario_type == "node_scenarios": logging.info("Running node scenarios") - failed_post_scenarios, scenario_telemetries = nodeaction.run(scenarios_list, config, wait_duration, kubecli, telemetry) + failed_post_scenarios, scenario_telemetries = nodeaction.run(scenarios_list, config, wait_duration, kubecli, telemetry_k8s) chaos_telemetry.scenarios.extend(scenario_telemetries) # Inject managedcluster chaos scenarios specified in the config # krkn_lib @@ -300,7 +307,7 @@ def main(cfg): elif scenario_type == "time_scenarios": if distribution == "openshift": logging.info("Running time skew scenarios") - failed_post_scenarios, scenario_telemetries = time_actions.run(scenarios_list, config, wait_duration, kubecli, telemetry) + failed_post_scenarios, scenario_telemetries = time_actions.run(scenarios_list, config, wait_duration, kubecli, telemetry_k8s) chaos_telemetry.scenarios.extend(scenario_telemetries) else: logging.error( @@ -351,7 +358,7 @@ def main(cfg): # Inject cluster shutdown scenarios # krkn_lib elif scenario_type == "cluster_shut_down_scenarios": - failed_post_scenarios, scenario_telemetries = shut_down.run(scenarios_list, config, wait_duration, kubecli, telemetry) + failed_post_scenarios, scenario_telemetries = shut_down.run(scenarios_list, config, wait_duration, kubecli, telemetry_k8s) chaos_telemetry.scenarios.extend(scenario_telemetries) # Inject namespace chaos scenarios @@ -365,34 +372,34 @@ def main(cfg): failed_post_scenarios, kubeconfig_path, kubecli, - telemetry + telemetry_k8s ) chaos_telemetry.scenarios.extend(scenario_telemetries) # Inject zone failures elif scenario_type == "zone_outages": logging.info("Inject zone outages") - failed_post_scenarios, scenario_telemetries = zone_outages.run(scenarios_list, config, wait_duration, telemetry) + failed_post_scenarios, scenario_telemetries = zone_outages.run(scenarios_list, config, wait_duration, telemetry_k8s) chaos_telemetry.scenarios.extend(scenario_telemetries) # Application outages elif scenario_type == "application_outages": logging.info("Injecting application outage") failed_post_scenarios, scenario_telemetries = application_outage.run( - scenarios_list, config, wait_duration, telemetry) + scenarios_list, config, wait_duration, telemetry_k8s) chaos_telemetry.scenarios.extend(scenario_telemetries) # PVC scenarios # krkn_lib elif scenario_type == "pvc_scenarios": logging.info("Running PVC scenario") - failed_post_scenarios, scenario_telemetries = pvc_scenario.run(scenarios_list, config, kubecli, telemetry) + failed_post_scenarios, scenario_telemetries = pvc_scenario.run(scenarios_list, config, kubecli, telemetry_k8s) chaos_telemetry.scenarios.extend(scenario_telemetries) # Network scenarios # krkn_lib elif scenario_type == "network_chaos": logging.info("Running Network Chaos") - failed_post_scenarios, scenario_telemetries = network_chaos.run(scenarios_list, config, wait_duration, kubecli, telemetry) + failed_post_scenarios, scenario_telemetries = network_chaos.run(scenarios_list, config, wait_duration, kubecli, telemetry_k8s) # Check for critical alerts when enabled if check_critical_alerts: @@ -416,21 +423,31 @@ def main(cfg): # is disabled, it's necessary to serialize the ChaosRunTelemetry object # to json, and recreate a new object from it. end_time = int(time.time()) - telemetry.collect_cluster_metadata(chaos_telemetry) + + # if platform is openshift will be collected + # Cloud platform and network plugins metadata + # through OCP specific APIs + if config["kraken"]["distribution"] == "openshift": + telemetry_ocp.collect_cluster_metadata(chaos_telemetry) + else: + telemetry_k8s.collect_cluster_metadata(chaos_telemetry) + decoded_chaos_run_telemetry = ChaosRunTelemetry(json.loads(chaos_telemetry.to_json())) logging.info(f"Telemetry data:\n{decoded_chaos_run_telemetry.to_json()}") + if config["telemetry"]["enabled"]: logging.info(f"telemetry data will be stored on s3 bucket folder: {telemetry_request_id}") logging.info(f"telemetry upload log: {safe_logger.log_file_name}") try: - telemetry.send_telemetry(config["telemetry"], telemetry_request_id, chaos_telemetry) - if config["telemetry"]["prometheus_backup"]: + telemetry_k8s.send_telemetry(config["telemetry"], telemetry_request_id, chaos_telemetry) + # prometheus data collection is available only on Openshift + if config["telemetry"]["prometheus_backup"] and config["kraken"]["distribution"] == "openshift": safe_logger.info("archives download started:") - prometheus_archive_files = telemetry.get_ocp_prometheus_data(config["telemetry"], telemetry_request_id) + prometheus_archive_files = telemetry_ocp.get_ocp_prometheus_data(config["telemetry"], telemetry_request_id) safe_logger.info("archives upload started:") - telemetry.put_ocp_prometheus_data(config["telemetry"], prometheus_archive_files, telemetry_request_id) + telemetry_k8s.put_prometheus_data(config["telemetry"], prometheus_archive_files, telemetry_request_id) if config["telemetry"]["logs_backup"]: - telemetry.put_ocp_logs(telemetry_request_id, config["telemetry"], start_time, end_time) + telemetry_ocp.put_ocp_logs(telemetry_request_id, config["telemetry"], start_time, end_time) except Exception as e: logging.error(f"failed to send telemetry data: {str(e)}") else: