diff --git a/components/clp-package-utils/clp_package_utils/controller.py b/components/clp-package-utils/clp_package_utils/controller.py index eba003784c..c39f50aaca 100644 --- a/components/clp-package-utils/clp_package_utils/controller.py +++ b/components/clp-package-utils/clp_package_utils/controller.py @@ -10,6 +10,7 @@ import time import uuid from abc import ABC, abstractmethod +from collections.abc import Callable from typing import Any from clp_py_utils.clp_config import ( @@ -26,9 +27,11 @@ ClpConfig, ClpDbNameType, ClpDbUserType, + COMPRESSION_COORDINATOR_COMPONENT_NAME, COMPRESSION_JOBS_TABLE_NAME, COMPRESSION_SCHEDULER_COMPONENT_NAME, COMPRESSION_WORKER_COMPONENT_NAME, + CompressionOrchestration, CONTAINER_INPUT_LOGS_ROOT_DIR, DatabaseEngine, DB_COMPONENT_NAME, @@ -44,6 +47,9 @@ REDIS_COMPONENT_NAME, REDUCER_COMPONENT_NAME, RESULTS_CACHE_COMPONENT_NAME, + SPIDER_COMPONENT_NAME, + SpiderScheduler, + SpiderSchedulerPolicy, StorageEngine, StorageType, WEBUI_COMPONENT_NAME, @@ -70,6 +76,7 @@ validate_queue_config, validate_redis_config, validate_results_cache_config, + validate_spider_config, validate_webui_config, ) @@ -523,6 +530,125 @@ def _set_up_env_for_query_scheduler(self) -> EnvVarsDict: return env_vars + def _set_up_env_for_spider(self) -> EnvVarsDict: + """ + Sets up environment variables for Spider's components. + + :return: Dictionary of environment variables necessary to launch the components. + """ + component_name = SPIDER_COMPONENT_NAME + spider_config = self._clp_config.spider + if ( + CompressionOrchestration.SPIDER != self._clp_config.package.scheduler + or spider_config is None + ): + logger.info("%s is not configured, skipping environment setup...", component_name) + return EnvVarsDict() + + logger.info("Setting up environment for %s...", component_name) + + validate_spider_config(self._clp_config) + + storage = spider_config.storage + scheduler = spider_config.scheduler + worker = spider_config.worker + + env_vars = EnvVarsDict() + + # Connection config + env_vars |= { + "SPIDER_STORAGE_PUBLISHED_IP": _get_ip_from_hostname(spider_config.host), + "SPIDER_STORAGE_PUBLISHED_PORT": str(spider_config.port), + } + + # Storage config + env_vars |= { + "SPIDER_STORAGE_LOG_LEVEL": storage.log_level, + "SPIDER_STORAGE_DB_MAX_CONNECTIONS": _str_or_none(storage.db_max_connections), + "SPIDER_STORAGE_INBOUND_QUEUE_CLEANUP_CAPACITY": _str_or_none( + storage.inbound_queue.cleanup_capacity + ), + "SPIDER_STORAGE_INBOUND_QUEUE_COMMIT_CAPACITY": _str_or_none( + storage.inbound_queue.commit_capacity + ), + "SPIDER_STORAGE_INBOUND_QUEUE_TASK_CAPACITY": _str_or_none( + storage.inbound_queue.task_capacity + ), + "SPIDER_STORAGE_JOB_CACHE_GC_INTERVAL_SEC": _str_or_none( + storage.job_cache_gc.gc_interval_sec + ), + "SPIDER_STORAGE_JOB_CACHE_TERMINATED_JOB_RETENTION_SEC": _str_or_none( + storage.job_cache_gc.terminated_job_retention_sec + ), + "SPIDER_STORAGE_TASK_INSTANCE_POOL_EXECUTION_MANAGER_STALE_CUTOFF_SEC": _str_or_none( + storage.task_instance_pool.execution_manager_stale_cutoff_sec + ), + "SPIDER_STORAGE_TASK_INSTANCE_POOL_GC_INTERVAL_SEC": _str_or_none( + storage.task_instance_pool.gc_interval_sec + ), + "SPIDER_STORAGE_TASK_INSTANCE_POOL_MESSAGE_CHANNEL_CAPACITY": _str_or_none( + storage.task_instance_pool.message_channel_capacity + ), + } + + # Scheduler config + env_vars |= { + "SPIDER_SCHEDULER_LOG_LEVEL": scheduler.log_level, + "SPIDER_SCHEDULER_POLICY": scheduler.policy, + "SPIDER_SCHEDULER_CONNECTION_POOL_SIZE": _str_or_none(scheduler.connection_pool_size), + "SPIDER_SCHEDULER_STOP_TIMEOUT_SEC": _str_or_none(scheduler.stop_timeout_sec), + "SPIDER_SCHEDULER_EM_REGISTRY_DEAD_EM_CUTOFF_SEC": _str_or_none( + scheduler.em_registry.dead_em_cutoff_sec + ), + "SPIDER_SCHEDULER_EM_REGISTRY_LIVENESS_TRACKING_INTERVAL_MS": _str_or_none( + scheduler.em_registry.liveness_tracking_interval_ms + ), + } + env_vars |= _get_spider_scheduler_policy_env_vars(scheduler) + + # Worker config + env_vars |= { + "SPIDER_WORKER_REPLICAS": _str_or_none(worker.replicas), + "SPIDER_WORKER_LOG_LEVEL": worker.log_level, + "SPIDER_WORKER_CONNECTION_POOL_SIZE": _str_or_none(worker.connection_pool_size), + "SPIDER_WORKER_SCHEDULER_POLL_WAIT_MS": _str_or_none(worker.scheduler_poll_wait_ms), + "SPIDER_WORKER_MAX_LOG_LINE_BYTES": _str_or_none(worker.max_log_line_bytes), + "SPIDER_WORKER_LIVENESS_SCHEDULER_HEARTBEAT_INTERVAL_SEC": _str_or_none( + worker.liveness.scheduler_heartbeat_interval_sec + ), + "SPIDER_WORKER_LIVENESS_STORAGE_HEARTBEAT_INTERVAL_SEC": _str_or_none( + worker.liveness.storage_heartbeat_interval_sec + ), + } + + return env_vars + + def _set_up_env_for_compression_coordinator(self) -> EnvVarsDict: + """ + Sets up environment variables for the compression coordinator component. + + :return: Dictionary of environment variables necessary to launch the component. + """ + component_name = COMPRESSION_COORDINATOR_COMPONENT_NAME + coordinator_config = self._clp_config.compression_coordinator + if ( + CompressionOrchestration.SPIDER != self._clp_config.package.scheduler + or coordinator_config is None + ): + 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() + + # Runtime config + env_vars |= { + "CLP_COMPRESSION_COORDINATOR_LOGGING_LEVEL": coordinator_config.logging_level, + } + + return env_vars + def _set_up_env_for_compression_worker(self, num_workers: int) -> EnvVarsDict: """ Sets up environment variables for the compression worker component. @@ -1097,6 +1223,7 @@ def set_up_env(self) -> None: # Restart Policy env_vars |= { "CLP_RESTART_POLICY": self._restart_policy, + "SPIDER_RESTART_POLICY": self._restart_policy, } # Credentials @@ -1132,9 +1259,8 @@ def set_up_env(self) -> None: env_vars |= self._set_up_env_for_telemetry() # Paths - aws_config_dir = self._clp_config.aws_config_directory env_vars |= { - "CLP_AWS_CONFIG_DIR_HOST": (None if aws_config_dir is None else str(aws_config_dir)), + "CLP_AWS_CONFIG_DIR_HOST": _str_or_none(self._clp_config.aws_config_directory), "CLP_DATA_DIR_HOST": str(self._clp_config.data_directory), "CLP_LOGS_DIR_HOST": str(self._clp_config.logs_directory), "CLP_TMP_DIR_HOST": str(self._clp_config.tmp_directory), @@ -1171,6 +1297,8 @@ def set_up_env(self) -> None: 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() + env_vars |= self._set_up_env_for_compression_coordinator() + env_vars |= self._set_up_env_for_spider() env_vars |= self._set_up_env_for_query_scheduler() env_vars |= self._set_up_env_for_compression_worker(num_workers) env_vars |= self._set_up_env_for_query_worker(num_workers) @@ -1202,7 +1330,7 @@ def start(self) -> None: logger.info("Starting CLP using Docker Compose...") cmd = ["docker", "compose", "--project-name", self._project_name] - cmd += ["--file", "docker-compose.yaml"] + cmd += ["--file", self._get_docker_file_name()] cmd += ["up", "--detach", "--wait"] subprocess.run( cmd, @@ -1239,12 +1367,29 @@ def stop(self) -> None: logger.info("Stopping all CLP containers using Docker Compose...") subprocess.run( - ["docker", "compose", "--project-name", self._project_name, "down"], + [ + "docker", + "compose", + "--project-name", + self._project_name, + "--file", + self._get_docker_file_name(), + "down", + "--remove-orphans", + ], cwd=self._clp_home, check=True, ) logger.info("Stopped CLP.") + def _get_docker_file_name(self) -> str: + """ + :return: The Docker Compose file name to use based on the config. + """ + if CompressionOrchestration.SPIDER == self._clp_config.package.scheduler: + return "docker-compose-spider.yaml" + return "docker-compose.yaml" + @staticmethod def _get_num_workers() -> int: """ @@ -1362,6 +1507,73 @@ def _chown_recursively( subprocess.run(chown_cmd, stdout=subprocess.DEVNULL, check=True) +def _get_spider_round_robin_scheduler_env_vars(scheduler: SpiderScheduler) -> EnvVarsDict: + """ + :param scheduler: + :return: Dictionary of environment variables necessary to tune the round-robin policy. + """ + round_robin = scheduler.round_robin + return EnvVarsDict( + { + "SPIDER_SCHEDULER_ROUND_ROBIN_ACTIVE_JOB_QUEUE_CAPACITY": _str_or_none( + round_robin.active_job_queue_capacity + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_CLEANUP_READY_TASK_CAPACITY": _str_or_none( + round_robin.cleanup_ready_task_capacity + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_COMMIT_READY_TASK_CAPACITY": _str_or_none( + round_robin.commit_ready_task_capacity + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_DISPATCH_QUEUE_CAPACITY": _str_or_none( + round_robin.dispatch_queue_capacity + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_FINALIZING_JOB_EXPIRATION_TIMEOUT_SEC": _str_or_none( + round_robin.finalizing_job_expiration_timeout_sec + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_READY_TASK_CAPACITY": _str_or_none( + round_robin.ready_task_capacity + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_STORAGE_POLL_TIMEOUT_MS": _str_or_none( + round_robin.storage_poll_timeout_ms + ), + "SPIDER_SCHEDULER_ROUND_ROBIN_TICK_INTERVAL_MS": _str_or_none( + round_robin.tick_interval_ms + ), + } + ) + + +_SPIDER_SCHEDULER_POLICY_ENV_GETTERS: dict[ + SpiderSchedulerPolicy, Callable[[SpiderScheduler], EnvVarsDict] +] = { + SpiderSchedulerPolicy.ROUND_ROBIN: _get_spider_round_robin_scheduler_env_vars, +} + + +def _get_spider_scheduler_policy_env_vars(scheduler: SpiderScheduler) -> EnvVarsDict: + """ + :param scheduler: + :return: Dictionary of environment variables necessary to tune the policy selected by + `scheduler.policy`, or by `SpiderScheduler.DEFAULT_POLICY` if it's unset. + :raise NotImplementedError: If no getter is registered for the selected policy. + """ + policy = SpiderScheduler.DEFAULT_POLICY if scheduler.policy is None else scheduler.policy + get_env = _SPIDER_SCHEDULER_POLICY_ENV_GETTERS.get(policy) + if get_env is None: + msg = f"No environment getter is registered for Spider scheduler policy `{policy}`." + raise NotImplementedError(msg) + + return get_env(scheduler) + + +def _str_or_none(value: int | pathlib.Path | None) -> str | None: + """ + :param value: + :return: `value` as a string, or None if `value` is None. + """ + return None if value is None else str(value) + + def _get_ip_from_hostname(hostname: str) -> str: """ Resolves a hostname to an IPv4 IP address. diff --git a/components/clp-package-utils/clp_package_utils/general.py b/components/clp-package-utils/clp_package_utils/general.py index 2f2ab255b5..b9fd0ab718 100644 --- a/components/clp-package-utils/clp_package_utils/general.py +++ b/components/clp-package-utils/clp_package_utils/general.py @@ -31,6 +31,7 @@ REDIS_COMPONENT_NAME, REDUCER_COMPONENT_NAME, RESULTS_CACHE_COMPONENT_NAME, + SPIDER_COMPONENT_NAME, StorageType, WEBUI_COMPONENT_NAME, WorkerConfig, @@ -658,6 +659,10 @@ def validate_webui_config( validate_port(f"{WEBUI_COMPONENT_NAME}.port", clp_config.webui.host, clp_config.webui.port) +def validate_spider_config(clp_config: ClpConfig): + validate_port(f"{SPIDER_COMPONENT_NAME}.port", clp_config.spider.host, clp_config.spider.port) + + def validate_mcp_server_config(clp_config: ClpConfig, logs_dir: pathlib.Path): _validate_log_directory(logs_dir, MCP_SERVER_COMPONENT_NAME) 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 78607db272..53b89b072e 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -16,7 +16,7 @@ PrivateAttr, StringConstraints, ) -from strenum import KebabCaseStrEnum, LowercaseStrEnum +from strenum import KebabCaseStrEnum, LowercaseStrEnum, SnakeCaseStrEnum from clp_py_utils.clp_logging import LoggingLevel from clp_py_utils.core import ( @@ -36,6 +36,7 @@ REDUCER_COMPONENT_NAME = "reducer" RESULTS_CACHE_COMPONENT_NAME = "results_cache" OTEL_COLLECTOR_COMPONENT_NAME = "otel-collector" +COMPRESSION_COORDINATOR_COMPONENT_NAME = "compression_coordinator" COMPRESSION_SCHEDULER_COMPONENT_NAME = "compression_scheduler" QUERY_SCHEDULER_COMPONENT_NAME = "query_scheduler" PRESTO_COORDINATOR_COMPONENT_NAME = "presto-coordinator" @@ -45,6 +46,8 @@ LOG_INGESTOR_COMPONENT_NAME = "log_ingestor" WEBUI_COMPONENT_NAME = "webui" MCP_SERVER_COMPONENT_NAME = "mcp_server" +SPIDER_COMPONENT_NAME = "spider" +SPIDER_STORAGE_COMPONENT_NAME = "spider-storage" GARBAGE_COLLECTOR_COMPONENT_NAME = "garbage_collector" # Action names @@ -135,6 +138,21 @@ class DatabaseEngine(KebabCaseStrEnum): DatabaseEngineStr = Annotated[DatabaseEngine, StrEnumSerializer] +class CompressionOrchestration(KebabCaseStrEnum): + CELERY = auto() + SPIDER = auto() + + +CompressionOrchestrationStr = Annotated[CompressionOrchestration, StrEnumSerializer] + + +class SpiderSchedulerPolicy(SnakeCaseStrEnum): + ROUND_ROBIN = auto() + + +SpiderSchedulerPolicyStr = Annotated[SpiderSchedulerPolicy, StrEnumSerializer] + + class QueryEngine(KebabCaseStrEnum): CLP = auto() CLP_S = auto() @@ -161,6 +179,7 @@ class AwsAuthType(LowercaseStrEnum): class Package(BaseModel): storage_engine: StorageEngineStr = StorageEngine.CLP_S + scheduler: CompressionOrchestrationStr = CompressionOrchestration.CELERY class ClpDbUserType(KebabCaseStrEnum): @@ -750,9 +769,89 @@ class LogIngestor(BaseModel): logging_level: LoggingLevelRust = "INFO" +class SpiderInboundQueue(BaseModel): + cleanup_capacity: PositiveInt | None = None + commit_capacity: PositiveInt | None = None + task_capacity: PositiveInt | None = None + + +class SpiderJobCacheGc(BaseModel): + gc_interval_sec: PositiveInt | None = None + terminated_job_retention_sec: PositiveInt | None = None + + +class SpiderTaskInstancePool(BaseModel): + execution_manager_stale_cutoff_sec: PositiveInt | None = None + gc_interval_sec: PositiveInt | None = None + message_channel_capacity: PositiveInt | None = None + + +class SpiderStorage(BaseModel): + log_level: LoggingLevelRust | None = None + db_max_connections: PositiveInt | None = None + inbound_queue: SpiderInboundQueue = SpiderInboundQueue() + job_cache_gc: SpiderJobCacheGc = SpiderJobCacheGc() + task_instance_pool: SpiderTaskInstancePool = SpiderTaskInstancePool() + + +class SpiderEmRegistry(BaseModel): + dead_em_cutoff_sec: PositiveInt | None = None + liveness_tracking_interval_ms: PositiveInt | None = None + + +class SpiderRoundRobin(BaseModel): + active_job_queue_capacity: PositiveInt | None = None + cleanup_ready_task_capacity: PositiveInt | None = None + commit_ready_task_capacity: PositiveInt | None = None + dispatch_queue_capacity: PositiveInt | None = None + finalizing_job_expiration_timeout_sec: PositiveInt | None = None + ready_task_capacity: PositiveInt | None = None + storage_poll_timeout_ms: PositiveInt | None = None + tick_interval_ms: PositiveInt | None = None + + +class SpiderScheduler(BaseModel): + DEFAULT_POLICY: ClassVar[SpiderSchedulerPolicy] = SpiderSchedulerPolicy.ROUND_ROBIN + + log_level: LoggingLevelRust | None = None + policy: SpiderSchedulerPolicyStr | None = None + connection_pool_size: PositiveInt | None = None + stop_timeout_sec: PositiveInt | None = None + em_registry: SpiderEmRegistry = SpiderEmRegistry() + round_robin: SpiderRoundRobin = SpiderRoundRobin() + + +class SpiderLiveness(BaseModel): + scheduler_heartbeat_interval_sec: PositiveInt | None = None + storage_heartbeat_interval_sec: PositiveInt | None = None + + +class SpiderWorker(BaseModel): + replicas: PositiveInt | None = None + log_level: LoggingLevelRust | None = None + connection_pool_size: PositiveInt | None = None + scheduler_poll_wait_ms: PositiveInt | None = None + max_log_line_bytes: PositiveInt | None = None + liveness: SpiderLiveness = SpiderLiveness() + + class Spider(BaseModel): + DEFAULT_PORT: ClassVar[int] = 50051 + host: DomainStr = "localhost" - port: Port = 6000 + port: Port = DEFAULT_PORT + + storage: SpiderStorage = SpiderStorage() + scheduler: SpiderScheduler = SpiderScheduler() + worker: SpiderWorker = SpiderWorker() + + def dump_to_primitive_dict(self): + """:return: A dictionary representation of this model, excluding Spider's own settings.""" + return self.model_dump(exclude={"storage", "scheduler", "worker"}) + + def transform_for_container(self): + self.host = SPIDER_STORAGE_COMPONENT_NAME + self.port = self.DEFAULT_PORT class SpiderResourceGroup(BaseModel): @@ -765,6 +864,7 @@ class PollingBackoff(BaseModel): class CompressionCoordinator(BaseModel): + logging_level: LoggingLevelRust = "INFO" resource_group: SpiderResourceGroup = SpiderResourceGroup(name="compression-coordinator") job_polling_interval_millisecs: PositiveInt = 100 max_concurrent_jobs: PositiveInt = 1000 @@ -1008,7 +1108,7 @@ def get_shared_config_file_path(self) -> pathlib.Path: return self.logs_directory / CLP_SHARED_CONFIG_FILENAME def dump_to_primitive_dict(self): - custom_serialized_fields = {"database", "queue", "redis"} + custom_serialized_fields = {"database", "queue", "redis", "spider"} d = self.model_dump(exclude=custom_serialized_fields) for key in custom_serialized_fields: value = getattr(self, key) @@ -1025,6 +1125,26 @@ def validate_log_ingestor_config(self): raise ValueError(msg) return self + @model_validator(mode="after") + def validate_compression_orchestration_config(self): + if CompressionOrchestration.SPIDER != self.package.scheduler: + # Neither service is deployed in this mode, so no `spider` or + # `compression_coordinator` config check is necessary. + return self + if self.spider is None: + msg = ( + "`spider` must be configured when `package.scheduler` is" + f" `{CompressionOrchestration.SPIDER}`." + ) + raise ValueError(msg) + if self.compression_coordinator is None: + msg = ( + "`compression_coordinator` must be configured when `package.scheduler` is" + f" `{CompressionOrchestration.SPIDER}`." + ) + raise ValueError(msg) + return self + @model_validator(mode="after") def validate_compression_coordinator_config(self): if self.compression_coordinator is None: @@ -1109,6 +1229,8 @@ def transform_for_container(self): self.reducer.transform_for_container() if self.presto is not None: self.presto.transform_for_container() + if self.spider is not None: + self.spider.transform_for_container() class WorkerConfig(BaseModel): 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 923447f432..98a8a9753c 100644 --- a/components/package-template/src/etc/clp-config.template.json.yaml +++ b/components/package-template/src/etc/clp-config.template.json.yaml @@ -8,6 +8,7 @@ telemetry: #package: # storage_engine: "clp-s" +# scheduler: "celery" ## API server config #api_server: @@ -57,6 +58,36 @@ telemetry: # logging_level: "INFO" # telemetry_update_interval_ms: 60000 # +## `spider` and `compression_coordinator` default to `null`. To use them, set `package.scheduler` +## to "spider", remove `null` from both, and uncomment their fields below. +#spider: null +## host: "localhost" +## port: 50051 +## storage: +## log_level: "INFO" +## scheduler: +## log_level: "INFO" +## policy: "round_robin" +## worker: +## replicas: 4 +## log_level: "INFO" +# +#compression_coordinator: null +## logging_level: "INFO" +## resource_group: +## name: "compression-coordinator" +## job_polling_interval_millisecs: 100 +## max_concurrent_jobs: 1000 +## result_polling: +## init_backoff_millisecs: 100 +## max_backoff_millisecs: 1000 +## compression_task_max_retry: 1 +## commit_task_max_retry: 1 +## database_connection_pool_size: 10 +## termination_timeout_secs: 30 +## commit_task_soft_timeout_secs: 45 +## commit_task_hard_timeout_secs: 60 +# #query_scheduler: # host: "localhost" # port: 7000 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 f2f8c16bf8..9b1264905e 100644 --- a/components/package-template/src/etc/clp-config.template.text.yaml +++ b/components/package-template/src/etc/clp-config.template.text.yaml @@ -61,6 +61,10 @@ log_ingestor: null # logging_level: "INFO" # telemetry_update_interval_ms: 60000 # +## `spider` and `compression_coordinator` are unsupported with CLP Text and must stay `null`. +#spider: null +#compression_coordinator: null +# #query_scheduler: # host: "localhost" # port: 7000