From 463906073e15e3ea42ea7246e46db70dfec8a0f0 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Mon, 24 Nov 2025 23:08:29 -0300 Subject: [PATCH] SIENTIAPDE-1350: Refactor: Improve logging, configuration, and resource management. Includes .env updates, client initialization logging, and resource closing. --- .env.example | 29 ++++++--------- model_manager/activities/activities.py | 2 + .../activities/experiment_tracking.py | 2 + model_manager/utils/connectors_config.py | 1 + .../utils/repository/model_repository.py | 2 +- .../utils/repository/storage_repository.py | 6 ++- model_manager/worker/worker.py | 17 ++++----- run_local.sh | 8 ---- scripts/run_training_test.py | 37 +++++++++++-------- todo-list.txt | 2 +- 10 files changed, 53 insertions(+), 53 deletions(-) diff --git a/.env.example b/.env.example index f977184..6efb9cb 100644 --- a/.env.example +++ b/.env.example @@ -19,6 +19,7 @@ TEMPORAL_HOST=temporal-frontend.temporal.svc.cluster.local:7233 TEMPORAL_NAMESPACE=model-manager TRAIN_TASK_QUEUE=train_model-queue CLEANUP_TASK_QUEUE=cleanup-queue +TEMPORAL_USE_TLS=false MONGODB_USERNAME=mongo_user MONGODB_PASSWORD=mongo_db_password @@ -36,24 +37,16 @@ MINIO_RETRY_MODE=adaptive MINIO_CONNECT_TIMEOUT=10 MINIO_READ_TIMEOUT=60 -# Workflow Activity Timeouts (in seconds) -# These timeouts are designed to handle large files (up to 200MB) -TIMEOUT_VALIDATE_PARAMS=30 # Parameter validation (fast operation) -TIMEOUT_TRAIN_MODEL=2700 # Model training (30 min for large datasets) -TIMEOUT_DELETE_FILE=120 # Delete file from MinIO (1 min) -TIMEOUT_UPDATE_DATABASE=30 # Database update operations (30 sec) +TIMEOUT_VALIDATE_PARAMS=30 +TIMEOUT_TRAIN_MODEL=2700 +TIMEOUT_DELETE_FILE=120 +TIMEOUT_UPDATE_DATABASE=30 -# Cleanup Configuration -# Cleanup retention period in hours (files older than this will be deleted) -CLEANUP_RETENTION_HOURS=24 # Default: 24 hours -# Enable dry run mode to test without actually deleting files -CLEANUP_DRY_RUN=false # Set to true for testing without deletion -# Cleanup operation timeouts -TIMEOUT_CLEANUP_MINIO=300 # MinIO cleanup timeout (5 minutes) -TIMEOUT_CLEANUP_LOCAL=120 # Local directory cleanup timeout (2 minutes) -# MinIO list operation page size for cleanup -MAX_KEYS_CLEANUP=1000 # Maximum keys per page when listing objects -# Default bucket for cleanup operations -DEFAULT_CLEANUP_BUCKET=model-training # Default bucket to clean +CLEANUP_RETENTION_HOURS=24 +CLEANUP_DRY_RUN=false +TIMEOUT_CLEANUP_MINIO=300 +TIMEOUT_CLEANUP_LOCAL=120 +MAX_KEYS_CLEANUP=1000 +DEFAULT_CLEANUP_BUCKET=model-training EXTRA_PIP_REQUIREMENTS=git+https://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@github.com/Aignosi/sientia-mlops-library.git diff --git a/model_manager/activities/activities.py b/model_manager/activities/activities.py index 177f9ed..972c27e 100644 --- a/model_manager/activities/activities.py +++ b/model_manager/activities/activities.py @@ -152,3 +152,5 @@ class Activities(ExperimentTracking, Training, Cleanup): Prefer calling this method explicitly rather than relying on __del__. """ ExperimentTracking.close(self) + self.info('Postgres client closed') + self.storage_repository.close() diff --git a/model_manager/activities/experiment_tracking.py b/model_manager/activities/experiment_tracking.py index c9686de..7ab6d59 100644 --- a/model_manager/activities/experiment_tracking.py +++ b/model_manager/activities/experiment_tracking.py @@ -90,6 +90,8 @@ class ExperimentTracking(Postgres): metrics_controller=metrics_controller, ) + self.info(f'Postgres client initialized at {host}:{port}') + def __del__(self): """ Destructor to safely handle cleanup during garbage collection. diff --git a/model_manager/utils/connectors_config.py b/model_manager/utils/connectors_config.py index b42d017..9a1e40a 100644 --- a/model_manager/utils/connectors_config.py +++ b/model_manager/utils/connectors_config.py @@ -84,6 +84,7 @@ def build_mongodb_config() -> dict[str, Any]: 'connection_string': connection_string, 'database_name': getenv('MONGODB_DATABASE_NAME', 'sientia'), 'ttl_index_seconds': int(getenv('MONGODB_TTL_INDEX_HOURS', '1')) * 3600, + 'uri': uri, } diff --git a/model_manager/utils/repository/model_repository.py b/model_manager/utils/repository/model_repository.py index 5154f84..803e7da 100644 --- a/model_manager/utils/repository/model_repository.py +++ b/model_manager/utils/repository/model_repository.py @@ -31,8 +31,8 @@ warnings.filterwarnings('ignore', category=FutureWarning, message=".*'squared' i class ModelRepository: def __init__(self, url, username, password, logger: Logger): self.model_serving = ModelServing(tracking_uri=url, username=username, password=password) - self.logger = logger + self.logger.info(f'MLFlow client initialized at {url}') def save_model(self, train_result: TrainModelResult) -> TrainModelResult: """ diff --git a/model_manager/utils/repository/storage_repository.py b/model_manager/utils/repository/storage_repository.py index 3c655f4..c8639fb 100644 --- a/model_manager/utils/repository/storage_repository.py +++ b/model_manager/utils/repository/storage_repository.py @@ -86,7 +86,11 @@ class StorageRepository: use_ssl=use_ssl, ) - self.logger.info(f'MinIO client initialized successfully: {endpoint_url}') + self.logger.info(f'MinIO client initialized at {endpoint_url}') + + def close(self) -> None: + self.minio_client.close() + self.logger.info('MinIO client closed') def fetch_file(self, bucket_name: str, file_name: str) -> BytesIO: """ diff --git a/model_manager/worker/worker.py b/model_manager/worker/worker.py index d290deb..29200aa 100644 --- a/model_manager/worker/worker.py +++ b/model_manager/worker/worker.py @@ -73,15 +73,14 @@ async def main(): SystemExit: On graceful shutdown or error conditions """ host = os.getenv('TEMPORAL_HOST', 'localhost:7233') + use_tls = os.getenv('TEMPORAL_USE_TLS', 'false').lower() == 'true' logger = get_logger(__name__) metadata = { 'pod_id': POD_ID, } - logger.custom_info(f'Starting Worker with POD_ID: {POD_ID}', metadata) start_prometheus_server(logger, metadata) - logger.custom_info('Starting Notification Handler...', metadata) mongo_config = build_mongodb_config() notification_handler = NotificationHandler( @@ -91,7 +90,7 @@ async def main(): project_name=os.getenv('PROJECT_NAME', 'model-manager'), ) - logger.custom_info('Starting Activities...', metadata) + logger.custom_info(f'MongoDB client initialized at {mongo_config["uri"]}', metadata) activities = Activities( postgres_config=build_postgres_config(), @@ -101,23 +100,22 @@ async def main(): notification_handler=notification_handler, ) - logger.custom_info(f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata) - new_runtime = Runtime( telemetry=TelemetryConfig( metrics=PrometheusConfig(bind_address=f'0.0.0.0:{SDK_METRICS_PORT}') ) ) - logger.custom_info(f'Starting Temporal Client at {host}...', metadata) + logger.custom_info(f'SDK metrics server initialized on port {SDK_METRICS_PORT}', metadata) temporal_client = await client.Client.connect( target_host=host, namespace=os.getenv('TEMPORAL_NAMESPACE', 'model-manager'), runtime=new_runtime, + tls=use_tls, ) - logger.custom_info('Starting Workers...', metadata) + logger.custom_info(f'Temporal client initialized at {host}', metadata) workers = [ Worker( @@ -159,7 +157,7 @@ async def main(): for w in workers: handlers.append(w.run()) - logger.custom_info('Workers started successfully', metadata) + logger.custom_info('Model manager workers initialized', metadata) try: # This will run the workers and wait for them to complete. @@ -169,6 +167,7 @@ async def main(): logger.custom_error(f'An unhandled exception occurred: {e}', metadata) finally: notification_handler.shutdown() + logger.custom_info('MongoDB client closed', metadata) await activities.shutdown() # Exit with a non-zero status code to indicate failure to Kubernetes metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN @@ -195,7 +194,7 @@ def start_prometheus_server(logger: SientiaLogger, metadata: dict[str, str | Non try: port = int(os.getenv('HTTP_METRICS_PORT', 9090)) start_http_server(port) - logger.custom_info(f'Prometheus server started on port {port}.', metadata) + logger.custom_info(f'Prometheus server initialized on port {port}.', metadata) metrics.APP_UP.labels(pod_id=POD_ID).set(1) # Mark app as UP except Exception as e: # noqa: BLE001 logger.custom_critical(f'Failed to start Prometheus server: {e}', metadata) diff --git a/run_local.sh b/run_local.sh index d4c8c10..4c88728 100755 --- a/run_local.sh +++ b/run_local.sh @@ -3,12 +3,6 @@ # Exit on any error set -e -#echo "Activating virtual environment..." - -#conda activate ./venv - -echo "Loading environment variables from .env..." - if [ -f .env ]; then export $(cat .env | grep -v '^#' | xargs) echo "Environment variables loaded from .env" @@ -16,6 +10,4 @@ else echo "Warning: .env file not found. Continuing without environment variables." fi -echo "Starting ingestor application..." - python -m model_manager.worker.worker diff --git a/scripts/run_training_test.py b/scripts/run_training_test.py index dbb641f..a978daa 100644 --- a/scripts/run_training_test.py +++ b/scripts/run_training_test.py @@ -24,29 +24,35 @@ import uuid from datetime import datetime, timedelta from pathlib import Path +from dotenv import load_dotenv import psycopg2 from psycopg2.extras import Json from temporalio import client + +# Carrega variáveis de ambiente do arquivo .env na raiz do projeto +PROJECT_ROOT = Path(__file__).resolve().parent.parent +ENV_PATH = PROJECT_ROOT / '.env' +if ENV_PATH.exists(): + load_dotenv(dotenv_path=ENV_PATH) + + DOCS_PATH = Path('docs/test-model-data.csv') -MINIO_ALIAS = os.getenv('MINIO_ALIAS', 'suse') -MINIO_BUCKET = os.getenv('MINIO_BUCKET', 'model-training') +MINIO_ALIAS = 'suse' +MINIO_BUCKET = 'model-training' POSTGRES_CONFIG = { - 'host': os.getenv('POSTGRES_HOST', 'localhost'), - 'port': os.getenv('POSTGRES_PORT', '55432'), - 'user': os.getenv('POSTGRES_USER', 'postgres'), - 'password': os.getenv( - 'POSTGRES_PASSWORD', - 'nFqc81y6kwmr2zuAIx43DhiOosFCVPpeEfTtTWZflkNjB2j1KtEeIANkhFR9mAX3', - ), - 'dbname': os.getenv('POSTGRES_DBNAME', 'sientia-core-mlops-bff'), + 'host': os.getenv('POSTGRES_HOST'), + 'port': os.getenv('POSTGRES_PORT'), + 'user': os.getenv('POSTGRES_USER'), + 'password': os.getenv('POSTGRES_PASSWORD'), + 'dbname': os.getenv('POSTGRES_DBNAME'), } -TEMPORAL_HOST = os.getenv('TEMPORAL_HOST', 'localhost:37463') -TEMPORAL_NAMESPACE = os.getenv('TEMPORAL_NAMESPACE', 'model-manager') -TEMPORAL_TASK_QUEUE = os.getenv('TEMPORAL_TASK_QUEUE', 'train_model-queue') -TEMPORAL_WORKFLOW = os.getenv('TEMPORAL_WORKFLOW', 'train_model') +TEMPORAL_HOST = os.getenv('TEMPORAL_HOST') +TEMPORAL_NAMESPACE = os.getenv('TEMPORAL_NAMESPACE') +TRAIN_TASK_QUEUE = os.getenv('TRAIN_TASK_QUEUE') +TEMPORAL_WORKFLOW = 'train_model' BASE_REQUEST_DATA = { 'experimentName': 'model-manager-test-01', @@ -169,6 +175,7 @@ async def trigger_temporal_workflow(workflow_input: dict) -> str: temporal_client = await client.Client.connect( target_host=TEMPORAL_HOST, namespace=TEMPORAL_NAMESPACE, + tls=os.getenv('TEMPORAL_USE_TLS', False), ) workflow_id = f'train-model-test-{uuid.uuid4()}' @@ -176,7 +183,7 @@ async def trigger_temporal_workflow(workflow_input: dict) -> str: TEMPORAL_WORKFLOW, workflow_input, id=workflow_id, - task_queue=TEMPORAL_TASK_QUEUE, + task_queue=TRAIN_TASK_QUEUE, execution_timeout=timedelta(minutes=5), run_timeout=timedelta(minutes=5), task_timeout=timedelta(minutes=5), diff --git a/todo-list.txt b/todo-list.txt index ba89408..a2d8850 100644 --- a/todo-list.txt +++ b/todo-list.txt @@ -1,7 +1,7 @@ -- Atualizar o sientia-dataops-library para a versão 1.6.1 - Refatorar o arquivo .dockerignore para só deixar copiar os arquivos que forem necessários para a execução do container, pois ele está copiando muitos arquivos desnecessários. - atualizar as variáveis de ambiente no helm chart - Criar um gráfico no grafana para cada nova atividade. +- Atualizar a documentação dos métodos alterados. - Atualizar a documentação do projeto. - Atualizar o .github/workflows/quality-gate.yml para usar os pipelines genéricos do github;