diff --git a/components/clp-package-utils/clp_package_utils/controller.py b/components/clp-package-utils/clp_package_utils/controller.py index 15a7672a0..fac77f2d1 100644 --- a/components/clp-package-utils/clp_package_utils/controller.py +++ b/components/clp-package-utils/clp_package_utils/controller.py @@ -35,7 +35,6 @@ GARBAGE_COLLECTOR_COMPONENT_NAME, LOG_INGESTOR_COMPONENT_NAME, MCP_SERVER_COMPONENT_NAME, - OrchestrationType, OTEL_COLLECTOR_COMPONENT_NAME, QUERY_JOBS_TABLE_NAME, QUERY_SCHEDULER_COMPONENT_NAME, @@ -45,9 +44,6 @@ REDIS_COMPONENT_NAME, REDUCER_COMPONENT_NAME, RESULTS_CACHE_COMPONENT_NAME, - SPIDER_DB_PASS_ENV_VAR_NAME, - SPIDER_DB_USER_ENV_VAR_NAME, - SPIDER_SCHEDULER_COMPONENT_NAME, StorageEngine, StorageType, WEBUI_COMPONENT_NAME, @@ -192,9 +188,6 @@ def _set_up_env_for_database(self) -> EnvVarsDict: env_vars |= { "CLP_DB_NAME": self._clp_config.database.names[ClpDbNameType.CLP], } - if self._clp_config.compression_scheduler.type == OrchestrationType.SPIDER: - env_vars["SPIDER_DB_NAME"] = self._clp_config.database.names[ClpDbNameType.SPIDER] - if BundledService.DATABASE not in self._clp_config.bundled: env_vars |= { "CLP_DB_CONNECT_PORT": str(self._clp_config.database.port), @@ -214,10 +207,8 @@ def _set_up_env_for_database(self) -> EnvVarsDict: env_vars |= { CLP_DB_PASS_ENV_VAR_NAME: credentials[ClpDbUserType.CLP].password, CLP_DB_ROOT_PASS_ENV_VAR_NAME: credentials[ClpDbUserType.ROOT].password, - SPIDER_DB_PASS_ENV_VAR_NAME: credentials[ClpDbUserType.SPIDER].password, CLP_DB_ROOT_USER_ENV_VAR_NAME: credentials[ClpDbUserType.ROOT].username, CLP_DB_USER_ENV_VAR_NAME: credentials[ClpDbUserType.CLP].username, - SPIDER_DB_USER_ENV_VAR_NAME: credentials[ClpDbUserType.SPIDER].username, } return env_vars @@ -383,32 +374,6 @@ def _set_up_env_for_redis(self) -> EnvVarsDict: return env_vars - def _set_up_env_for_spider_scheduler(self) -> EnvVarsDict: - """ - Sets up environment variables for the Spider scheduler component. - - :return: Dictionary of environment variables necessary to launch the component. - """ - component_name = SPIDER_SCHEDULER_COMPONENT_NAME - if self._clp_config.compression_scheduler.type != OrchestrationType.SPIDER: - logger.info( - "%s is not configured, skipping environment setup...", - component_name, - ) - return EnvVarsDict() - - logger.info("Setting up environment for %s...", component_name) - - env_vars = EnvVarsDict() - - # Connection config - env_vars |= { - "SPIDER_SCHEDULER_HOST": _get_ip_from_hostname(self._clp_config.spider_scheduler.host), - "SPIDER_SCHEDULER_PORT": str(self._clp_config.spider_scheduler.port), - } - - return env_vars - def _set_up_env_for_results_cache_bundling(self) -> EnvVarsDict: """ Sets up environment variables and directories for bundling the results cache component. @@ -1199,7 +1164,6 @@ def set_up_env(self) -> None: env_vars |= self._set_up_env_for_database() env_vars |= self._set_up_env_for_queue() env_vars |= self._set_up_env_for_redis() - env_vars |= self._set_up_env_for_spider_scheduler() env_vars |= self._set_up_env_for_results_cache() env_vars |= self._set_up_env_for_otel_collector() env_vars |= self._set_up_env_for_compression_scheduler() @@ -1231,11 +1195,10 @@ def start(self) -> None: should_compose_project_be_running=False, project_name=self._project_name ) - orchestration_type = self._clp_config.compression_scheduler.type - logger.info("Starting CLP using Docker Compose (%s orchestration)...", orchestration_type) + logger.info("Starting CLP using Docker Compose...") cmd = ["docker", "compose", "--project-name", self._project_name] - cmd += ["--file", self._get_docker_file_name()] + cmd += ["--file", "docker-compose.yaml"] cmd += ["up", "--detach", "--wait"] subprocess.run( cmd, @@ -1286,14 +1249,6 @@ def _get_num_workers() -> int: # This will change when we move from single to multi-container workers. See y-scope/clp#1424 return max(1, multiprocessing.cpu_count() // 2) - def _get_docker_file_name(self) -> str: - """ - :return: The Docker Compose file name to use based on the config. - """ - if self._clp_config.compression_scheduler.type == OrchestrationType.SPIDER: - return "docker-compose-spider.yaml" - return "docker-compose.yaml" - def _emit_topology_metrics(self) -> None: timestamp_ns = int(time.time() * 1e9) metrics = [] diff --git a/components/clp-package-utils/clp_package_utils/general.py b/components/clp-package-utils/clp_package_utils/general.py index b20de97c1..bedeebc09 100644 --- a/components/clp-package-utils/clp_package_utils/general.py +++ b/components/clp-package-utils/clp_package_utils/general.py @@ -516,8 +516,6 @@ def generate_credentials_file(credentials_file_path: pathlib.Path): "password": secrets.token_urlsafe(8), "root_username": "root", "root_password": secrets.token_urlsafe(8), - "spider_username": "spider-user", - "spider_password": secrets.token_urlsafe(8), }, QUEUE_COMPONENT_NAME: {"username": "clp-user", "password": secrets.token_urlsafe(8)}, REDIS_COMPONENT_NAME: {"password": secrets.token_urlsafe(16)}, diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index b1ac2396d..e0b5b04f7 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -33,7 +33,6 @@ DB_COMPONENT_NAME = "database" QUEUE_COMPONENT_NAME = "queue" REDIS_COMPONENT_NAME = "redis" -SPIDER_SCHEDULER_COMPONENT_NAME = "spider_scheduler" REDUCER_COMPONENT_NAME = "reducer" RESULTS_CACHE_COMPONENT_NAME = "results_cache" OTEL_COLLECTOR_COMPONENT_NAME = "otel-collector" @@ -83,8 +82,6 @@ CLP_QUEUE_USER_ENV_VAR_NAME = "CLP_QUEUE_USER" CLP_QUEUE_PASS_ENV_VAR_NAME = "CLP_QUEUE_PASS" CLP_REDIS_PASS_ENV_VAR_NAME = "CLP_REDIS_PASS" -SPIDER_DB_USER_ENV_VAR_NAME = "SPIDER_DB_USER" -SPIDER_DB_PASS_ENV_VAR_NAME = "SPIDER_DB_PASS" # Serializer StrEnumSerializer = PlainSerializer(serialize_str_enum) @@ -138,14 +135,6 @@ class DatabaseEngine(KebabCaseStrEnum): DatabaseEngineStr = Annotated[DatabaseEngine, StrEnumSerializer] -class OrchestrationType(KebabCaseStrEnum): - CELERY = auto() - SPIDER = auto() - - -OrchestrationTypeStr = Annotated[OrchestrationType, StrEnumSerializer] - - class QueryEngine(KebabCaseStrEnum): CLP = auto() CLP_S = auto() @@ -179,21 +168,18 @@ class ClpDbUserType(KebabCaseStrEnum): CLP = auto() ROOT = auto() - SPIDER = auto() class ClpDbNameType(KebabCaseStrEnum): """Database name types used by CLP components.""" CLP = auto() - SPIDER = auto() _DB_USER_TYPE_TO_DB_NAME_TYPE: MappingProxyType[ClpDbUserType, ClpDbNameType] = MappingProxyType( { ClpDbUserType.CLP: ClpDbNameType.CLP, ClpDbUserType.ROOT: ClpDbNameType.CLP, - ClpDbUserType.SPIDER: ClpDbNameType.SPIDER, } ) @@ -219,7 +205,6 @@ class Database(BaseModel): port: Port = DEFAULT_PORT names: dict[ClpDbNameType, NonEmptyStr] = { ClpDbNameType.CLP: "clp-db", - ClpDbNameType.SPIDER: "spider-db", } ssl_cert: NonEmptyStr | None = None auto_commit: bool = False @@ -339,10 +324,6 @@ def load_credentials_from_file(self, credentials_file_path: pathlib.Path): username=get_config_value(config, f"{DB_COMPONENT_NAME}.root_username"), password=get_config_value(config, f"{DB_COMPONENT_NAME}.root_password"), ) - self.credentials[ClpDbUserType.SPIDER] = DbUserCredentials( - username=get_config_value(config, f"{DB_COMPONENT_NAME}.spider_username"), - password=get_config_value(config, f"{DB_COMPONENT_NAME}.spider_password"), - ) except KeyError as ex: raise ValueError( f"Credentials file '{credentials_file_path}' does not contain key '{ex}'." @@ -362,9 +343,6 @@ def load_credentials_from_env(self, user_type: ClpDbUserType = ClpDbUserType.CLP elif user_type == ClpDbUserType.ROOT: user_env_var = CLP_DB_ROOT_USER_ENV_VAR_NAME pass_env_var = CLP_DB_ROOT_PASS_ENV_VAR_NAME - elif user_type == ClpDbUserType.SPIDER: - user_env_var = SPIDER_DB_USER_ENV_VAR_NAME - pass_env_var = SPIDER_DB_PASS_ENV_VAR_NAME else: err_msg = f"Unsupported user type '{user_type}'." raise ValueError(err_msg) @@ -380,24 +358,12 @@ def transform_for_container(self, is_bundled: bool): self.port = self.DEFAULT_PORT -class SpiderScheduler(BaseModel): - DEFAULT_PORT: ClassVar[int] = 6000 - - host: DomainStr = "localhost" - port: Port = DEFAULT_PORT - - def transform_for_container(self): - self.host = SPIDER_SCHEDULER_COMPONENT_NAME - self.port = self.DEFAULT_PORT - - class CompressionScheduler(BaseModel): UNLIMITED_CONCURRENT_TASKS_PER_JOB: ClassVar[NonNegativeInt] = 0 jobs_poll_delay: PositiveFloat = 0.1 # seconds max_concurrent_tasks_per_job: NonNegativeInt = UNLIMITED_CONCURRENT_TASKS_PER_JOB logging_level: LoggingLevel = "INFO" - type: OrchestrationTypeStr = OrchestrationType.CELERY telemetry_update_interval_ms: PositiveInt = 60000 @@ -850,7 +816,6 @@ class ClpConfig(BaseModel): results_cache: ResultsCache = ResultsCache() otel_collector: OtelCollector = OtelCollector() compression_scheduler: CompressionScheduler = CompressionScheduler() - spider_scheduler: SpiderScheduler | None = None query_scheduler: QueryScheduler | None = QueryScheduler() compression_worker: CompressionWorker = CompressionWorker() query_worker: QueryWorker | None = QueryWorker() @@ -1106,24 +1071,8 @@ def validate_query_engine_package_compatibility(self): return self - @model_validator(mode="after") - def validate_spider_config(self): - orchestration_type = self.compression_scheduler.type - if orchestration_type != OrchestrationType.SPIDER: - return self - if self.spider_scheduler is None: - raise ValueError( - "`spider_scheduler` must be configured when using Spider orchestration." - ) - if self.database.type != DatabaseEngine.MARIADB: - raise ValueError("Spider only supports MariaDB for the metadata database.") - return self - @model_validator(mode="after") def validate_celery_config(self): - orchestration_type = self.compression_scheduler.type - if orchestration_type != OrchestrationType.CELERY: - return self if self.queue is None: raise ValueError("`queue` must be configured when using Celery orchestration.") if self.redis is None: @@ -1149,8 +1098,6 @@ def transform_for_container(self): self.queue.transform_for_container(BundledService.QUEUE in self.bundled) if self.redis is not None: self.redis.transform_for_container(BundledService.REDIS in self.bundled) - if self.spider_scheduler is not None: - self.spider_scheduler.transform_for_container() self.results_cache.transform_for_container(BundledService.RESULTS_CACHE in self.bundled) self.otel_collector.transform_for_container(BundledService.OTEL_COLLECTOR in self.bundled) if self.query_scheduler is not None: diff --git a/components/clp-py-utils/clp_py_utils/create-db-tables.py b/components/clp-py-utils/clp_py_utils/create-db-tables.py index ce88d6fa9..0b5d29e4a 100644 --- a/components/clp-py-utils/clp_py_utils/create-db-tables.py +++ b/components/clp-py-utils/clp_py_utils/create-db-tables.py @@ -4,8 +4,7 @@ import subprocess import sys -from clp_py_utils.clp_config import ClpConfig, OrchestrationType, StorageEngine -from clp_py_utils.core import read_yaml_config_file +from clp_py_utils.clp_config import StorageEngine # Setup logging # Create logger @@ -50,23 +49,6 @@ def main(argv): # fmt: on subprocess.run(cmd, check=True) - try: - clp_config = ClpConfig.model_validate(read_yaml_config_file(pathlib.Path(config_file_path))) - clp_config.database.load_credentials_from_env() - if clp_config.compression_scheduler.type != OrchestrationType.SPIDER: - logger.info("No spider database configured. Skipping Spider database initialization.") - return 0 - except Exception as e: - logger.error(f"Failed to load CLP configuration: {e}") - return 1 - # fmt: off - cmd = [ - "python3", "-m", "clp_py_utils.initialize-spider-db", - "--config", str(config_file_path), - ] - # fmt: on - subprocess.run(cmd, check=True) - return 0 diff --git a/components/clp-py-utils/clp_py_utils/initialize-spider-db.py b/components/clp-py-utils/clp_py_utils/initialize-spider-db.py deleted file mode 100644 index f589ac28a..000000000 --- a/components/clp-py-utils/clp_py_utils/initialize-spider-db.py +++ /dev/null @@ -1,301 +0,0 @@ -#!/usr/bin/env python3 - -"""Script to initialize Spider database.""" - -import argparse -import logging -import pathlib -import re -import sys -from contextlib import closing - -from pydantic import ValidationError - -from clp_py_utils.clp_config import ClpConfig, ClpDbNameType, ClpDbUserType -from clp_py_utils.core import read_yaml_config_file -from clp_py_utils.sql_adapter import SqlAdapter - -# Setup logging -# Create logger -logger = logging.getLogger("initialize-spider-db") -logger.setLevel(logging.INFO) -# Setup console logging -logging_console_handler = logging.StreamHandler() -logging_formatter = logging.Formatter("%(asctime)s [%(levelname)s] %(message)s") -logging_console_handler.setFormatter(logging_formatter) -logger.addHandler(logging_console_handler) - - -table_creators = [ - """ -CREATE TABLE IF NOT EXISTS `drivers` -( - `id` BINARY(16) NOT NULL, - `heartbeat` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, - PRIMARY KEY (`id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `schedulers` -( - `id` BINARY(16) NOT NULL, - `address` VARCHAR(40) NOT NULL, - `port` INT UNSIGNED NOT NULL, - CONSTRAINT `scheduler_driver_id` FOREIGN KEY (`id`) REFERENCES `drivers` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - PRIMARY KEY (`id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS jobs -( - `id` BINARY(16) NOT NULL, - `client_id` BINARY(16) NOT NULL, - `creation_time` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, - `state` ENUM ('running', 'success', 'fail', 'cancel') NOT NULL DEFAULT 'running', - KEY (`client_id`) USING BTREE, - INDEX idx_jobs_creation_time (`creation_time`), - INDEX idx_jobs_state (`state`), - PRIMARY KEY (`id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS tasks -( - `id` BINARY(16) NOT NULL, - `job_id` BINARY(16) NOT NULL, - `func_name` VARCHAR(64) NOT NULL, - `language` ENUM('cpp', 'python') NOT NULL, - `state` ENUM ('pending', 'ready', 'running', 'success', 'cancel', 'fail') NOT NULL, - `timeout` FLOAT, - `max_retry` INT UNSIGNED DEFAULT 0, - `retry` INT UNSIGNED DEFAULT 0, - `instance_id` BINARY(16), - CONSTRAINT `task_job_id` FOREIGN KEY (`job_id`) REFERENCES `jobs` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - PRIMARY KEY (`id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS input_tasks -( - `job_id` BINARY(16) NOT NULL, - `task_id` BINARY(16) NOT NULL, - `position` INT UNSIGNED NOT NULL, - CONSTRAINT `input_task_job_id` FOREIGN KEY (`job_id`) REFERENCES `jobs` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `input_task_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - INDEX (`job_id`, `position`), - PRIMARY KEY (`task_id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS output_tasks -( - `job_id` BINARY(16) NOT NULL, - `task_id` BINARY(16) NOT NULL, - `position` INT UNSIGNED NOT NULL, - CONSTRAINT `output_task_job_id` FOREIGN KEY (`job_id`) REFERENCES `jobs` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `output_task_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - INDEX (`job_id`, `position`), - PRIMARY KEY (`task_id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `data` -( - `id` BINARY(16) NOT NULL, - `value` BLOB(60000) NOT NULL, - `hard_locality` BOOL DEFAULT FALSE, - `persisted` BOOL DEFAULT FALSE, - PRIMARY KEY (`id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `task_outputs` -( - `task_id` BINARY(16) NOT NULL, - `position` INT UNSIGNED NOT NULL, - `type` VARCHAR(999) NOT NULL, - `value` BLOB(60000), - `data_id` BINARY(16), - CONSTRAINT `output_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `output_data_id` FOREIGN KEY (`data_id`) REFERENCES `data` (`id`) ON UPDATE NO ACTION ON DELETE NO ACTION, - PRIMARY KEY (`task_id`, `position`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `task_inputs` -( - `task_id` BINARY(16) NOT NULL, - `position` INT UNSIGNED NOT NULL, - `type` VARCHAR(999) NOT NULL, - `output_task_id` BINARY(16), - `output_task_position` INT UNSIGNED, - `value` BLOB(60000), - `data_id` BINARY(16), - CONSTRAINT `input_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `input_task_output_match` FOREIGN KEY (`output_task_id`, `output_task_position`) REFERENCES task_outputs (`task_id`, `position`) ON UPDATE NO ACTION ON DELETE SET NULL, - CONSTRAINT `input_data_id` FOREIGN KEY (`data_id`) REFERENCES `data` (`id`) ON UPDATE NO ACTION ON DELETE NO ACTION, - PRIMARY KEY (`task_id`, `position`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `task_dependencies` -( - `parent` BINARY(16) NOT NULL, - `child` BINARY(16) NOT NULL, - KEY (`parent`) USING BTREE, - KEY (`child`) USING BTREE, - CONSTRAINT `task_dep_parent` FOREIGN KEY (`parent`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `task_dep_child` FOREIGN KEY (`child`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE -); -""", - """ -CREATE TABLE IF NOT EXISTS `task_instances` -( - `id` BINARY(16) NOT NULL, - `task_id` BINARY(16) NOT NULL, - `start_time` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, - CONSTRAINT `instance_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - PRIMARY KEY (`id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `scheduler_leases` -( - `scheduler_id` BINARY(16) NOT NULL, - `task_id` BINARY(16) NOT NULL, - `lease_time` TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, - CONSTRAINT `lease_scheduler_id` FOREIGN KEY (`scheduler_id`) REFERENCES `schedulers` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `lease_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - INDEX (`scheduler_id`), - PRIMARY KEY (`scheduler_id`, `task_id`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `data_locality` -( - `id` BINARY(16) NOT NULL, - `address` VARCHAR(40) NOT NULL, - KEY (`id`) USING BTREE, - CONSTRAINT `locality_data_id` FOREIGN KEY (`id`) REFERENCES `data` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE -); -""", - """ -CREATE TABLE IF NOT EXISTS `data_ref_driver` -( - `id` BINARY(16) NOT NULL, - `driver_id` BINARY(16) NOT NULL, - KEY (`id`) USING BTREE, - KEY (`driver_id`) USING BTREE, - CONSTRAINT `data_driver_ref_id` FOREIGN KEY (`id`) REFERENCES `data` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `data_ref_driver_id` FOREIGN KEY (`driver_id`) REFERENCES `drivers` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE -); -""", - """ -CREATE TABLE IF NOT EXISTS `data_ref_task` -( - `id` BINARY(16) NOT NULL, - `task_id` BINARY(16) NOT NULL, - KEY (`id`) USING BTREE, - KEY (`task_id`) USING BTREE, - CONSTRAINT `data_task_ref_id` FOREIGN KEY (`id`) REFERENCES `data` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE, - CONSTRAINT `data_ref_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE -); -""", - """ -CREATE TABLE IF NOT EXISTS `client_kv_data` -( - `kv_key` VARCHAR(64) NOT NULL, - `value` BLOB(60000) NOT NULL, - `client_id` BINARY(16) NOT NULL, - PRIMARY KEY (`client_id`, `kv_key`) -); -""", - """ -CREATE TABLE IF NOT EXISTS `task_kv_data` -( - `kv_key` VARCHAR(64) NOT NULL, - `value` BLOB(60000) NOT NULL, - `task_id` BINARY(16) NOT NULL, - PRIMARY KEY (`task_id`, `kv_key`), - CONSTRAINT `kv_data_task_id` FOREIGN KEY (`task_id`) REFERENCES `tasks` (`id`) ON UPDATE NO ACTION ON DELETE CASCADE -); -""", -] - - -def main(argv: list[str]) -> int: - args_parser = argparse.ArgumentParser(description="Sets up Spider database.") - args_parser.add_argument("--config", "-c", required=True, help="CLP configuration file.") - parsed_args = args_parser.parse_args(argv[1:]) - - config_path = pathlib.Path(parsed_args.config) - try: - clp_config = ClpConfig.model_validate(read_yaml_config_file(config_path)) - clp_config.database.load_credentials_from_env(user_type=ClpDbUserType.CLP) - clp_config.database.load_credentials_from_env(user_type=ClpDbUserType.ROOT) - clp_config.database.load_credentials_from_env(user_type=ClpDbUserType.SPIDER) - except (ValidationError, ValueError): - logger.exception("Invalid CLP configuration.") - return -1 - except Exception: - logger.exception("Failed to load CLP configuration.") - return -1 - - try: - sql_adapter = SqlAdapter(clp_config.database) - with ( - closing(sql_adapter.create_connection(user_type=ClpDbUserType.ROOT)) as db_conn, - closing(db_conn.cursor()) as db_cursor, - ): - clp_db_user = clp_config.database.credentials[ClpDbUserType.CLP].username - spider_db_name = clp_config.database.names[ClpDbNameType.SPIDER] - spider_db_user = clp_config.database.credentials[ClpDbUserType.SPIDER].username - spider_db_password = clp_config.database.credentials[ClpDbUserType.SPIDER].password - if not _validate_name(spider_db_name): - logger.exception("Invalid database name: %s.", spider_db_name) - return -1 - if not _validate_name(spider_db_user): - logger.exception("Invalid database user name: %s.", spider_db_user) - return -1 - if not _validate_name(clp_db_user): - logger.exception("Invalid CLP database user name: %s.", clp_db_user) - return -1 - - db_cursor.execute(f"""CREATE DATABASE IF NOT EXISTS `{spider_db_name}`""") - if spider_db_password is None: - logger.exception("Password must be set for Spider database user.") - return -1 - db_cursor.execute( - f"""CREATE USER IF NOT EXISTS '{spider_db_user}'@'%' IDENTIFIED BY '{spider_db_password}'""" - ) - db_cursor.execute( - f"""GRANT ALL PRIVILEGES ON `{spider_db_name}`.* TO '{spider_db_user}'@'%'""" - ) - db_cursor.execute( - f"""GRANT ALL PRIVILEGES ON `{spider_db_name}`.* TO '{clp_db_user}'@'%'""" - ) - - db_cursor.execute(f"""USE `{spider_db_name}`""") - for table_creator in table_creators: - db_cursor.execute(table_creator) - except Exception: - logger.exception("Failed to setup Spider database.") - return -1 - - return 0 - - -_name_pattern = re.compile(r"^[A-Za-z0-9_-]+$") - - -def _validate_name(name: str) -> bool: - """ - Validates that the input string contains only alphanumeric characters, underscores, or hyphens. - :param name: The input string to validate. - :return: If the input string is valid. - """ - return _name_pattern.match(name) is not None - - -if "__main__" == __name__: - sys.exit(main(sys.argv)) diff --git a/components/clp-py-utils/clp_py_utils/sql_adapter.py b/components/clp-py-utils/clp_py_utils/sql_adapter.py index 2fe080d98..d579d0df0 100644 --- a/components/clp-py-utils/clp_py_utils/sql_adapter.py +++ b/components/clp-py-utils/clp_py_utils/sql_adapter.py @@ -10,7 +10,7 @@ from sqlalchemy import pool from sqlalchemy.dialects.mysql import mariadbconnector, mysqlconnector -from clp_py_utils.clp_config import ClpDbUserType, Database, DatabaseEngine +from clp_py_utils.clp_config import ClpDbNameType, ClpDbUserType, Database, DatabaseEngine class DummyCloseableObject: @@ -139,7 +139,8 @@ def _create_mysql_connection( logging.exception("Database access denied.") elif err.errno == errorcode.ER_BAD_DB_ERROR: logging.exception( - f'Specified database "{self.database_config.name}" does not exist.' + 'Specified database "%s" does not exist.', + self.database_config.names[ClpDbNameType.CLP], ) else: logging.exception(err) diff --git a/components/clp-rust-utils/src/clp_config/package/config.rs b/components/clp-rust-utils/src/clp_config/package/config.rs index badacf380..d445ba6c5 100644 --- a/components/clp-rust-utils/src/clp_config/package/config.rs +++ b/components/clp-rust-utils/src/clp_config/package/config.rs @@ -118,14 +118,12 @@ impl Default for SpiderTaskExecutorConfig { #[serde(deny_unknown_fields)] pub struct ClpDbNames { pub clp: String, - pub spider: String, } impl Default for ClpDbNames { fn default() -> Self { Self { clp: "clp-db".to_owned(), - spider: "spider-db".to_owned(), } } } @@ -704,7 +702,6 @@ mod tests { "port": 3306, "names": { "clp": "clp-db", - "spider": "spider-db", }, "table_prefix": "custom_" }); diff --git a/components/clp-tdl-package/src/task/compression/compress.rs b/components/clp-tdl-package/src/task/compression/compress.rs index eda4c0415..a1ccc35b7 100644 --- a/components/clp-tdl-package/src/task/compression/compress.rs +++ b/components/clp-tdl-package/src/task/compression/compress.rs @@ -1136,7 +1136,6 @@ mod tests { port: 3306, names: ClpDbNames { clp: "clp-db".to_string(), - spider: "spider-db".to_string(), }, table_prefix: "clp_".to_string(), }; diff --git a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py index 8093c6c1f..b41755642 100644 --- a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py +++ b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py @@ -14,11 +14,9 @@ from clp_package_utils.general import CONTAINER_INPUT_LOGS_ROOT_DIR from clp_py_utils.clp_config import ( ClpConfig, - ClpDbUserType, COMPRESSION_JOBS_TABLE_NAME, COMPRESSION_SCHEDULER_COMPONENT_NAME, COMPRESSION_TASKS_TABLE_NAME, - OrchestrationType, StorageEngine, ) from clp_py_utils.clp_logging import configure_logging, get_logger @@ -40,7 +38,6 @@ from job_orchestration.scheduler.compress.partition import PathsToCompressBuffer from job_orchestration.scheduler.compress.task_manager.celery_task_manager import CeleryTaskManager -from job_orchestration.scheduler.compress.task_manager.spider_task_manager import SpiderTaskManager from job_orchestration.scheduler.compress.task_manager.task_manager import TaskManager from job_orchestration.scheduler.constants import ( CompressionJobStatus, @@ -636,19 +633,7 @@ def main(argv) -> int | None: atexit.register(shutdown_telemetry) sql_adapter = SqlAdapter(clp_config.database) - task_manager: CeleryTaskManager | SpiderTaskManager - if clp_config.compression_scheduler.type == OrchestrationType.CELERY: - task_manager = CeleryTaskManager() - elif clp_config.compression_scheduler.type == OrchestrationType.SPIDER: - clp_config.database.load_credentials_from_env(ClpDbUserType.SPIDER) - task_manager = SpiderTaskManager( - clp_config.database.get_container_url(ClpDbUserType.SPIDER) - ) - else: - logger.error( - f"Unsupported compression scheduler type: {clp_config.compression_scheduler.type}" - ) - return -1 + task_manager = CeleryTaskManager() try: killed_jobs = kill_hanging_jobs(sql_adapter, SchedulerType.COMPRESSION) diff --git a/components/package-template/src/etc/clp-config.template.json.yaml b/components/package-template/src/etc/clp-config.template.json.yaml index 63c8af433..923447f43 100644 --- a/components/package-template/src/etc/clp-config.template.json.yaml +++ b/components/package-template/src/etc/clp-config.template.json.yaml @@ -50,19 +50,12 @@ telemetry: # port: 3306 # names: # clp: "clp-db" -# spider: "spider-db" # #compression_scheduler: # jobs_poll_delay: 0.1 # seconds # max_concurrent_tasks_per_job: 0 # A value of 0 disables the limit # logging_level: "INFO" # telemetry_update_interval_ms: 60000 -# type: "celery" # "celery" or "spider" -# -#spider_scheduler: -# host: "localhost" -# port: 6000 -# logging_level: "INFO" # #query_scheduler: # host: "localhost" diff --git a/components/package-template/src/etc/clp-config.template.text.yaml b/components/package-template/src/etc/clp-config.template.text.yaml index 9241b493a..2d288692f 100644 --- a/components/package-template/src/etc/clp-config.template.text.yaml +++ b/components/package-template/src/etc/clp-config.template.text.yaml @@ -49,18 +49,11 @@ log_ingestor: null # port: 3306 # names: # clp: "clp-db" -# spider: "spider-db" # #compression_scheduler: # jobs_poll_delay: 0.1 # seconds # logging_level: "INFO" # telemetry_update_interval_ms: 60000 -# type: "celery" # "celery" or "spider" -# -#spider_scheduler: -# host: "localhost" -# port: 6000 -# logging_level: "INFO" # #query_scheduler: # host: "localhost" diff --git a/components/package-template/src/etc/credentials.template.yaml b/components/package-template/src/etc/credentials.template.yaml index fa1bfd473..e4259600b 100644 --- a/components/package-template/src/etc/credentials.template.yaml +++ b/components/package-template/src/etc/credentials.template.yaml @@ -3,8 +3,6 @@ # password: "pass" # root_username: "root" # root_password: "root-pass" -# spider_username: "spider-user" -# spider_password: "spider-pass" # #queue: # username: "clp-user" diff --git a/docs/src/user-docs/guides-docker-compose-deployment.md b/docs/src/user-docs/guides-docker-compose-deployment.md index 0526934bb..d431cb6d3 100755 --- a/docs/src/user-docs/guides-docker-compose-deployment.md +++ b/docs/src/user-docs/guides-docker-compose-deployment.md @@ -165,13 +165,13 @@ docker compose \ up db-table-creator \ --no-deps -# Start queue (if using Celery) +# Start queue docker compose \ --project-name "clp-package-$(cat var/log/instance-id)" \ up queue \ --no-deps --wait -# Start redis (if using Celery) +# Start redis docker compose \ --project-name "clp-package-$(cat var/log/instance-id)" \ up redis \ diff --git a/tools/deployment/package-helm/Chart.yaml b/tools/deployment/package-helm/Chart.yaml index f96fd11bc..1bb17bc37 100644 --- a/tools/deployment/package-helm/Chart.yaml +++ b/tools/deployment/package-helm/Chart.yaml @@ -1,6 +1,6 @@ apiVersion: "v2" name: "clp" -version: "0.4.1-dev.4" +version: "0.4.1-dev.5" description: "A Helm chart for CLP's (Compressed Log Processor) package deployment" type: "application" appVersion: "0.13.1-dev" diff --git a/tools/deployment/package-helm/templates/configmap.yaml b/tools/deployment/package-helm/templates/configmap.yaml index 48d7aeef8..012982dcf 100644 --- a/tools/deployment/package-helm/templates/configmap.yaml +++ b/tools/deployment/package-helm/templates/configmap.yaml @@ -117,7 +117,6 @@ data: host: "{{ include "clp.databaseHost" . }}" names: clp: {{ .Values.clpConfig.database.names.clp | quote }} - spider: {{ .Values.clpConfig.database.names.spider | quote }} port: {{ include "clp.databasePort" . | int }} ssl_cert: null type: {{ .Values.clpConfig.database.type | quote }} diff --git a/tools/deployment/package-helm/values.yaml b/tools/deployment/package-helm/values.yaml index 4e38dc3ae..6c897d701 100644 --- a/tools/deployment/package-helm/values.yaml +++ b/tools/deployment/package-helm/values.yaml @@ -191,7 +191,6 @@ clpConfig: port: 30306 names: clp: "clp-db" - spider: "spider-db" compression_coordinator: commit_task_hard_timeout_secs: 60 diff --git a/tools/deployment/package/docker-compose-all.yaml b/tools/deployment/package/docker-compose-all.yaml index 7942171d8..bb833e56c 100644 --- a/tools/deployment/package/docker-compose-all.yaml +++ b/tools/deployment/package/docker-compose-all.yaml @@ -109,8 +109,6 @@ services: CLP_DB_ROOT_USER: "${CLP_DB_ROOT_USER:-root}" CLP_DB_USER: "${CLP_DB_USER:-clp-user}" PYTHONPATH: "/opt/clp/lib/python3/site-packages" - SPIDER_DB_PASS: "${SPIDER_DB_PASS:?Please set a value.}" - SPIDER_DB_USER: "${SPIDER_DB_USER:-spider-user}" volumes: - *volume_clp_config_readonly depends_on: @@ -245,28 +243,6 @@ services: "--stream-collection", "${CLP_RESULTS_CACHE_STREAM_COLLECTION_NAME:-stream-files}", ] - spider-scheduler: - <<: *service_defaults - hostname: "spider_scheduler" - environment: - SPIDER_LOG_FILE: "/var/log/spider_scheduler.log" - depends_on: - db-table-creator: - condition: "service_completed_successfully" - ports: - - host_ip: "${SPIDER_SCHEDULER_HOST:-127.0.0.1}" - published: "${SPIDER_SCHEDULER_PORT:-6000}" - target: 6000 - command: [ - "/opt/clp/bin/spider_scheduler", - "--host", "spider_scheduler", - "--port", "6000", - "--storage_url", "jdbc:mariadb://database:${CLP_DB_CONNECT_PORT:-3306}/\ - ${SPIDER_DB_NAME:-spider-db}?\ - user=${SPIDER_DB_USER:-spider-user}\ - &password=${SPIDER_DB_PASS:?Please set a value.}", - ] - compression-scheduler: <<: *service_defaults hostname: "compression_scheduler" @@ -341,38 +317,6 @@ services: "-n", "compression-worker@%h" ] - spider-compression-worker: - <<: *service_defaults - hostname: "compression_worker" - environment: - CLP_CONFIG_PATH: "/etc/clp-config.yaml" - CLP_HOME: "/opt/clp" - CLP_LOGGING_LEVEL: "${CLP_COMPRESSION_WORKER_LOGGING_LEVEL:-INFO}" - CLP_LOGS_DIR: "/var/log/compression_worker" - PYTHONPATH: "/opt/clp/lib/python3/site-packages" - SPIDER_LOG_DIR: "/var/log/compression_worker" - volumes: - - *volume_clp_config_readonly - - *volume_clp_logs - - *volume_clp_tmp - - "${CLP_ARCHIVE_OUTPUT_DIR_HOST:-empty}:/var/data/archives" - - "${CLP_AWS_CONFIG_DIR_HOST:-empty}:/opt/clp/.aws:ro" - - "${CLP_LOGS_INPUT_DIR_HOST:-empty}:${CLP_LOGS_INPUT_DIR_CONTAINER:-/mnt/logs}" - - "${CLP_STAGED_ARCHIVE_OUTPUT_DIR_HOST:-empty}:/var/data/staged-archives" - depends_on: - db-table-creator: - condition: "service_completed_successfully" - command: [ - "python3", "-u", - "-m", "job_orchestration.executor.start-spider-worker", - "--host", "compression_worker", - "--num-workers", "${CLP_COMPRESSION_WORKER_CONCURRENCY:-1}", - "--storage-url", "jdbc:mariadb://database:${CLP_DB_CONNECT_PORT:-3306}/\ - ${SPIDER_DB_NAME:-spider-db}?\ - user=${SPIDER_DB_USER:-spider-user}\ - &password=${SPIDER_DB_PASS:?Please set a value.}", - ] - webui: <<: *service_defaults hostname: "webui" diff --git a/tools/deployment/package/docker-compose-spider.yaml b/tools/deployment/package/docker-compose-spider.yaml deleted file mode 100644 index 4b7c20365..000000000 --- a/tools/deployment/package/docker-compose-spider.yaml +++ /dev/null @@ -1,92 +0,0 @@ -name: "clp-package-spider" - -include: - - "docker-compose.utils.yaml" - -services: - database: - extends: - file: "docker-compose-all.yaml" - service: "database" - - db-table-creator: - extends: - file: "docker-compose-all.yaml" - service: "db-table-creator" - - spider-scheduler: - extends: - file: "docker-compose-all.yaml" - service: "spider-scheduler" - - compression-scheduler: - extends: - file: "docker-compose-all.yaml" - service: "compression-scheduler" - # Override: Spider does NOT require `queue` or `redis`. - depends_on: - db-table-creator: - condition: "service_completed_successfully" - environment: - SPIDER_DB_PASS: "${SPIDER_DB_PASS:?Please set a value.}" - SPIDER_DB_USER: "${SPIDER_DB_USER:?Please set a value.}" - - compression-worker: - extends: - file: "docker-compose-all.yaml" - service: "spider-compression-worker" - - webui: - extends: - file: "docker-compose-all.yaml" - service: "webui" - - garbage-collector: - extends: - file: "docker-compose-all.yaml" - service: "garbage-collector" - - queue: - extends: - file: "docker-compose-all.yaml" - service: "queue" - - redis: - extends: - file: "docker-compose-all.yaml" - service: "redis" - - results-cache: - extends: - file: "docker-compose-all.yaml" - service: "results-cache" - - results-cache-indices-creator: - extends: - file: "docker-compose-all.yaml" - service: "results-cache-indices-creator" - - query-scheduler: - extends: - file: "docker-compose-all.yaml" - service: "query-scheduler" - - query-worker: - extends: - file: "docker-compose-all.yaml" - service: "query-worker" - - reducer: - extends: - file: "docker-compose-all.yaml" - service: "reducer" - - log-ingestor: - extends: - file: "docker-compose-all.yaml" - service: "log-ingestor" - - otel-collector: - extends: - file: "docker-compose-all.yaml" - service: "otel-collector"