From 2624102d65ec602137dcb48c1ac1e73c78610b96 Mon Sep 17 00:00:00 2001 From: Tullio Sebastiani Date: Thu, 10 Apr 2025 15:47:29 +0100 Subject: [PATCH] Node Network Filtering Scenario + Network Chaos NG modular architecture (#766) * network chaos NG modular architecture error handling * first working version (missing protocols, number of instances, wait duration) * added instance_count + sleep + methods documentation Signed-off-by: Tullio Sebastiani --------- Signed-off-by: Tullio Sebastiani Co-authored-by: Naga Ravi Chaitanya Elluri --- .../network_chaos_ng/__init__.py | 0 .../network_chaos_ng/models.py | 41 ++++++ .../network_chaos_ng/modules/__init__.py | 0 .../modules/abstract_network_chaos_module.py | 58 ++++++++ .../modules/node_network_filter.py | 136 ++++++++++++++++++ .../modules/templates/network-chaos.j2 | 17 +++ .../network_chaos_ng/network_chaos_factory.py | 24 ++++ .../network_chaos_ng_scenario_plugin.py | 116 +++++++++++++++ scenarios/kube/network_filter.yml | 13 ++ 9 files changed, 405 insertions(+) create mode 100644 krkn/scenario_plugins/network_chaos_ng/__init__.py create mode 100644 krkn/scenario_plugins/network_chaos_ng/models.py create mode 100644 krkn/scenario_plugins/network_chaos_ng/modules/__init__.py create mode 100644 krkn/scenario_plugins/network_chaos_ng/modules/abstract_network_chaos_module.py create mode 100644 krkn/scenario_plugins/network_chaos_ng/modules/node_network_filter.py create mode 100644 krkn/scenario_plugins/network_chaos_ng/modules/templates/network-chaos.j2 create mode 100644 krkn/scenario_plugins/network_chaos_ng/network_chaos_factory.py create mode 100644 krkn/scenario_plugins/network_chaos_ng/network_chaos_ng_scenario_plugin.py create mode 100644 scenarios/kube/network_filter.yml diff --git a/krkn/scenario_plugins/network_chaos_ng/__init__.py b/krkn/scenario_plugins/network_chaos_ng/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/krkn/scenario_plugins/network_chaos_ng/models.py b/krkn/scenario_plugins/network_chaos_ng/models.py new file mode 100644 index 00000000..e8efbcfe --- /dev/null +++ b/krkn/scenario_plugins/network_chaos_ng/models.py @@ -0,0 +1,41 @@ +from dataclasses import dataclass +from enum import Enum + + +class NetworkChaosScenarioType(Enum): + Node = 1 + Pod = 2 + +@dataclass +class BaseNetworkChaosConfig: + supported_execution = ["serial", "parallel"] + id: str + wait_duration: int + test_duration: int + label_selector: str + instance_count: int + execution: str + namespace: str + + def validate(self) -> list[str]: + errors = [] + if self.execution is None: + errors.append(f"execution cannot be None, supported values are: {','.join(self.supported_execution)}") + if self.execution not in self.supported_execution: + errors.append(f"{self.execution} is not in supported execution mod: {','.join(self.supported_execution)}") + if self.label_selector is None: + errors.append("label_selector cannot be None") + return errors + +@dataclass +class NetworkFilterConfig(BaseNetworkChaosConfig): + ingress: bool + egress: bool + interfaces: list[str] + target: str + ports: list[int] + + def validate(self) -> list[str]: + errors = super().validate() + # here further validations + return errors diff --git a/krkn/scenario_plugins/network_chaos_ng/modules/__init__.py b/krkn/scenario_plugins/network_chaos_ng/modules/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/krkn/scenario_plugins/network_chaos_ng/modules/abstract_network_chaos_module.py b/krkn/scenario_plugins/network_chaos_ng/modules/abstract_network_chaos_module.py new file mode 100644 index 00000000..072d9afe --- /dev/null +++ b/krkn/scenario_plugins/network_chaos_ng/modules/abstract_network_chaos_module.py @@ -0,0 +1,58 @@ +import abc +import logging +import queue + +from krkn_lib.telemetry.ocp import KrknTelemetryOpenshift +from krkn.scenario_plugins.network_chaos_ng.models import BaseNetworkChaosConfig, NetworkChaosScenarioType + + +class AbstractNetworkChaosModule(abc.ABC): + """ + The abstract class that needs to be implemented by each Network Chaos Scenario + """ + @abc.abstractmethod + def run(self, target: str, kubecli: KrknTelemetryOpenshift, error_queue: queue.Queue = None): + """ + the entrypoint method for the Network Chaos Scenario + :param target: The resource name that will be targeted by the scenario (Node Name, Pod Name etc.) + :param kubecli: The `KrknTelemetryOpenshift` needed by the scenario to access to the krkn-lib methods + :param error_queue: A queue that will be used by the plugin to push the errors raised during the execution of parallel modules + """ + pass + + @abc.abstractmethod + def get_config(self) -> (NetworkChaosScenarioType, BaseNetworkChaosConfig): + """ + returns the common subset of settings shared by all the scenarios `BaseNetworkChaosConfig` and the type of Network + Chaos Scenario that is running (Pod Scenario or Node Scenario) + """ + pass + + + def log_info(self, message: str, parallel: bool = False, node_name: str = ""): + """ + log helper method for INFO severity to be used in the scenarios + """ + if parallel: + logging.info(f"[{node_name}]: {message}") + else: + logging.info(message) + + def log_warning(self, message: str, parallel: bool = False, node_name: str = ""): + """ + log helper method for WARNING severity to be used in the scenarios + """ + if parallel: + logging.warning(f"[{node_name}]: {message}") + else: + logging.warning(message) + + + def log_error(self, message: str, parallel: bool = False, node_name: str = ""): + """ + log helper method for ERROR severity to be used in the scenarios + """ + if parallel: + logging.error(f"[{node_name}]: {message}") + else: + logging.error(message) \ No newline at end of file diff --git a/krkn/scenario_plugins/network_chaos_ng/modules/node_network_filter.py b/krkn/scenario_plugins/network_chaos_ng/modules/node_network_filter.py new file mode 100644 index 00000000..96d1b399 --- /dev/null +++ b/krkn/scenario_plugins/network_chaos_ng/modules/node_network_filter.py @@ -0,0 +1,136 @@ +import os +import queue +import time + +import yaml +from jinja2 import Environment, FileSystemLoader + + +from krkn_lib.telemetry.ocp import KrknTelemetryOpenshift +from krkn_lib.utils import get_random_string +from krkn.scenario_plugins.network_chaos_ng.models import ( + BaseNetworkChaosConfig, + NetworkFilterConfig, + NetworkChaosScenarioType, +) +from krkn.scenario_plugins.network_chaos_ng.modules.abstract_network_chaos_module import ( + AbstractNetworkChaosModule, +) + + +class NodeNetworkFilterModule(AbstractNetworkChaosModule): + config: NetworkFilterConfig + + def run( + self, + target: str, + kubecli: KrknTelemetryOpenshift, + error_queue: queue.Queue = None, + ): + parallel = False + if error_queue: + parallel = True + try: + file_loader = FileSystemLoader(os.path.abspath(os.path.dirname(__file__))) + env = Environment(loader=file_loader, autoescape=True) + pod_name = f"node-filter-{get_random_string(5)}" + pod_template = env.get_template("templates/network-chaos.j2") + pod_body = yaml.safe_load( + pod_template.render( + pod_name=pod_name, + namespace=self.config.namespace, + host_network=True, + target=target, + ) + ) + self.log_info( + f"creating pod to filter " + f"ports {','.join([str(port) for port in self.config.ports])}, " + f"ingress:{str(self.config.ingress)}, " + f"egress:{str(self.config.egress)}", + parallel, + target, + ) + kubecli.get_lib_kubernetes().create_pod( + pod_body, self.config.namespace, 300 + ) + + if len(self.config.interfaces) == 0: + interfaces = [ + self.get_default_interface(pod_name, self.config.namespace, kubecli) + ] + self.log_info(f"detected default interface {interfaces[0]}") + else: + interfaces = self.config.interfaces + + input_rules, output_rules = self.generate_rules(interfaces) + + for rule in input_rules: + self.log_info(f"applying iptables INPUT rule: {rule}", parallel, target) + kubecli.get_lib_kubernetes().exec_cmd_in_pod( + [rule], pod_name, self.config.namespace + ) + for rule in output_rules: + self.log_info( + f"applying iptables OUTPUT rule: {rule}", parallel, target + ) + kubecli.get_lib_kubernetes().exec_cmd_in_pod( + [rule], pod_name, self.config.namespace + ) + self.log_info( + f"waiting {self.config.test_duration} seconds before removing the iptables rules" + ) + time.sleep(self.config.test_duration) + self.log_info("removing iptables rules") + for _ in input_rules: + # always deleting the first rule since has been inserted from the top + kubecli.get_lib_kubernetes().exec_cmd_in_pod( + [f"iptables -D INPUT 1"], pod_name, self.config.namespace + ) + for _ in output_rules: + # always deleting the first rule since has been inserted from the top + kubecli.get_lib_kubernetes().exec_cmd_in_pod( + [f"iptables -D OUTPUT 1"], pod_name, self.config.namespace + ) + self.log_info( + f"deleting network chaos pod {pod_name} from {self.config.namespace}" + ) + + kubecli.get_lib_kubernetes().delete_pod(pod_name, self.config.namespace) + + except Exception as e: + if error_queue is None: + raise e + else: + error_queue.put(str(e)) + + def __init__(self, config: NetworkFilterConfig): + self.config = config + + def get_config(self) -> (NetworkChaosScenarioType, BaseNetworkChaosConfig): + return NetworkChaosScenarioType.Node, self.config + + def get_default_interface( + self, pod_name: str, namespace: str, kubecli: KrknTelemetryOpenshift + ) -> str: + cmd = "ip r | grep default | awk '/default/ {print $5}'" + output = kubecli.get_lib_kubernetes().exec_cmd_in_pod( + [cmd], pod_name, namespace + ) + return output.replace("\n", "") + + def generate_rules(self, interfaces: list[str]) -> (list[str], list[str]): + input_rules = [] + output_rules = [] + for interface in interfaces: + for port in self.config.ports: + if self.config.egress: + output_rules.append( + f"iptables -I OUTPUT 1 -p tcp --dport {port} -m state --state NEW,RELATED,ESTABLISHED -j DROP" + ) + + if self.config.ingress: + input_rules.append( + f"iptables -I INPUT 1 -i {interface} -p tcp --dport {port} -m state --state NEW,RELATED,ESTABLISHED -j DROP" + ) + return input_rules, output_rules diff --git a/krkn/scenario_plugins/network_chaos_ng/modules/templates/network-chaos.j2 b/krkn/scenario_plugins/network_chaos_ng/modules/templates/network-chaos.j2 new file mode 100644 index 00000000..5d14b098 --- /dev/null +++ b/krkn/scenario_plugins/network_chaos_ng/modules/templates/network-chaos.j2 @@ -0,0 +1,17 @@ +apiVersion: v1 +kind: Pod +metadata: + name: {{pod_name}} + namespace: {{namespace}} +spec: + {% if host_network %} + hostNetwork: true + {%endif%} + nodeSelector: + kubernetes.io/hostname: {{target}} + containers: + - name: fedora + imagePullPolicy: Always + image: quay.io/krkn-chaos/krkn-network-chaos:latest + securityContext: + privileged: true diff --git a/krkn/scenario_plugins/network_chaos_ng/network_chaos_factory.py b/krkn/scenario_plugins/network_chaos_ng/network_chaos_factory.py new file mode 100644 index 00000000..9e301fd6 --- /dev/null +++ b/krkn/scenario_plugins/network_chaos_ng/network_chaos_factory.py @@ -0,0 +1,24 @@ +from krkn.scenario_plugins.network_chaos_ng.models import NetworkFilterConfig +from krkn.scenario_plugins.network_chaos_ng.modules.abstract_network_chaos_module import AbstractNetworkChaosModule +from krkn.scenario_plugins.network_chaos_ng.modules.node_network_filter import NodeNetworkFilterModule + + +supported_modules = ["node_network_filter"] + +class NetworkChaosFactory: + + @staticmethod + def get_instance(config: dict[str, str]) -> AbstractNetworkChaosModule: + if config["id"] is None: + raise Exception("network chaos id cannot be None") + if config["id"] not in supported_modules: + raise Exception(f"{config['id']} is not a supported network chaos module") + + if config["id"] == "node_network_filter": + config = NetworkFilterConfig(**config) + errors = config.validate() + if len(errors) > 0: + raise Exception(f"config validation errors: [{';'.join(errors)}]") + return NodeNetworkFilterModule(config) + + diff --git a/krkn/scenario_plugins/network_chaos_ng/network_chaos_ng_scenario_plugin.py b/krkn/scenario_plugins/network_chaos_ng/network_chaos_ng_scenario_plugin.py new file mode 100644 index 00000000..2da27784 --- /dev/null +++ b/krkn/scenario_plugins/network_chaos_ng/network_chaos_ng_scenario_plugin.py @@ -0,0 +1,116 @@ +import logging +import queue +import random +import threading +import time + +import yaml +from krkn_lib.models.telemetry import ScenarioTelemetry +from krkn_lib.telemetry.ocp import KrknTelemetryOpenshift + +from krkn.scenario_plugins.abstract_scenario_plugin import AbstractScenarioPlugin +from krkn.scenario_plugins.network_chaos_ng.models import ( + NetworkChaosScenarioType, + BaseNetworkChaosConfig, +) +from krkn.scenario_plugins.network_chaos_ng.modules.abstract_network_chaos_module import ( + AbstractNetworkChaosModule, +) +from krkn.scenario_plugins.network_chaos_ng.network_chaos_factory import ( + NetworkChaosFactory, +) + + +class NetworkChaosNgScenarioPlugin(AbstractScenarioPlugin): + def run( + self, + run_uuid: str, + scenario: str, + krkn_config: dict[str, any], + lib_telemetry: KrknTelemetryOpenshift, + scenario_telemetry: ScenarioTelemetry, + ) -> int: + try: + with open(scenario, "r") as file: + scenario_config = yaml.safe_load(file) + if not isinstance(scenario_config, list): + logging.error( + "network chaos scenario config must be a list of objects" + ) + return 1 + for config in scenario_config: + network_chaos = NetworkChaosFactory.get_instance(config) + network_chaos_config = network_chaos.get_config() + logging.info( + f"running network_chaos scenario: {network_chaos_config[1].id}" + ) + if network_chaos_config[0] == NetworkChaosScenarioType.Node: + targets = lib_telemetry.get_lib_kubernetes().list_nodes( + network_chaos_config[1].label_selector + ) + else: + targets = lib_telemetry.get_lib_kubernetes().list_pods( + network_chaos_config[1].namespace, + network_chaos_config[1].label_selector, + ) + if len(targets) == 0: + logging.warning( + f"no targets found for {network_chaos_config[1].id} " + f"network chaos scenario with selector {network_chaos_config[1].label_selector} " + f"with target type {network_chaos_config[0]}" + ) + + if network_chaos_config[1].instance_count != 0 and network_chaos_config[1].instance_count > len(targets): + targets = random.sample(targets, network_chaos_config[1].instance_count) + + if network_chaos_config[1].execution == "parallel": + self.run_parallel(targets, network_chaos, lib_telemetry) + else: + self.run_serial(targets, network_chaos, lib_telemetry) + if len(config) > 1: + logging.info(f"waiting {network_chaos_config[1].wait_duration} seconds before running the next " + f"Network Chaos NG Module") + time.sleep(network_chaos_config[1].wait_duration) + except Exception as e: + logging.error(str(e)) + return 1 + return 0 + + def run_parallel( + self, + targets: list[str], + module: AbstractNetworkChaosModule, + lib_telemetry: KrknTelemetryOpenshift, + ): + error_queue = queue.Queue() + threads = [] + errors = [] + for target in targets: + thread = threading.Thread( + target=module.run, args=[target, lib_telemetry, error_queue] + ) + thread.start() + threads.append(thread) + for thread in threads: + thread.join() + while True: + try: + errors.append(error_queue.get_nowait()) + except queue.Empty: + break + if len(errors) > 0: + raise Exception( + f"module {module.get_config()[1].id} execution failed: [{';'.join(errors)}]" + ) + + def run_serial( + self, + targets: list[str], + module: AbstractNetworkChaosModule, + lib_telemetry: KrknTelemetryOpenshift, + ): + for target in targets: + module.run(target, lib_telemetry) + + def get_scenario_types(self) -> list[str]: + return ["network_chaos_ng_scenarios"] diff --git a/scenarios/kube/network_filter.yml b/scenarios/kube/network_filter.yml new file mode 100644 index 00000000..68259e24 --- /dev/null +++ b/scenarios/kube/network_filter.yml @@ -0,0 +1,13 @@ +- id: node_network_filter + wait_duration: 300 + test_duration: 100 + label_selector: "kubernetes.io/hostname=ip-10-0-39-182.us-east-2.compute.internal" + namespace: 'default' + instance_count: 1 + execution: parallel + ingress: false + egress: true + target: node + interfaces: [] + ports: + - 2049 \ No newline at end of file