Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 2 additions & 47 deletions components/clp-package-utils/clp_package_utils/controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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),
Expand All @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 = []
Expand Down
2 changes: 0 additions & 2 deletions components/clp-package-utils/clp_package_utils/general.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)},
Expand Down
53 changes: 0 additions & 53 deletions components/clp-py-utils/clp_py_utils/clp_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -179,21 +168,18 @@ class ClpDbUserType(KebabCaseStrEnum):

CLP = auto()
ROOT = auto()
SPIDER = auto()


class ClpDbNameType(KebabCaseStrEnum):

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rust mirrors this config in clp-rust-utils which is not updated:

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The rust side has already removed the spider option.

"""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,
}
)

Expand All @@ -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
Expand Down Expand Up @@ -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}'."
Expand All @@ -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)
Expand All @@ -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


Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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:
Expand All @@ -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:
Expand Down
20 changes: 1 addition & 19 deletions components/clp-py-utils/clp_py_utils/create-db-tables.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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


Expand Down
Loading
Loading