diff options
| author | Vitalii Gridnev <[email protected]> | 2024-08-16 20:03:11 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-08-16 17:03:11 +0000 |
| commit | 28512db81ae560484d090e6a6796016e86e9c8a9 (patch) | |
| tree | a9150f8bfc37a285b25fdb2d0c716c93ed4ecf40 | |
| parent | 89826f07428fdf325c2b0f79da0deb09c3628706 (diff) | |
nemesis driver and ut (#7937)
| -rw-r--r-- | ydb/tests/tools/nemesis/driver/__main__.py | 162 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/driver/ya.make | 14 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/__init__.py | 0 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/base.py | 41 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/catalog.py | 101 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/disk.py | 129 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/monitor.py | 48 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/node.py | 275 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/tablet.py | 256 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/library/ya.make | 22 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/ut/__init__.py | 0 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/ut/test_disk.py | 42 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/ut/test_tablet.py | 74 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/ut/ya.make | 26 | ||||
| -rw-r--r-- | ydb/tests/tools/nemesis/ya.make | 5 | ||||
| -rw-r--r-- | ydb/tests/tools/ya.make | 1 |
16 files changed, 1196 insertions, 0 deletions
diff --git a/ydb/tests/tools/nemesis/driver/__main__.py b/ydb/tests/tools/nemesis/driver/__main__.py new file mode 100644 index 00000000000..e7e9093a091 --- /dev/null +++ b/ydb/tests/tools/nemesis/driver/__main__.py @@ -0,0 +1,162 @@ +# -*- coding: utf-8 -*- +import argparse +import logging.config +import subprocess as sp +import os +import tempfile + +import logging + +from ydb.tests.tools.nemesis.library import monitor +from ydb.tests.tools.nemesis.library import catalog +from ydb.tests.library.harness.kikimr_cluster import ExternalKiKiMRCluster + + +def setup_logging_config(filename=None): + handler = {'class': 'logging.StreamHandler', 'level': 'DEBUG', 'formatter': 'base'} + if filename: + handler = { + 'class': 'logging.handlers.TimedRotatingFileHandler', + 'filename': filename, 'when': 'midnight', 'level': 'DEBUG', 'formatter': 'base' + } + return { + 'version': 1, + 'formatters': { + 'base': { + 'format': '%(asctime)s - %(name)s - %(levelname)s - %(message)s', + }, + }, + 'handlers': { + 'handler': handler, + }, + 'root': { + 'level': 'DEBUG', + 'handlers': ( + 'handler', + ) + }, + 'ydb.tests.library.harness.kikimr_runner': { + 'level': 'DEBUG', + 'handlers': ( + 'handler', + ) + } + } + + +logger = logging.getLogger(__name__) + + +class SshAgent(object): + def __init__(self): + self._env = {} + self._env_backup = {} + self._keys = {} + self.start() + + @property + def pid(self): + return int(self._env["SSH_AGENT_PID"]) + + def start(self): + self._env_backup["SSH_AUTH_SOCK"] = os.environ.get("SSH_AUTH_SOCK") + self._env_backup["SSH_OPTIONS"] = os.environ.get("SSH_OPTIONS") + + for line in self._run(["ssh-agent"]).splitlines(): + name, _, value = line.decode('utf-8').partition("=") + if _ == "=": + value = value.split(";", 1)[0] + self._env[name] = value + os.environ[name] = value + + os.environ["SSH_OPTIONS"] = "{}UserKnownHostsFile=/dev/null,StrictHostKeyChecking=no".format( + "," + os.environ["SSH_OPTIONS"] if os.environ.get("SSH_OPTIONS") else "" + ) + + def stop(self): + self._run(['kill', '-9', str(self.pid)]) + + def add(self, key): + key_pub = self._key_pub(key) + self._run(["ssh-add", "-"], stdin=key) + return key_pub + + def remove(self, key_pub): + with tempfile.NamedTemporaryFile() as f: + f.write(key_pub) + f.flush() + self._run(["ssh-add", "-d", f.name]) + + def _key_pub(self, key): + with tempfile.NamedTemporaryFile() as f: + f.write(key) + f.flush() + return self._run(["ssh-keygen", "-y", "-f", f.name]) + + @staticmethod + def _run(cmd, stdin=None): + p = sp.Popen(cmd, stdout=sp.PIPE, stderr=sp.PIPE, stdin=sp.PIPE if stdin else None) + stdout, stderr = p.communicate(stdin) + + # Listing keys from empty ssh-agent results in exit code 1 + if stdout.strip() == "The agent has no identities.": + return "" + + if p.returncode: + message = stderr.strip() + "\n" + stdout.strip() + raise RuntimeError(message.strip()) + + return stdout + + +class Key(object): + def __init__(self, key_file): + self.key_file = key_file + with open(key_file) as fd: + self.key = fd.read() + self._key_pub = None + self._ssh_agent = SshAgent() + + def __enter__(self): + self._key_pub = self._ssh_agent.add(self.key.encode('utf-8')) + + def __exit__(self, exc_type, exc_val, exc_tb): + self._ssh_agent.remove(self._key_pub) + self._ssh_agent.stop() + + +def nemesis_logic(arguments): + logging.config.dictConfig(setup_logging_config(arguments.log_file)) + nemesis = catalog.nemesis_factory( + ExternalKiKiMRCluster( + arguments.ydb_cluster_template, + binary_path=arguments.ydb_binary_path, + output_path=tempfile.gettempdir(), + ), + enable_nemesis_list_filter_by_hostname=arguments.enable_nemesis_list_filter_by_hostname, + ) + nemesis.start() + monitor.setup_page(arguments.mon_host, arguments.mon_port) + nemesis.stop() + + +def main(): + parser = argparse.ArgumentParser(formatter_class=argparse.RawDescriptionHelpFormatter) + parser.add_argument('--ydb-cluster-template', required=True, help='Path to the Yandex DB cluster template') + parser.add_argument('--ydb-binary-path', required=True, help='Path to the Yandex DB binary') + parser.add_argument('--private-key-file', default='') + parser.add_argument('--log-file', default=None) + parser.add_argument('--mon-port', default=8666, type=lambda x: int(x)) + parser.add_argument('--mon-host', default='::', type=lambda x: str(x)) + parser.add_argument('--enable-nemesis-list-filter-by-hostname', action='store_true') + arguments = parser.parse_args() + + if arguments.private_key_file: + with Key(arguments.private_key_file): + nemesis_logic(arguments) + else: + nemesis_logic(arguments) + + +if __name__ == '__main__': + main() diff --git a/ydb/tests/tools/nemesis/driver/ya.make b/ydb/tests/tools/nemesis/driver/ya.make new file mode 100644 index 00000000000..585ac99e95e --- /dev/null +++ b/ydb/tests/tools/nemesis/driver/ya.make @@ -0,0 +1,14 @@ +SUBSCRIBER(g:kikimr) +PY3_PROGRAM(nemesis) + +PY_SRCS( + __main__.py +) + +PEERDIR( + ydb/tests/library + ydb/tests/tools/nemesis/library + ydb/tools/cfg +) + +END() diff --git a/ydb/tests/tools/nemesis/library/__init__.py b/ydb/tests/tools/nemesis/library/__init__.py new file mode 100644 index 00000000000..e69de29bb2d --- /dev/null +++ b/ydb/tests/tools/nemesis/library/__init__.py diff --git a/ydb/tests/tools/nemesis/library/base.py b/ydb/tests/tools/nemesis/library/base.py new file mode 100644 index 00000000000..d9e9955d5bd --- /dev/null +++ b/ydb/tests/tools/nemesis/library/base.py @@ -0,0 +1,41 @@ +# -*- coding: utf-8 -*- +import abc + +import six +from ydb.tests.tools.nemesis.library import monitor + + [email protected]_metaclass(abc.ABCMeta) +class AbstractMonitoredNemesis(object): + def __init__(self, scope=None): + self.inject_completed = None + self.inject_in_flight = None + self.inject_in_flight_value = 0 + self.extract_completed = None + self.registry = monitor.monitor() + self.register_counters(scope) + + @property + def name(self): + return self.__class__.__name__ + + def register_counters(self, scope=None): + labels = {'nemesis': self.name} + if scope is not None: + labels.update({'scope': scope}) + self.inject_completed = self.registry.rate('InjectCompleted', labels) + self.inject_in_flight = self.registry.int_gauge('InjectInFlight', labels) + self.extract_completed = self.registry.rate('ExtractCompleted', labels) + + def start_inject_fault(self): + self.inject_in_flight_value += 1 + self.inject_in_flight.set(self.inject_in_flight_value) + + def on_success_extract_fault(self): + self.extract_completed.inc() + + def on_success_inject_fault(self): + if self.inject_in_flight_value > 0: + self.inject_in_flight_value -= 1 + self.inject_in_flight.set(self.inject_in_flight_value) + self.inject_completed.inc() diff --git a/ydb/tests/tools/nemesis/library/catalog.py b/ydb/tests/tools/nemesis/library/catalog.py new file mode 100644 index 00000000000..1ac228f72cc --- /dev/null +++ b/ydb/tests/tools/nemesis/library/catalog.py @@ -0,0 +1,101 @@ +# -*- coding: utf-8 -*- +import socket + +from ydb.tests.library.nemesis.nemesis_core import NemesisProcess +from ydb.tests.library.nemesis.nemesis_network import NetworkNemesis + +from ydb.tests.library.harness import param_constants + +from ydb.tests.tools.nemesis.library.node import nodes_nemesis_list + +from ydb.tests.tools.nemesis.library.tablet import change_tablet_group_nemesis_list +from ydb.tests.tools.nemesis.library.tablet import ReBalanceTabletsNemesis +from ydb.tests.tools.nemesis.library.tablet import KillTenantSlotBrokerNemesis +from ydb.tests.tools.nemesis.library.tablet import KillPersQueueNemesis +from ydb.tests.tools.nemesis.library.tablet import KickTabletsFromNode +from ydb.tests.tools.nemesis.library.tablet import KillKeyValueNemesis +from ydb.tests.tools.nemesis.library.tablet import KillHiveNemesis +from ydb.tests.tools.nemesis.library.tablet import KillBsControllerNemesis +from ydb.tests.tools.nemesis.library.tablet import KillCoordinatorNemesis +from ydb.tests.tools.nemesis.library.tablet import KillSchemeShardNemesis +from ydb.tests.tools.nemesis.library.tablet import KillMediatorNemesis +from ydb.tests.tools.nemesis.library.tablet import KillDataShardNemesis +from ydb.tests.tools.nemesis.library.tablet import KillTxAllocatorNemesis +from ydb.tests.tools.nemesis.library.tablet import KillNodeBrokerNemesis +from ydb.tests.tools.nemesis.library.tablet import KillBlocktoreVolume +from ydb.tests.tools.nemesis.library.tablet import KillBlocktorePartition +from ydb.tests.tools.nemesis.library.disk import data_storage_nemesis_list + + +def is_first_cluster_node(cluster): + if len(cluster.hostnames) > 0: + return cluster.hostnames[0] == socket.gethostname().strip() + return False + + +def basic_kikimr_nemesis_list( + cluster, num_of_pq_nemesis=10, network_nemesis=False, + enable_nemesis_list_filter_by_hostname=False): + harmful_nemesis_list = [] + harmful_nemesis_list.extend(data_storage_nemesis_list(cluster)) + harmful_nemesis_list.extend(nodes_nemesis_list(cluster)) + harmful_nemesis_list.extend( + [ + KickTabletsFromNode(cluster), + ReBalanceTabletsNemesis(cluster), + ] + ) + + if network_nemesis: + harmful_nemesis_list.append( + NetworkNemesis( + cluster, + ssh_username=param_constants.ssh_username + ) + ) + + light_nemesis_list = [] + light_nemesis_list.extend([ + KillCoordinatorNemesis(cluster), + KillHiveNemesis(cluster), + KillBsControllerNemesis(cluster), + KillNodeBrokerNemesis(cluster), + KillSchemeShardNemesis(cluster), + KillMediatorNemesis(cluster), + KillTxAllocatorNemesis(cluster), + KillKeyValueNemesis(cluster), + KillTenantSlotBrokerNemesis(cluster), + ]) + + light_nemesis_list.extend(change_tablet_group_nemesis_list(cluster)) + light_nemesis_list.extend([KillPersQueueNemesis(cluster) for _ in range(num_of_pq_nemesis)]) + light_nemesis_list.extend([KillDataShardNemesis(cluster) for _ in range(num_of_pq_nemesis)]) + light_nemesis_list.extend([KillBlocktoreVolume(cluster) for _ in range(num_of_pq_nemesis)]) + light_nemesis_list.extend([KillBlocktorePartition(cluster) for _ in range(num_of_pq_nemesis)]) + + nemesis_list = [] + if enable_nemesis_list_filter_by_hostname: + hostnames = cluster.hostnames + self_hostname = socket.gethostname() + self_id = None + + for host_id, hostname in enumerate(hostnames): + if self_hostname == hostname: + self_id = host_id + + for nemesis_actor_id, nemesis_actor in enumerate(light_nemesis_list): + if self_id is not None and nemesis_actor_id % len(hostnames) == self_id: + nemesis_list.append(nemesis_actor) + + if is_first_cluster_node(cluster): + nemesis_list.extend( + harmful_nemesis_list) + + return nemesis_list + nemesis_list.extend(light_nemesis_list) + nemesis_list.extend(harmful_nemesis_list) + return nemesis_list + + +def nemesis_factory(kikimr_cluster, num_of_pq_nemesis=10, **kwargs): + return NemesisProcess(basic_kikimr_nemesis_list(kikimr_cluster, num_of_pq_nemesis, **kwargs)) diff --git a/ydb/tests/tools/nemesis/library/disk.py b/ydb/tests/tools/nemesis/library/disk.py new file mode 100644 index 00000000000..c8e6f8f3985 --- /dev/null +++ b/ydb/tests/tools/nemesis/library/disk.py @@ -0,0 +1,129 @@ +# -*- coding: utf-8 -*- +import abc +import random +import subprocess +import time + +from ydb.tests.library.nemesis.nemesis_core import Nemesis +from ydb.tests.library.predicates import blobstorage +from ydb.tests.tools.nemesis.library.base import AbstractMonitoredNemesis +from ydb.tests.library.common.msgbus_types import EDriveStatus +from ydb.core.protos.blobstorage_config_pb2 import TConfigResponse + + +class AbstractSafeEraseDataOnDisk(Nemesis, AbstractMonitoredNemesis): + def __init__(self, cluster, schedule=(90, 150)): + super(AbstractSafeEraseDataOnDisk, self).__init__(schedule) + AbstractMonitoredNemesis.__init__(self, 'disk') + self._cluster = cluster + self._successful_data_erase_count = 0 + + def all_vdisks_are_replicated(self): + return blobstorage.cluster_has_no_unreplicated_vdisks(self._cluster, 1, self.logger) + + @property + def state_ready(self): + if len(self._cluster.nodes.values()) < 8: + self.logger.info("Data erase is prohibited.") + return False + return self.all_vdisks_are_replicated() + + def extract_fault(self): + pass + + def inject_fault(self): + if self.state_ready: + self._erase_data() + self.logger.debug("Successful data erase") + else: + self.prepare_state() + self.logger.debug("Do not erase data this time") + return False + + def prepare_state(self): + pass + + @property + def successful_data_erase_count(self): + return self._successful_data_erase_count + + @abc.abstractmethod + def _erase_data(self): + pass + + +class SafelyCleanupDisks(AbstractSafeEraseDataOnDisk): + def __init__(self, cluster, schedule=(90, 150)): + super(SafelyCleanupDisks, self).__init__(cluster, schedule) + self._node_ids = self._cluster.nodes.keys() + + def _erase_data(self): + try: + node_id = random.choice(list(self._node_ids)) + self._cluster.nodes[node_id].kill_process_and_daemon() + self._cluster.nodes[node_id].cleanup_disks() + self._cluster.nodes[node_id].start() + self._successful_data_erase_count += 1 + self.on_success_inject_fault() + return True + except subprocess.CalledProcessError as e: + self.logger.error("Failed to cleanup disks, %s", str(e)) + return False + + +class SafelyBreakDisk(AbstractSafeEraseDataOnDisk): + def __init__(self, cluster, schedule=(90, 150)): + super(SafelyBreakDisk, self).__init__(cluster, schedule) + self._currently_broken_drive = None + self._states = {} + self._broken_drives = set() + + def prepare_state(self): + super(SafelyBreakDisk, self).prepare_state() + self._broken_drives = set() + for node_id, node in self._cluster.nodes.items(): + drives = self._cluster.client.read_drive_status(node.host, node.ic_port).BlobStorageConfigResponse + self._states[node_id] = {} + for status in drives.Status: + for drive in status.DriveStatus: + self._states[node_id][drive.Path] = drive.Status + if drive.Status != EDriveStatus.ACTIVE: + self._broken_drives.add( + (node_id, drive.Path)) + + def _change_drive_status(self, node_id, path, status): + host, ic_port = self._cluster.nodes[node_id].host, self._cluster.nodes[node_id].ic_port + self.logger.info("Change drive status host %s:%s, path %s, status %s", host, ic_port, path, status.name) + for _ in range(6): + response = self._cluster.client.update_drive_status(host, ic_port, path, status).BlobStorageConfigResponse + if not response.Success and len(response.Status) == 1 and response.Status[0].FailReason == \ + TConfigResponse.TStatus.EFailReason.kMayLoseData: + time.sleep(10) + else: + break + self.logger.info("Received response from controller %s" % response) + return response.Success + + def _erase_data(self): + self.extract_fault() + + if len(self._states.keys()) > 0: + node_id = random.choice(list(self._states.keys())) + if len(self._states[node_id].keys()) > 0: + path = random.choice(list(self._states[node_id].keys())) + if self._change_drive_status(node_id, path, EDriveStatus.BROKEN): + self._successful_data_erase_count += 1 + self.on_success_inject_fault() + self.prepare_state() + + def extract_fault(self): + for node_id, path in self._broken_drives: + self._change_drive_status( + node_id, path, EDriveStatus.ACTIVE) + + +def data_storage_nemesis_list(cluster): + return [ + SafelyBreakDisk(cluster), + SafelyCleanupDisks(cluster), + ] diff --git a/ydb/tests/tools/nemesis/library/monitor.py b/ydb/tests/tools/nemesis/library/monitor.py new file mode 100644 index 00000000000..f8a70aee5f2 --- /dev/null +++ b/ydb/tests/tools/nemesis/library/monitor.py @@ -0,0 +1,48 @@ +# -*- coding: utf-8 -*- +import flask +import copy + +from library.python.monlib.metric_registry import MetricRegistry +from library.python.monlib import encoder + + +CONTENT_TYPE_SPACK = 'application/x-solomon-spack' +CONTENT_TYPE_JSON = 'application/json' +app = flask.Flask(__name__) + + +class Monitor(object): + def __init__(self): + self._registry = MetricRegistry() + + @property + def registry(self): + return self._registry + + def int_gauge(self, sensor, labels): + all_labels = copy.deepcopy(labels) + all_labels.update({'sensor': sensor}) + return self._registry.int_gauge(all_labels) + + def rate(self, sensor, labels): + all_labels = copy.deepcopy(labels) + all_labels.update({'sensor': sensor}) + return self._registry.rate(all_labels) + + +_MONITOR = Monitor() + + [email protected]('/sensors') +def sensors(): + if flask.request.headers['accept'] == CONTENT_TYPE_SPACK: + return flask.Response(encoder.dumps(monitor().registry), mimetype=CONTENT_TYPE_SPACK) + return flask.Response(encoder.dumps(monitor().registry, format='json'), mimetype=CONTENT_TYPE_JSON) + + +def monitor(): + return _MONITOR + + +def setup_page(host, port): + app.run(host, port) diff --git a/ydb/tests/tools/nemesis/library/node.py b/ydb/tests/tools/nemesis/library/node.py new file mode 100644 index 00000000000..02d44383dc9 --- /dev/null +++ b/ydb/tests/tools/nemesis/library/node.py @@ -0,0 +1,275 @@ +# -*- coding: utf-8 -*- +import random +import signal +import abc +import six +import itertools +import collections + +from ydb.tests.library.nemesis.nemesis_core import Nemesis, Schedule +from ydb.tests.tools.nemesis.library import base + + [email protected]_metaclass(abc.ABCMeta) +class AbstractKillDaemonNemesis(Nemesis, base.AbstractMonitoredNemesis): + def __init__(self, cluster, schedule): + base.AbstractMonitoredNemesis.__init__(self, scope='node') + Nemesis.__init__(self, schedule=schedule) + self.cluster = cluster + + def prepare_state(self): + self.logger.info("Daemons to kill = {}".format(str(self.daemons))) + + @property + @abc.abstractmethod + def daemons(self): + pass + + def extract_fault(self): + pass + + def inject_fault(self): + if len(self.daemons) == 0: + return self.logger.info("Cannot inject the fault. List of daemons is empty, %s", str(self.daemons)) + + daemon = random.choice(self.daemons) + self.logger.info("Kill daemon %s", str(daemon)) + daemon.kill() + self.on_success_inject_fault() + + +class KillSlotNemesis(AbstractKillDaemonNemesis): + def __init__(self, cluster, schedule=(60, 180)): + super(KillSlotNemesis, self).__init__(cluster, schedule=schedule) + + @property + def daemons(self): + return list(self.cluster.slots.values()) + + +class KillNodeNemesis(AbstractKillDaemonNemesis): + def __init__(self, cluster, schedule=(120, 180)): + super(KillNodeNemesis, self).__init__(cluster, schedule=schedule) + + @property + def daemons(self): + return list(self.cluster.nodes.values()) + + +class KillBlockStoreNodeNemesis(AbstractKillDaemonNemesis): + def __init__(self, cluster, schedule=(120, 180)): + super(KillBlockStoreNodeNemesis, self).__init__(cluster, schedule=schedule) + + @property + def daemons(self): + return list(self.cluster.nbs.values()) + + [email protected]_metaclass(abc.ABCMeta) +class AbstractSerialDaemonKillNemesis(Nemesis, base.AbstractMonitoredNemesis): + def __init__(self, cluster, schedule, schedule_between_kills): + base.AbstractMonitoredNemesis.__init__(self, scope='node') + Nemesis.__init__(self, schedule=schedule) + self.cluster = cluster + self.target = None + self.target_cnt = 0 + self.schedule_between_kills = Schedule.from_tuple_or_int(schedule_between_kills) + + def next_schedule(self): + if self.target is not None: + return next(self.schedule_between_kills) + else: + return super(AbstractSerialDaemonKillNemesis, self).next_schedule() + + def prepare_state(self): + self.logger.info("Daemons to kill = {}".format(str(self.daemons))) + + @property + @abc.abstractmethod + def daemons(self): + pass + + def extract_fault(self): + pass + + def inject_fault(self): + if len(self.daemons) == 0: + return self.logger.info( + "Cannot inject the fault. List of daemons is empty, %s", str( + self.daemons + ) + ) + + if self.target is None: + self.target = random.choice(self.daemons) + self.target_cnt = random.randint(1, 4) + + self.kill_target_daemon() + self.on_success_inject_fault() + + def kill_target_daemon(self): + self.logger.info("Killing node = " + str(self.target)) + self.target.kill() + self.target_cnt -= 1 + if self.target_cnt <= 0: + self.target = None + + +class SerialKillNodeNemesis(AbstractSerialDaemonKillNemesis): + def __init__(self, cluster, schedule=(240, 300), schedule_between_kills=(30, 60)): + super(SerialKillNodeNemesis, self).__init__(cluster, schedule, schedule_between_kills) + + @property + def daemons(self): + return list(self.cluster.nodes.values()) + + +class SerialKillSlotsNemesis(AbstractSerialDaemonKillNemesis): + def __init__(self, cluster, schedule=(240, 300), schedule_between_kills=(30, 60)): + super(SerialKillSlotsNemesis, self).__init__(cluster, schedule, schedule_between_kills) + + @property + def daemons(self): + return list(self.cluster.slots.values()) + + +def nodes_nemesis_list(cluster): + scale_per_cluster = max(1, int(len(cluster.nodes.values()) / 8)) + nemesis_list = [ + KillNodeNemesis(cluster), + + SerialKillNodeNemesis(cluster), + SerialKillSlotsNemesis(cluster), + + RollingUpdateClusterNemesis(cluster), + ] + + for _ in range(scale_per_cluster): + nemesis_list.extend([ + KillSlotNemesis(cluster), + KillBlockStoreNodeNemesis(cluster), + ]) + return nemesis_list + + +class StopStartNodeNemesis(Nemesis, base.AbstractMonitoredNemesis): + def __init__(self, cluster, schedule=(300, 600)): + super(StopStartNodeNemesis, self).__init__(schedule=schedule) + base.AbstractMonitoredNemesis.__init__(self, scope='node') + self._cluster = cluster + self._processes = self._cluster.nodes.values() + self._current_process = None + self.__stop_interval_schedule = Schedule.from_tuple_or_int(30) + self._can_stop = len(self._cluster.nodes.values()) >= 8 + + def next_schedule(self): + if self._current_process is not None: + return next(self.__stop_interval_schedule) + return super(StopStartNodeNemesis, self).next_schedule() + + def prepare_state(self): + self.logger.info("Nodes to stop/start = " + str(self._processes)) + + def inject_fault(self): + if not self._can_stop: + self.logger.info("Node stop is prohibited.") + return + + # we keep exactly one process stopped + if self.extract_fault(): + return + + self.start_inject_fault() + self._current_process = random.choice(self._processes) + self.logger.info("Stopping node = %s", str(self._current_process)) + self._current_process.stop() + self.on_success_inject_fault() + + def extract_fault(self): + if self._current_process is not None: + self.logger.info("Starting node = %s", str(self._current_process)) + self._current_process.start() + self._current_process = None + self.on_success_extract_fault() + return True + return False + + +class SuspendNodeNemesis(Nemesis, base.AbstractMonitoredNemesis): + def __init__(self, cluster, schedule=(600, 1200)): + super(SuspendNodeNemesis, self).__init__(schedule=schedule) + base.AbstractMonitoredNemesis.__init__(self, scope='node') + self._cluster = cluster + self._processes = self._cluster.nodes.values() + self._cluster.slots.values() + self._process_to_wake_up = None + self.__stop_time_schedule = Schedule.from_tuple_or_int((10, 30)) + self._can_suspend = True + + def next_schedule(self): + if self._process_to_wake_up is not None: + return next(self.__stop_time_schedule) + return super(SuspendNodeNemesis, self).next_schedule() + + def prepare_state(self): + self.logger.info("Nodes to suspend = " + str(self._processes)) + self._can_suspend = len(self._cluster.nodes.values()) >= 8 + + def inject_fault(self): + if not self._can_suspend: + self.logger.info("Node suspend is prohibited") + return + + if self.extract_fault(): + return + + self.start_inject_fault() + self._process_to_wake_up = random.choice(self._processes) + self.logger.info("Suspending node = " + str(self._process_to_wake_up)) + self._process_to_wake_up.send_signal(signal.SIGSTOP) + self.on_success_inject_fault() + + def extract_fault(self): + if self._process_to_wake_up is not None: + self.logger.info("Continuing node = " + str(self._process_to_wake_up)) + self._process_to_wake_up.send_signal(signal.SIGCONT) + self._process_to_wake_up = None + self.on_success_extract_fault() + return True + return False + + +class RollingUpdateClusterNemesis(Nemesis, base.AbstractMonitoredNemesis): + def __init__(self, cluster, schedule=(60, 70)): + super(RollingUpdateClusterNemesis, self).__init__(schedule=schedule) + base.AbstractMonitoredNemesis.__init__(self, scope='node') + self.cluster = cluster + self.buckets = {0: collections.deque(self.cluster.nodes.values()), 1: collections.deque()} + self.step_id = itertools.count(start=1) + self.slots_by_host = collections.defaultdict(list) + + def prepare_state(self): + self.logger.info("Initializing rolling update nemesis....") + self.logger.info("Nodes to update = %s", str(self.cluster.nodes.values())) + self.logger.info("Slots to update = %s", str(self.cluster.slots.values())) + + for slot in self.cluster.slots.values(): + self.slots_by_host[slot.host].append(slot) + + def extract_fault(self): + pass + + def inject_fault(self): + self.logger.info("Starting next (%d-th) iteration of rolling update process...." % next(self.step_id)) + bucket_id = 0 if len(self.buckets[0]) >= len(self.buckets[1]) else 1 + + node = self.buckets[bucket_id].popleft() + self.buckets[bucket_id ^ 1].append(node) + self.logger.info("Update nodes on host %s, direction is %s -> %s" % (node.host, bucket_id, bucket_id ^ 1)) + node.switch_version() + node.kill() + + self.logger.info("Successfully updated version on host %s" % node.host) + for slot in self.slots_by_host.get(node.host, []): + slot.kill() + + self.on_success_inject_fault() diff --git a/ydb/tests/tools/nemesis/library/tablet.py b/ydb/tests/tools/nemesis/library/tablet.py new file mode 100644 index 00000000000..8ca5e52625d --- /dev/null +++ b/ydb/tests/tools/nemesis/library/tablet.py @@ -0,0 +1,256 @@ +# -*- coding: utf-8 -*- +import abc +import random +import six + +from ydb.tests.library.nemesis.nemesis_core import Nemesis +from ydb.tests.library.common.types import TabletTypes +from ydb.tests.library.harness.kikimr_client import kikimr_client_factory +from ydb.tests.library.harness.kikimr_http_client import HiveClient +from ydb.tests.tools.nemesis.library.base import AbstractMonitoredNemesis + + [email protected]_metaclass(abc.ABCMeta) +class AbstractTabletByTypeNemesis(Nemesis, AbstractMonitoredNemesis): + def __init__(self, tablet_type, cluster, schedule): + AbstractMonitoredNemesis.__init__(self, scope='tablets') + Nemesis.__init__(self, schedule=schedule) + self.cluster = cluster + self.__client = None + self.__tablet_type = tablet_type + self.__tablet_ids = [] + + @property + def tablet_type(self): + return self.__tablet_type + + @property + def client(self): + if self.__client is None: + self.__client = kikimr_client_factory( + self.cluster.nodes[1].host, self.cluster.nodes[1].grpc_port, retry_count=10) + return self.__client + + def extract_fault(self): + pass + + @property + def tablet_ids(self): + return self.__tablet_ids + + def prepare_state(self): + self.logger.info('Preparing state for nemesis = ' + str(self)) + response = self.client.tablet_state(self.__tablet_type) + + self.__tablet_ids = [ + info.TabletId for info in response.TabletStateInfo + ] + + def __str__(self): + return "{class_name}(tablet_type={tablet_type})".format( + class_name=self.__class__.__name__, + tablet_type=self.__tablet_type + ) + + +class KillSystemTabletByTypeNemesis(AbstractTabletByTypeNemesis): + + def __init__(self, tablet_type, cluster, schedule=(45, 90)): + super(KillSystemTabletByTypeNemesis, self).__init__(tablet_type, cluster, schedule) + self.__cluster = cluster + self.__client = None + self.__tablet_type = tablet_type + + def inject_fault(self): + if self.tablet_ids: + tablet_id = random.choice(self.tablet_ids) + self.logger.info( + "Killing {tablet_type}, tablet_id = {tablet_id}".format( + tablet_type=self.__tablet_type, + tablet_id=tablet_id + ) + ) + try: + self.client.tablet_kill(tablet_id) + self.on_success_inject_fault() + except RuntimeError: + self.logger.error( + "Failed to kill {tablet_type}, tablet_id = {tablet_id}".format( + tablet_type=self.__tablet_type, + tablet_id=tablet_id + ) + ) + else: + self.prepare_state() + + +class KillCoordinatorNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster): + super(KillCoordinatorNemesis, self).__init__(TabletTypes.FLAT_TX_COORDINATOR, cluster) + + +class KillMediatorNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster): + super(KillMediatorNemesis, self).__init__(TabletTypes.TX_MEDIATOR, cluster) + + +class KillDataShardNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(30, 60)): + super(KillDataShardNemesis, self).__init__(TabletTypes.FLAT_DATASHARD, cluster, schedule) + + +class KillHiveNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(120, 240)): + super(KillHiveNemesis, self).__init__(TabletTypes.FLAT_HIVE, cluster, schedule) + + +class KillBsControllerNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(300, 600)): + super(KillBsControllerNemesis, self).__init__(TabletTypes.FLAT_BS_CONTROLLER, cluster, schedule) + + +class KillSchemeShardNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(120, 240)): + super(KillSchemeShardNemesis, self).__init__(TabletTypes.FLAT_SCHEMESHARD, cluster, schedule) + + +class KillPersQueueNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(30, 60)): + super(KillPersQueueNemesis, self).__init__(TabletTypes.PERSQUEUE, cluster, schedule) + + +class KillKeyValueNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(30, 60)): + super(KillKeyValueNemesis, self).__init__(TabletTypes.KEYVALUEFLAT, cluster, schedule) + + +class KillTxAllocatorNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(120, 240)): + super(KillTxAllocatorNemesis, self).__init__(TabletTypes.TX_ALLOCATOR, cluster, schedule) + + +class KillNodeBrokerNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(120, 240)): + super(KillNodeBrokerNemesis, self).__init__(TabletTypes.NODE_BROKER, cluster, schedule) + + +class KillTenantSlotBrokerNemesis(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(120, 240)): + super(KillTenantSlotBrokerNemesis, self).__init__(TabletTypes.TENANT_SLOT_BROKER, cluster, schedule) + + +class KillBlocktoreVolume(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(30, 60)): + super(KillBlocktoreVolume, self).__init__(TabletTypes.BLOCKSTORE_VOLUME, cluster, schedule) + + +class KillBlocktorePartition(KillSystemTabletByTypeNemesis): + def __init__(self, cluster, schedule=(30, 60)): + super(KillBlocktorePartition, self).__init__(TabletTypes.BLOCKSTORE_PARTITION, cluster, schedule) + + +class KickTabletsFromNode(Nemesis, AbstractMonitoredNemesis): + def __init__(self, cluster, schedule=(150, 450)): + Nemesis.__init__(self, schedule=schedule) + AbstractMonitoredNemesis.__init__(self, scope='tablets') + self.__nodes = cluster.nodes.values() + self.hive = HiveClient(cluster.nodes[1].host, cluster.nodes[1].mon_port) + + def prepare_state(self): + pass + + def inject_fault(self): + if self.__nodes: + try: + node = random.choice(self.__nodes) + node_id = node.node_id + self.hive.block_node(node_id) + self.hive.kick_tablets_from_node(node_id) + self.hive.unblock_node(node_id) + self.on_success_inject_fault() + except Exception as e: + self.logger.error( + "Failed to inject fault, %s", str( + e + ) + ) + + def extract_fault(self): + pass + + +class ChangeTabletGroupNemesis(AbstractTabletByTypeNemesis): + + def __init__(self, cluster, tablet_type, schedule=(30, 60), channels=()): + super(ChangeTabletGroupNemesis, self).__init__(tablet_type=tablet_type, cluster=cluster, schedule=schedule) + self.hive = HiveClient(cluster.nodes[1].host, cluster.nodes[1].mon_port) + self.channels = channels + + def inject_fault(self): + if self.tablet_ids: + tablet_id = random.choice(self.tablet_ids) + try: + self.hive.change_tablet_group(tablet_id, channels=self.channels) + self.on_success_inject_fault() + except Exception as e: + self.logger.error( + "Failed to inject fault, %s", str( + e + ) + ) + else: + self.prepare_state() + + +class BulkChangeTabletGroupNemesis(AbstractTabletByTypeNemesis): + def __init__(self, cluster, tablet_type, schedule=(60, 180), channels=(), percent=None): + super(BulkChangeTabletGroupNemesis, self).__init__(tablet_type=tablet_type, cluster=cluster, schedule=schedule) + self.hive = HiveClient(cluster.nodes[1].host, cluster.nodes[1].mon_port) + self.channels = channels + self.percent = percent + + def inject_fault(self): + try: + self.logger.info('Injecting fault for nemesis = ' + str(self)) + self.hive.change_tablet_group_by_tablet_type(self.tablet_type, percent=self.percent, channels=self.channels) + self.on_success_inject_fault() + except Exception as e: + self.logger.error( + "Failed to inject fault, %s", str( + e + ) + ) + + +class ReBalanceTabletsNemesis(Nemesis, AbstractMonitoredNemesis): + def __init__(self, cluster, schedule=(60, 180)): + Nemesis.__init__(self, schedule=schedule) + AbstractMonitoredNemesis.__init__(self) + self.hive = HiveClient(cluster.nodes[1].host, cluster.nodes[1].mon_port) + + def prepare_state(self): + pass + + def inject_fault(self): + try: + self.hive.rebalance_all_tablets() + self.on_success_inject_fault() + except Exception as e: + self.logger.error( + "Failed to inject fault, %s", str( + e + ) + ) + + def extract_fault(self): + pass + + +def change_tablet_group_nemesis_list(cluster): + result = [] + scale_per_cluster = max(int(len(cluster.nodes.values()) / 8), 1) + for tablet_type in (TabletTypes.PERSQUEUE, TabletTypes.FLAT_DATASHARD, TabletTypes.KEYVALUEFLAT): + for _ in range(scale_per_cluster): + result.append(ChangeTabletGroupNemesis(cluster, tablet_type=tablet_type)) + result.append(BulkChangeTabletGroupNemesis(cluster, tablet_type)) + return result diff --git a/ydb/tests/tools/nemesis/library/ya.make b/ydb/tests/tools/nemesis/library/ya.make new file mode 100644 index 00000000000..a407c0e851b --- /dev/null +++ b/ydb/tests/tools/nemesis/library/ya.make @@ -0,0 +1,22 @@ +SUBSCRIBER(g:kikimr) + +PY3_LIBRARY() + +PY_SRCS( + __init__.py + base.py + catalog.py + disk.py + node.py + tablet.py + monitor.py +) + +PEERDIR( + contrib/python/Flask + ydb/tests/library + library/python/monlib + ydb/core/protos +) + +END() diff --git a/ydb/tests/tools/nemesis/ut/__init__.py b/ydb/tests/tools/nemesis/ut/__init__.py new file mode 100644 index 00000000000..e69de29bb2d --- /dev/null +++ b/ydb/tests/tools/nemesis/ut/__init__.py diff --git a/ydb/tests/tools/nemesis/ut/test_disk.py b/ydb/tests/tools/nemesis/ut/test_disk.py new file mode 100644 index 00000000000..5513b6055f4 --- /dev/null +++ b/ydb/tests/tools/nemesis/ut/test_disk.py @@ -0,0 +1,42 @@ +# -*- coding: utf-8 -*- +from hamcrest import assert_that, greater_than_or_equal_to + +from ydb.tests.library.harness.kikimr_cluster import kikimr_cluster_factory +from ydb.tests.library.harness.kikimr_config import KikimrConfigGenerator +from ydb.tests.library.common import types +from ydb.tests.library.common.wait_for import wait_for +from ydb.tests.tools.nemesis.library import disk + + +class TestSafeDiskBreak(object): + + @classmethod + def setup_class(cls): + cls.cluster = kikimr_cluster_factory( + KikimrConfigGenerator(erasure=types.Erasure.BLOCK_4_2, nodes=9, use_in_memory_pdisks=True)) + cls.cluster.start() + + @classmethod + def teardown_class(cls): + if hasattr(cls, 'cluster'): + cls.cluster.stop() + + def test_erase_method(self): + nemesis = disk.SafelyBreakDisk(self.cluster) + + def predicate(): + nemesis.prepare_state() + return nemesis.state_ready + + wait_for(predicate, 180) + for _ in range(100): + nemesis.inject_fault() + if nemesis.successful_data_erase_count >= 1: + break + + assert_that( + nemesis.successful_data_erase_count, + greater_than_or_equal_to( + 1 + ) + ) diff --git a/ydb/tests/tools/nemesis/ut/test_tablet.py b/ydb/tests/tools/nemesis/ut/test_tablet.py new file mode 100644 index 00000000000..a19eb28dbc3 --- /dev/null +++ b/ydb/tests/tools/nemesis/ut/test_tablet.py @@ -0,0 +1,74 @@ +#!/usr/bin/env python +# -*- coding: utf-8 -*- +from hamcrest import assert_that + +from ydb.tests.library.common.delayed import wait_tablets_state_by_id +from ydb.tests.library.common.protobuf import TCmdCreateTablet +from ydb.tests.library.common.types import Erasure, TabletTypes, TabletStates +from ydb.tests.library.harness.util import LogLevels +from ydb.tests.tools.nemesis.library.tablet import BulkChangeTabletGroupNemesis +from ydb.tests.tools.nemesis.library.tablet import KillHiveNemesis, KillBsControllerNemesis +from ydb.tests.library.harness.kikimr_cluster import kikimr_cluster_factory +from ydb.tests.library.harness.kikimr_config import KikimrConfigGenerator +from ydb.tests.library.matchers.tablets import all_tablets_are_created + +TIMEOUT_SECONDS = 480 +TABLETS_PER_NODE = 100 + + +class TestMassiveKills(object): + @classmethod + def setup_class(cls): + cls.configurator = KikimrConfigGenerator( + erasure=Erasure.BLOCK_4_2, + additional_log_configs={ + 'HIVE': LogLevels.DEBUG, + 'LOCAL': LogLevels.DEBUG, + 'BS_CONTROLLER': LogLevels.DEBUG, + } + ) + cls.cluster = kikimr_cluster_factory(configurator=cls.configurator) + cls.cluster.start() + + @classmethod + def teardown_class(cls): + if hasattr(cls, 'cluster'): + cls.cluster.stop() + + def test_tablets_are_ok_after_many_kills(self): + nodes_count = len(self.cluster.nodes.values()) + cmd_create_tablets = [ + TCmdCreateTablet(owner_id=234843, owner_idx=index, type=TabletTypes.KEYVALUEFLAT) + for index in range(nodes_count * TABLETS_PER_NODE) + ] + + create_response = self.cluster.client.hive_create_tablets(cmd_create_tablets) + all_tablet_ids = set([tablet.TabletId for tablet in create_response.CreateTabletResult]) + + assert_that(create_response, all_tablets_are_created(cmd_create_tablets)) + + wait_tablets_state_by_id( + self.cluster.client, TabletStates.Active, tablet_ids=all_tablet_ids, timeout_seconds=TIMEOUT_SECONDS + ) + nemesis = [ + BulkChangeTabletGroupNemesis(self.cluster, TabletTypes.KEYVALUEFLAT), + KillHiveNemesis(self.cluster), + KillBsControllerNemesis(self.cluster), + ] + + for element in nemesis: + element.prepare_state() + + for _ in range(3): + for element in nemesis: + element.inject_fault() + + for tablet_id in all_tablet_ids: + self.cluster.client.tablet_kill(tablet_id) + + wait_tablets_state_by_id( + self.cluster.client, + TabletStates.Active, + tablet_ids=all_tablet_ids, + timeout_seconds=TIMEOUT_SECONDS, + ) diff --git a/ydb/tests/tools/nemesis/ut/ya.make b/ydb/tests/tools/nemesis/ut/ya.make new file mode 100644 index 00000000000..0cd7bb09e7f --- /dev/null +++ b/ydb/tests/tools/nemesis/ut/ya.make @@ -0,0 +1,26 @@ +SUBSCRIBER(g:kikimr) + +PY3TEST() +ENV(YDB_DRIVER_BINARY="ydb/apps/ydbd/ydbd") + +TEST_SRCS( + test_disk.py + test_tablet.py +) + +TIMEOUT(600) +SIZE(MEDIUM) + + +DEPENDS( + ydb/apps/ydbd +) + +PEERDIR( + ydb/tests/tools/nemesis/library +) + +FORK_SUBTESTS() +FORK_TEST_FILES() + +END() diff --git a/ydb/tests/tools/nemesis/ya.make b/ydb/tests/tools/nemesis/ya.make new file mode 100644 index 00000000000..44ddfabe3e0 --- /dev/null +++ b/ydb/tests/tools/nemesis/ya.make @@ -0,0 +1,5 @@ +RECURSE( + driver + library + ut +) diff --git a/ydb/tests/tools/ya.make b/ydb/tests/tools/ya.make index ce62665f270..b7586673bba 100644 --- a/ydb/tests/tools/ya.make +++ b/ydb/tests/tools/ya.make @@ -5,6 +5,7 @@ RECURSE( idx_test kqprun mdb_mock + nemesis pq_read s3_recipe token_accessor_mock |
