summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authormregrock <[email protected]>2025-04-23 14:08:25 +0300
committerGitHub <[email protected]>2025-04-23 14:08:25 +0300
commit8df4e48482c60a29fed06be86ac1e987127a1bec (patch)
tree36bfc7d161e500afaa1fa0c6b87f1d7e2f521baa
parentf1d45a21704af583d8f061338aa97b9fbb579d0f (diff)
Tests for distconf (#14953)
Co-authored-by: Alexander Rutkovsky <[email protected]>
-rw-r--r--ydb/tests/functional/config/test_config_with_metadata.py1
-rw-r--r--ydb/tests/functional/config/test_distconf.py196
-rw-r--r--ydb/tests/functional/config/test_generate_dynamic_config.py2
-rw-r--r--ydb/tests/functional/config/ya.make4
-rw-r--r--ydb/tests/library/harness/kikimr_config.py124
-rw-r--r--ydb/tests/library/harness/kikimr_runner.py148
-rw-r--r--ydb/tests/library/harness/param_constants.py7
7 files changed, 441 insertions, 41 deletions
diff --git a/ydb/tests/functional/config/test_config_with_metadata.py b/ydb/tests/functional/config/test_config_with_metadata.py
index 1b4dae0ab10..454791e89ec 100644
--- a/ydb/tests/functional/config/test_config_with_metadata.py
+++ b/ydb/tests/functional/config/test_config_with_metadata.py
@@ -41,6 +41,7 @@ class AbstractKiKiMRTest(object):
nodes=nodes_count,
use_in_memory_pdisks=False,
metadata_section=cls.metadata_section,
+ simple_config=True,
)
cls.cluster = KiKiMR(configurator=configurator)
cls.cluster.start()
diff --git a/ydb/tests/functional/config/test_distconf.py b/ydb/tests/functional/config/test_distconf.py
new file mode 100644
index 00000000000..387020d168f
--- /dev/null
+++ b/ydb/tests/functional/config/test_distconf.py
@@ -0,0 +1,196 @@
+# -*- coding: utf-8 -*-
+import logging
+import yaml
+import tempfile
+from hamcrest import assert_that
+import time
+
+from ydb.tests.library.common.types import Erasure
+import ydb.tests.library.common.cms as cms
+from ydb.tests.library.clients.kikimr_http_client import SwaggerClient
+from ydb.tests.library.harness.kikimr_runner import KiKiMR
+from ydb.tests.library.clients.kikimr_config_client import ConfigClient
+from ydb.tests.library.harness.kikimr_config import KikimrConfigGenerator
+from ydb.tests.library.kv.helpers import create_kv_tablets_and_wait_for_start
+from ydb.public.api.protos.ydb_status_codes_pb2 import StatusIds
+from ydb.tests.library.harness.util import LogLevels
+
+import ydb.public.api.protos.ydb_config_pb2 as config
+
+logger = logging.getLogger(__name__)
+
+
+def value_for(key, tablet_id):
+ return "Value: <key = {key}, tablet_id = {tablet_id}>".format(
+ key=key, tablet_id=tablet_id)
+
+
+def fetch_config(config_client):
+ fetch_config_response = config_client.fetch_all_configs()
+ assert_that(fetch_config_response.operation.status == StatusIds.SUCCESS)
+
+ result = config.FetchConfigResult()
+ fetch_config_response.operation.result.Unpack(result)
+ return result.config[0].config
+
+
+def get_config_version(yaml_config):
+ config = yaml.safe_load(yaml_config)
+ return config.get('metadata', {}).get('version', 0)
+
+
+class DistConfKiKiMRTest(object):
+ erasure = Erasure.BLOCK_4_2
+ use_config_store = True
+ separate_node_configs = True
+ metadata_section = {
+ "kind": "MainConfig",
+ "version": 0,
+ "cluster": "",
+ }
+
+ @classmethod
+ def setup_class(cls):
+ nodes_count = 8 if cls.erasure == Erasure.BLOCK_4_2 else 9
+ log_configs = {
+ 'BS_NODE': LogLevels.DEBUG,
+ 'GRPC_SERVER': LogLevels.DEBUG,
+ 'GRPC_PROXY': LogLevels.DEBUG,
+ 'TX_PROXY': LogLevels.DEBUG,
+ 'TICKET_PARSER': LogLevels.DEBUG,
+ }
+ cls.configurator = KikimrConfigGenerator(
+ cls.erasure,
+ nodes=nodes_count,
+ use_in_memory_pdisks=False,
+ use_config_store=cls.use_config_store,
+ metadata_section=cls.metadata_section,
+ separate_node_configs=cls.separate_node_configs,
+ simple_config=True,
+ use_self_management=True,
+ extra_grpc_services=['config'],
+ additional_log_configs=log_configs)
+
+ cls.cluster = KiKiMR(configurator=cls.configurator)
+ cls.cluster.start()
+
+ cms.request_increase_ratio_limit(cls.cluster.client)
+ host = cls.cluster.nodes[1].host
+ grpc_port = cls.cluster.nodes[1].port
+ cls.swagger_client = SwaggerClient(host, cls.cluster.nodes[1].mon_port)
+ cls.config_client = ConfigClient(host, grpc_port)
+
+ @classmethod
+ def teardown_class(cls):
+ cls.cluster.stop()
+
+ def check_kikimr_is_operational(self, table_path, tablet_ids):
+ for partition_id, tablet_id in enumerate(tablet_ids):
+ write_resp = self.cluster.kv_client.kv_write(
+ table_path, partition_id, "key", value_for("key", tablet_id)
+ )
+ assert_that(write_resp.operation.status == StatusIds.SUCCESS)
+
+ read_resp = self.cluster.kv_client.kv_read(
+ table_path, partition_id, "key"
+ )
+ assert_that(read_resp.operation.status == StatusIds.SUCCESS)
+
+
+class TestKiKiMRDistConfBasic(DistConfKiKiMRTest):
+
+ def test_cluster_is_operational_with_distconf(self):
+ table_path = '/Root/mydb/mytable_with_metadata'
+ number_of_tablets = 5
+ tablet_ids = create_kv_tablets_and_wait_for_start(
+ self.cluster.client,
+ self.cluster.kv_client,
+ self.swagger_client,
+ number_of_tablets,
+ table_path,
+ timeout_seconds=3
+ )
+ self.check_kikimr_is_operational(table_path, tablet_ids)
+
+ def test_cluster_expand_with_distconf(self):
+ table_path = '/Root/mydb/mytable_with_expand'
+ number_of_tablets = 5
+
+ tablet_ids = create_kv_tablets_and_wait_for_start(
+ self.cluster.client,
+ self.cluster.kv_client,
+ self.swagger_client,
+ number_of_tablets,
+ table_path,
+ timeout_seconds=3
+ )
+
+ current_node_ids = list(self.cluster.nodes.keys())
+ expected_new_node_id = max(current_node_ids) + 1
+
+ node_port_allocator = self.configurator.port_allocator.get_node_port_allocator(expected_new_node_id)
+
+ fetched_config = fetch_config(self.config_client)
+ dumped_fetched_config = yaml.safe_load(fetched_config)
+ config_section = dumped_fetched_config["config"]
+
+ # create new pdisk
+ tmp_file = tempfile.NamedTemporaryFile(prefix="pdisk{}".format(1), suffix=".data",
+ dir=None)
+ pdisk_path = tmp_file.name
+
+ # add new host config
+ host_config_id = len(config_section["host_configs"]) + 1
+ config_section["host_configs"].append({
+ "drive": [
+ {
+ "path": pdisk_path,
+ "type": "ROT"
+ }
+ ],
+ "host_config_id": host_config_id,
+ })
+
+ # add new node in hosts
+ config_section["hosts"].append({
+ "host_config_id": host_config_id,
+ "host": "localhost",
+ "port": node_port_allocator.ic_port,
+ })
+ self.configurator.full_config = dumped_fetched_config
+
+ # prepare new node
+ new_node = self.cluster.prepare_node(self.configurator)
+ new_node.format_pdisk(pdisk_path, self.configurator.static_pdisk_size)
+ dumped_fetched_config["metadata"]["version"] = 1
+
+ # replace config
+ replace_config_response = self.config_client.replace_config(yaml.dump(dumped_fetched_config))
+ assert_that(replace_config_response.operation.status == StatusIds.SUCCESS)
+ # start new node
+ new_node.start()
+
+ self.check_kikimr_is_operational(table_path, tablet_ids)
+
+ time.sleep(5)
+
+ try:
+ pdisk_info = self.swagger_client.pdisk_info(new_node.node_id)
+
+ pdisks_list = pdisk_info['PDiskStateInfo']
+
+ found_pdisk_in_viewer = False
+ for pdisk_entry in pdisks_list:
+ node_id_in_entry = pdisk_entry.get('NodeId')
+ path_in_entry = pdisk_entry.get('Path')
+ state_in_entry = pdisk_entry.get('State')
+
+ if node_id_in_entry == new_node.node_id and path_in_entry == pdisk_path:
+ logger.info(f"Found matching PDisk in viewer: NodeId={node_id_in_entry}, Path={path_in_entry}, State={state_in_entry}")
+ found_pdisk_in_viewer = True
+ break
+ except Exception as e:
+ logger.error(f"Viewer API check failed: {e}", exc_info=True)
+ if 'pdisk_info' in locals():
+ logger.error(f"Viewer API response content: {pdisk_info}")
+ raise
diff --git a/ydb/tests/functional/config/test_generate_dynamic_config.py b/ydb/tests/functional/config/test_generate_dynamic_config.py
index 4b3884b504a..a51cd75db95 100644
--- a/ydb/tests/functional/config/test_generate_dynamic_config.py
+++ b/ydb/tests/functional/config/test_generate_dynamic_config.py
@@ -30,6 +30,7 @@ class AbstractKiKiMRTest(object):
configurator = KikimrConfigGenerator(cls.erasure,
nodes=nodes_count,
use_in_memory_pdisks=False,
+ simple_config=True,
)
cls.cluster = KiKiMR(configurator=configurator)
cls.cluster.start()
@@ -77,6 +78,7 @@ class TestGenerateDynamicConfigFromConfigDir(AbstractKiKiMRTest):
use_in_memory_pdisks=False,
use_config_store=True,
separate_node_configs=True,
+ simple_config=True,
)
cls.cluster = KiKiMR(configurator=configurator)
diff --git a/ydb/tests/functional/config/ya.make b/ydb/tests/functional/config/ya.make
index de0bde4fbf0..4626c8c968d 100644
--- a/ydb/tests/functional/config/ya.make
+++ b/ydb/tests/functional/config/ya.make
@@ -3,6 +3,7 @@ PY3TEST()
TEST_SRCS(
test_config_with_metadata.py
test_generate_dynamic_config.py
+ test_distconf.py
)
SPLIT_FACTOR(10)
@@ -22,8 +23,11 @@ ENDIF()
ENV(YDB_DRIVER_BINARY="ydb/apps/ydbd/ydbd")
+ENV(YDB_CLI_BINARY="ydb/apps/ydb/ydb")
+ENV(IAM_TOKEN="")
DEPENDS(
ydb/apps/ydbd
+ ydb/apps/ydb
)
PEERDIR(
diff --git a/ydb/tests/library/harness/kikimr_config.py b/ydb/tests/library/harness/kikimr_config.py
index 86bbd1fba3c..cd7d1d915c4 100644
--- a/ydb/tests/library/harness/kikimr_config.py
+++ b/ydb/tests/library/harness/kikimr_config.py
@@ -18,7 +18,7 @@ from ydb.tests.library.common.types import Erasure
from . import tls_tools
from .kikimr_port_allocator import KikimrPortManagerPortAllocator
-from .param_constants import kikimr_driver_path
+from .param_constants import kikimr_driver_path, ydb_cli_path
from .util import LogLevels
PDISK_SIZE_STR = os.getenv("YDB_PDISK_SIZE", str(64 * 1024 * 1024 * 1024))
@@ -166,6 +166,8 @@ class KikimrConfigGenerator(object):
grouped_memory_limiter_config=None,
query_service_config=None,
domain_login_only=None,
+ use_self_management=False,
+ simple_config=False,
):
if extra_feature_flags is None:
extra_feature_flags = []
@@ -173,7 +175,11 @@ class KikimrConfigGenerator(object):
extra_grpc_services = []
self.use_log_files = use_log_files
+ self.use_self_management = use_self_management
+ self.simple_config = simple_config
self.suppress_version_check = suppress_version_check
+ if use_self_management:
+ self.suppress_version_check = False
self._pdisk_store_path = pdisk_store_path
self.static_pdisk_size = static_pdisk_size
self.app_config = config_pb2.TAppConfig()
@@ -248,6 +254,11 @@ class KikimrConfigGenerator(object):
self.yaml_config = _load_default_yaml(self.__node_ids, self.domain_name, self.static_erasure, self.__additional_log_configs)
+ security_config_root = self.yaml_config["domains_config"]
+ if self.use_self_management:
+ self.yaml_config["self_management_config"] = dict()
+ self.yaml_config["self_management_config"]["enabled"] = True
+
if overrided_actor_system_config:
self.yaml_config["actor_system_config"] = overrided_actor_system_config
@@ -387,14 +398,14 @@ class KikimrConfigGenerator(object):
if default_users is not None:
# check for None for remove default users for empty dict
- if "security_config" not in self.yaml_config["domains_config"]:
- self.yaml_config["domains_config"]["security_config"] = dict()
+ if "security_config" not in security_config_root:
+ security_config_root["security_config"] = dict()
# remove existed default users
- self.yaml_config["domains_config"]["security_config"]["default_users"] = []
+ security_config_root["security_config"]["default_users"] = []
for user, password in default_users.items():
- self.yaml_config["domains_config"]["security_config"]["default_users"].append({
+ security_config_root["security_config"]["default_users"].append({
"name": user,
"password": password,
})
@@ -403,10 +414,10 @@ class KikimrConfigGenerator(object):
self.yaml_config["monitoring_config"] = {"allow_origin": str(os.getenv("YDB_ALLOW_ORIGIN"))}
if enforce_user_token_requirement:
- self.yaml_config["domains_config"]["security_config"]["enforce_user_token_requirement"] = True
+ security_config_root["security_config"]["enforce_user_token_requirement"] = True
if default_user_sid:
- self.yaml_config["domains_config"]["security_config"]["default_user_sids"] = [default_user_sid]
+ security_config_root["security_config"]["default_user_sids"] = [default_user_sid]
if os.getenv("YDB_HARD_MEMORY_LIMIT_BYTES"):
self.yaml_config["memory_controller_config"] = {"hard_limit_bytes": int(os.getenv("YDB_HARD_MEMORY_LIMIT_BYTES"))}
@@ -465,6 +476,26 @@ class KikimrConfigGenerator(object):
self.yaml_config["kafka_proxy_config"] = kafka_proxy_config
self.full_config = dict()
+ if self.use_self_management:
+ self.yaml_config["domains_config"].pop("security_config")
+ self.yaml_config["default_disk_type"] = "ROT"
+ self.yaml_config["fail_domain_type"] = "rack"
+ self._add_host_config_and_hosts()
+ self.yaml_config["erasure"] = self.yaml_config.pop("static_erasure")
+
+ for name in ['blob_storage_config', 'domains_config', 'nameservice_config', 'system_tablets', 'grpc_config',
+ 'channel_profile_config', 'interconnect_config']:
+ del self.yaml_config[name]
+ if self.simple_config:
+ self.yaml_config.pop("feature_flags")
+ self.yaml_config.pop("federated_query_config")
+ self.yaml_config.pop("pqconfig")
+ self.yaml_config.pop("pqcluster_discovery_config")
+ self.yaml_config.pop("net_classifier_config")
+ self.yaml_config.pop("sqs_config")
+ self.yaml_config.pop("table_service_config")
+ self.yaml_config.pop("kqpconfig")
+
if metadata_section:
self.full_config["metadata"] = metadata_section
self.full_config["config"] = self.yaml_config
@@ -568,6 +599,9 @@ class KikimrConfigGenerator(object):
binary_paths = [kikimr_driver_path()]
return binary_paths[node_id % len(binary_paths)]
+ def get_ydb_cli_path(self):
+ return ydb_cli_path()
+
def write_tls_data(self):
if self.__grpc_ssl_enable:
for fpath, data in (
@@ -650,25 +684,8 @@ class KikimrConfigGenerator(object):
{"node_id": node_id, "pdisk_id": pdisk_id, "pdisk_guid": pdisk_id, 'vdisk_slot_id': 0}]}
)
- def __build(self):
+ def _initialize_pdisks_info(self):
datacenter_id_generator = itertools.cycle(self._dcs)
- self.yaml_config["blob_storage_config"] = {}
- if self.__bs_cache_file_path:
- self.yaml_config["blob_storage_config"]["cache_file_path"] = \
- self.__bs_cache_file_path
- self.yaml_config["blob_storage_config"]["service_set"] = {}
- self.yaml_config["blob_storage_config"]["service_set"]["availability_domains"] = 1
- self.yaml_config["blob_storage_config"]["service_set"]["pdisks"] = []
- self.yaml_config["blob_storage_config"]["service_set"]["vdisks"] = []
- self.yaml_config["blob_storage_config"]["service_set"]["groups"] = [
- {"group_id": 0, 'group_generation': 1, 'erasure_species': int(self.static_erasure)}]
- self.yaml_config["blob_storage_config"]["service_set"]["groups"][0]["rings"] = []
-
- for dc in self._dcs:
- self.yaml_config["blob_storage_config"]["service_set"]["groups"][0]["rings"].append({"fail_domains": []})
-
- self._add_state_storage_config()
-
for node_id in self.__node_ids:
datacenter_id = next(datacenter_id_generator)
@@ -689,7 +706,7 @@ class KikimrConfigGenerator(object):
self._pdisks_info.append({'pdisk_path': pdisk_path, 'node_id': node_id, 'disk_size': disk_size,
'pdisk_user_kind': pdisk_user_kind})
- if pdisk_id == 1 and node_id <= self.static_erasure.min_fail_domains * self._rings_count:
+ if not self.use_self_management and pdisk_id == 1 and node_id <= self.static_erasure.min_fail_domains * self._rings_count:
self._add_pdisk_to_static_group(
pdisk_id,
pdisk_path,
@@ -697,3 +714,58 @@ class KikimrConfigGenerator(object):
pdisk_user_kind,
datacenter_id - 1,
)
+
+ def _add_host_config_and_hosts(self):
+ self._initialize_pdisks_info()
+ host_configs = []
+ hosts = []
+ host_config_id_counter = itertools.count(1)
+
+ for node_id in self.__node_ids:
+ host_config_id = next(host_config_id_counter)
+ drive = []
+ for pdisk_info in self._pdisks_info:
+ if pdisk_info['node_id'] == node_id:
+ drive.append(
+ {
+ "path": pdisk_info['pdisk_path'],
+ "type": pdisk_info.get('pdisk_type', 'ROT').upper(),
+ }
+ )
+
+ host_configs.append(
+ {
+ "host_config_id": host_config_id,
+ "drive": drive,
+ }
+ )
+ hosts.append(
+ {
+ "host": "localhost",
+ "port": self.port_allocator.get_node_port_allocator(node_id).ic_port,
+ "host_config_id": host_config_id,
+ }
+ )
+
+ self.yaml_config["host_configs"] = host_configs
+ self.yaml_config["hosts"] = hosts
+
+ def __build(self):
+ self.yaml_config["blob_storage_config"] = {}
+ if self.__bs_cache_file_path:
+ self.yaml_config["blob_storage_config"]["cache_file_path"] = \
+ self.__bs_cache_file_path
+ self.yaml_config["blob_storage_config"]["service_set"] = {}
+ self.yaml_config["blob_storage_config"]["service_set"]["availability_domains"] = 1
+ self.yaml_config["blob_storage_config"]["service_set"]["pdisks"] = []
+ self.yaml_config["blob_storage_config"]["service_set"]["vdisks"] = []
+ self.yaml_config["blob_storage_config"]["service_set"]["groups"] = [
+ {"group_id": 0, 'group_generation': 1, 'erasure_species': int(self.static_erasure)}]
+ self.yaml_config["blob_storage_config"]["service_set"]["groups"][0]["rings"] = []
+
+ for dc in self._dcs:
+ self.yaml_config["blob_storage_config"]["service_set"]["groups"][0]["rings"].append({"fail_domains": []})
+
+ self._add_state_storage_config()
+ if not self.use_self_management:
+ self._initialize_pdisks_info()
diff --git a/ydb/tests/library/harness/kikimr_runner.py b/ydb/tests/library/harness/kikimr_runner.py
index 88aed13839e..2ec99b8bb53 100644
--- a/ydb/tests/library/harness/kikimr_runner.py
+++ b/ydb/tests/library/harness/kikimr_runner.py
@@ -9,6 +9,7 @@ import threading
from importlib_resources import read_binary
from google.protobuf import text_format
import yaml
+import subprocess
from six.moves.queue import Queue
@@ -71,7 +72,7 @@ class KiKiMRNode(daemon.Daemon, kikimr_node_interface.NodeInterface):
self.grpc_ssl_port = port_allocator.grpc_ssl_port
self.pgwire_port = port_allocator.pgwire_port
self.sqs_port = None
- if configurator.sqs_service_enabled:
+ if not configurator.simple_config and configurator.sqs_service_enabled:
self.sqs_port = port_allocator.sqs_port
self.__role = role
@@ -107,6 +108,44 @@ class KiKiMRNode(daemon.Daemon, kikimr_node_interface.NodeInterface):
daemon.Daemon.__init__(self, self.command, cwd=self.__working_dir, timeout=180, stderr_on_error_lines=240, **kwargs)
+ def is_port_listening(self, port):
+ """Check if the port is listening after node startup"""
+ try:
+ cmd = ["netstat", "-tuln"]
+ result = subprocess.run(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
+ output = result.stdout.decode()
+ port_lines = [line for line in output.split('\n') if str(port) in line]
+ if port_lines:
+ for line in port_lines:
+ logger.info(f"Port {port} status: {line.strip()}")
+ is_listening = True
+ else:
+ logger.info(f"Port {port} is not found in netstat output")
+ is_listening = False
+ return is_listening
+ except Exception as e:
+ logger.error(f"Error checking port {port}: {e}")
+ return False
+
+ def check_ports(self):
+ """Check if all allocated ports are listening"""
+ ports_status = {
+ "grpc_port": self.is_port_listening(self.grpc_port),
+ "mon_port": self.is_port_listening(self.mon_port),
+ "ic_port": self.is_port_listening(self.ic_port)
+ }
+
+ if hasattr(self, 'grpc_ssl_port') and self.grpc_ssl_port:
+ ports_status["grpc_ssl_port"] = self.is_port_listening(self.grpc_ssl_port)
+
+ if hasattr(self, 'pgwire_port') and self.pgwire_port:
+ ports_status["pgwire_port"] = self.is_port_listening(self.pgwire_port)
+
+ if hasattr(self, 'sqs_port') and self.sqs_port:
+ ports_status["sqs_port"] = self.is_port_listening(self.sqs_port)
+
+ return ports_status
+
@property
def cwd(self):
return self.__working_dir
@@ -319,6 +358,28 @@ class KiKiMR(kikimr_cluster_interface.KiKiMRClusterInterface):
))
raise
+ def __call_ydb_cli(self, cmd, token=None):
+ endpoint = 'grpc://{server}:{port}'.format(server=self.server, port=self.nodes[1].port)
+ full_command = [self.__configurator.get_ydb_cli_path(), '--endpoint', endpoint, '-y'] + cmd
+
+ env = None
+ token = token or self.__configurator.default_clusteradmin
+ if token is not None:
+ env = os.environ.copy()
+ env['YDB_TOKEN'] = token
+
+ logger.debug("Executing command = {}".format(full_command))
+ try:
+ return yatest.common.execute(full_command)
+ except yatest.common.ExecutionError as e:
+ logger.exception("KiKiMR command '{cmd}' failed with error: {e}\n\tstdout: {out}\n\tstderr: {err}".format(
+ cmd=" ".join(str(x) for x in full_command),
+ e=str(e),
+ out=e.execution_result.std_out,
+ err=e.execution_result.std_err
+ ))
+ raise
+
def start(self):
"""
Safely starts kikimr instance.
@@ -350,11 +411,15 @@ class KiKiMR(kikimr_cluster_interface.KiKiMRClusterInterface):
for node_id in self.__configurator.all_node_ids():
self.__run_node(node_id)
- bs_needed = 'blob_storage_config' in self.__configurator.yaml_config
+ if self.__configurator.use_self_management:
+ self.__cluster_bootstrap()
+
+ bs_needed = ('blob_storage_config' in self.__configurator.yaml_config) or self.__configurator.use_self_management
if bs_needed:
self.__wait_for_bs_controller_to_start()
- self.__add_bs_box()
+ if not self.__configurator.use_self_management:
+ self.__add_bs_box()
pools = {}
@@ -387,10 +452,11 @@ class KiKiMR(kikimr_cluster_interface.KiKiMRClusterInterface):
self._nodes[node_id].start()
return self._nodes[node_id]
- def __register_node(self):
+ def __register_node(self, configurator=None):
+ configurator = configurator or self.__configurator
node_index = next(self._node_index_allocator)
- if self.__configurator.separate_node_configs:
+ if configurator.separate_node_configs:
node_config_path = ensure_path_exists(
os.path.join(self.__config_base_path, "node_{}".format(node_index))
)
@@ -398,18 +464,18 @@ class KiKiMR(kikimr_cluster_interface.KiKiMRClusterInterface):
node_config_path = self.__config_path
data_center = None
- if isinstance(self.__configurator.dc_mapping, dict):
- if node_index in self.__configurator.dc_mapping:
- data_center = self.__configurator.dc_mapping[node_index]
+ if isinstance(configurator.dc_mapping, dict):
+ if node_index in configurator.dc_mapping:
+ data_center = configurator.dc_mapping[node_index]
self._nodes[node_index] = KiKiMRNode(
node_id=node_index,
config_path=node_config_path,
port_allocator=self.__port_allocator.get_node_port_allocator(node_index),
cluster_name=self.__cluster_name,
- configurator=self.__configurator,
+ configurator=configurator,
udfs_dir=self.__common_udfs_dir,
- tenant_affiliation=self.__configurator.yq_tenant,
- binary_path=self.__configurator.get_binary_path(node_index),
+ tenant_affiliation=configurator.yq_tenant,
+ binary_path=configurator.get_binary_path(node_index),
data_center=data_center,
)
return self._nodes[node_index]
@@ -498,19 +564,47 @@ class KiKiMR(kikimr_cluster_interface.KiKiMRClusterInterface):
node.stop()
node.start()
+ def prepare_node(self, configurator=None):
+ try:
+ new_node_object = self.__register_node(configurator)
+ self.__write_node_config(new_node_object.node_id, configurator)
+ logger.info(f"Successfully registered new node object with ID: {new_node_object.node_id}")
+ return new_node_object
+ except Exception as e:
+ logger.error(f"Failed to register new node: {e}", exc_info=True)
+ raise RuntimeError(f"Failed to register new node: {e}")
+
+ def start_node(self, node_id):
+ if node_id not in self._nodes:
+ logger.error(f"Cannot start node: Node ID {node_id} not found in registered nodes.")
+ raise KeyError(f"Node ID {node_id} not found.")
+
+ logger.info(f"Starting registered node {node_id}.")
+ try:
+ self._KiKiMR__run_node(node_id)
+ logger.info(f"Successfully started node {node_id}.")
+ except Exception as e:
+ raise RuntimeError(f"Failed to start node {node_id}: {e}")
+
@property
def config_path(self):
if self.__configurator.separate_node_configs:
return self.__config_base_path
return self.__config_path
+ def __write_node_config(self, node_id, configurator=None):
+ configurator = configurator or self.__configurator
+ node_config_path = ensure_path_exists(
+ os.path.join(self.__config_base_path, "node_{}".format(node_id))
+ )
+ logger.info(f"Writing node config to {node_config_path}")
+ logger.info(f"Config: {configurator.yaml_config}")
+ configurator.write_proto_configs(node_config_path)
+
def __write_configs(self):
if self.__configurator.separate_node_configs:
for node_id in self.__configurator.all_node_ids():
- node_config_path = ensure_path_exists(
- os.path.join(self.__config_base_path, "node_{}".format(node_id))
- )
- self.__configurator.write_proto_configs(node_config_path)
+ self.__write_node_config(node_id)
else:
self.__configurator.write_proto_configs(self.__config_path)
@@ -616,6 +710,30 @@ class KiKiMR(kikimr_cluster_interface.KiKiMRClusterInterface):
)
assert bs_controller_started
+ def __cluster_bootstrap(self):
+ timeout = 240
+ sleep = 5
+ retries, success = timeout / sleep, False
+ while retries > 0 and not success:
+ try:
+ self.__call_ydb_cli(
+ [
+ "admin",
+ "cluster",
+ "bootstrap",
+ "--uuid", "test-cluster"
+ ]
+ )
+ success = True
+
+ except Exception as e:
+ logger.error("Failed to execute, %s", str(e))
+ retries -= 1
+ time.sleep(sleep)
+
+ if retries == 0:
+ raise
+
class KikimrExternalNode(daemon.ExternalNodeDaemon, kikimr_node_interface.NodeInterface):
kikimr_binary_deploy_path = '/Berkanavt/kikimr/bin/kikimr'
diff --git a/ydb/tests/library/harness/param_constants.py b/ydb/tests/library/harness/param_constants.py
index 9a1244384cc..ed5834463fb 100644
--- a/ydb/tests/library/harness/param_constants.py
+++ b/ydb/tests/library/harness/param_constants.py
@@ -8,3 +8,10 @@ def kikimr_driver_path():
return yatest.common.binary_path(os.getenv("YDB_DRIVER_BINARY"))
return yatest.common.binary_path("kikimr/driver/kikimr")
+
+
+def ydb_cli_path():
+ if os.getenv("YDB_CLI_BINARY"):
+ return yatest.common.binary_path(os.getenv("YDB_CLI_BINARY"))
+
+ return yatest.common.binary_path("ydb/apps/ydb/ydb")