From 2ee66552d2b55e85b22637223990359f52775ce6 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Mon, 24 Nov 2025 11:21:43 -0300 Subject: [PATCH] SIENTIAPDE-1350: Configure separate task queues for train and cleanup workers, update sientia-dataops-library to 1.6.1 and refactor prometheus server startup to use logger. --- .env.example | 2 ++ model_manager/worker/worker.py | 21 ++++++++------- requirements.txt | 2 +- tests/worker/test_worker.py | 48 ++++++++++++++++++++++------------ 4 files changed, 45 insertions(+), 28 deletions(-) diff --git a/.env.example b/.env.example index b106f15..f977184 100644 --- a/.env.example +++ b/.env.example @@ -17,6 +17,8 @@ PROJECT_NAME=sientia-model-manager TEMPORAL_HOST=temporal-frontend.temporal.svc.cluster.local:7233 TEMPORAL_NAMESPACE=model-manager +TRAIN_TASK_QUEUE=train_model-queue +CLEANUP_TASK_QUEUE=cleanup-queue MONGODB_USERNAME=mongo_user MONGODB_PASSWORD=mongo_db_password diff --git a/model_manager/worker/worker.py b/model_manager/worker/worker.py index ff98693..d290deb 100644 --- a/model_manager/worker/worker.py +++ b/model_manager/worker/worker.py @@ -33,6 +33,7 @@ with workflow.unsafe.imports_passed_through(): from prometheus_client import start_http_server from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler + from sientia_do.observability.logger import Logger as SientiaLogger from model_manager import metrics from model_manager.activities.activities import Activities @@ -48,6 +49,8 @@ with workflow.unsafe.imports_passed_through(): POD_ID = os.getenv('POD_ID') SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091')) +TRAIN_TASK_QUEUE = os.getenv('TRAIN_TASK_QUEUE', 'train_model-queue') +CLEANUP_TASK_QUEUE = os.getenv('CLEANUP_TASK_QUEUE', 'cleanup-queue') async def main(): @@ -77,13 +80,10 @@ async def main(): } logger.custom_info(f'Starting Worker with POD_ID: {POD_ID}', metadata) - - logger.custom_info('Starting prometheus client...', metadata) - start_prometheus_server() - + start_prometheus_server(logger, metadata) logger.custom_info('Starting Notification Handler...', metadata) - mongo_config = build_mongodb_config() + notification_handler = NotificationHandler( connection_string=mongo_config['connection_string'], database=mongo_config['database_name'], @@ -122,7 +122,7 @@ async def main(): workers = [ Worker( temporal_client, - task_queue='train_model-queue', + task_queue=TRAIN_TASK_QUEUE, workflows=[TrainModel], activities=[ activities.update_experiment_run, @@ -139,7 +139,7 @@ async def main(): ), Worker( temporal_client, - task_queue='cleanup-queue', + task_queue=CLEANUP_TASK_QUEUE, workflows=[CleanupFiles], activities=[ activities.cleanup_minio_files, @@ -155,6 +155,7 @@ async def main(): ] handlers = [] + for w in workers: handlers.append(w.run()) @@ -174,7 +175,7 @@ async def main(): sys.exit(1) -def start_prometheus_server(): +def start_prometheus_server(logger: SientiaLogger, metadata: dict[str, str | None]): """ Starts the Prometheus metrics server for monitoring and observability. @@ -194,10 +195,10 @@ def start_prometheus_server(): try: port = int(os.getenv('HTTP_METRICS_PORT', 9090)) start_http_server(port) - print(f'Prometheus server started on port {port}.') + logger.custom_info(f'Prometheus server started on port {port}.', metadata) metrics.APP_UP.labels(pod_id=POD_ID).set(1) # Mark app as UP except Exception as e: # noqa: BLE001 - print(f'Failed to start Prometheus server: {e}') + logger.custom_critical(f'Failed to start Prometheus server: {e}', metadata) os._exit(1) diff --git a/requirements.txt b/requirements.txt index be6b7b6..362a16b 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,7 +3,7 @@ psycopg2-binary==2.9.11 sqlalchemy==2.0.44 boto3==1.40.55 botocore==1.40.55 -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.2 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.6.1 prometheus-client==0.23.1 mlflow==2.10.1 evidently==0.4.21 diff --git a/tests/worker/test_worker.py b/tests/worker/test_worker.py index 8883824..0e036c2 100644 --- a/tests/worker/test_worker.py +++ b/tests/worker/test_worker.py @@ -109,14 +109,18 @@ def test_sdk_metrics_port_from_env(): @patch('model_manager.worker.worker.POD_ID', 'test-pod-123') @patch('model_manager.worker.worker.start_http_server') @patch('model_manager.worker.worker.metrics') -def test_start_prometheus_server_success(mock_metrics, mock_start_http_server, mock_env_vars): +def test_start_prometheus_server_success( + mock_metrics, mock_start_http_server, mock_env_vars, mock_logger +): """Test successful Prometheus server startup.""" from model_manager.worker.worker import start_prometheus_server mock_app_up = Mock() mock_metrics.APP_UP.labels.return_value = mock_app_up - start_prometheus_server() + metadata = {'pod_id': 'test-pod-123', 'workflow_name': 'train_model'} + + start_prometheus_server(mock_logger, metadata) # Verify HTTP server started mock_start_http_server.assert_called_once_with(9090) @@ -124,11 +128,12 @@ def test_start_prometheus_server_success(mock_metrics, mock_start_http_server, m # Verify APP_UP metric was set to 1 mock_metrics.APP_UP.labels.assert_called_once_with(pod_id='test-pod-123') mock_app_up.set.assert_called_once_with(1) + mock_logger.custom_info.assert_called_once() @patch('model_manager.worker.worker.start_http_server') @patch('model_manager.worker.worker.metrics') -def test_start_prometheus_server_custom_port(mock_metrics, mock_start_http_server): +def test_start_prometheus_server_custom_port(mock_metrics, mock_start_http_server, mock_logger): """Test Prometheus server startup with custom port.""" from model_manager.worker.worker import start_prometheus_server @@ -136,7 +141,9 @@ def test_start_prometheus_server_custom_port(mock_metrics, mock_start_http_serve mock_app_up = Mock() mock_metrics.APP_UP.labels.return_value = mock_app_up - start_prometheus_server() + metadata = {'pod_id': 'custom-pod', 'workflow_name': 'train_model'} + + start_prometheus_server(mock_logger, metadata) mock_start_http_server.assert_called_once_with(8080) @@ -145,17 +152,20 @@ def test_start_prometheus_server_custom_port(mock_metrics, mock_start_http_serve @patch('model_manager.worker.worker.metrics') @patch('model_manager.worker.worker.os._exit') def test_start_prometheus_server_failure( - mock_exit, mock_metrics, mock_start_http_server, mock_env_vars + mock_exit, mock_metrics, mock_start_http_server, mock_env_vars, mock_logger ): """Test Prometheus server startup failure.""" from model_manager.worker.worker import start_prometheus_server mock_start_http_server.side_effect = OSError('Port already in use') - start_prometheus_server() + metadata = {'pod_id': 'test-pod-123', 'workflow_name': 'train_model'} - # Verify exit was called with code 1 + start_prometheus_server(mock_logger, metadata) + + # Verify exit was called with code 1 e log crítico emitido mock_exit.assert_called_once_with(1) + mock_logger.custom_critical.assert_called_once() @pytest.mark.asyncio @@ -234,7 +244,7 @@ async def test_main_successful_startup( mock_notification_handler_class.assert_called_once() mock_activities_class.assert_called_once() mock_client_class.connect.assert_called_once() - # Agora sao criados dois Workers: um para train_model-queue e outro para cleanup-queue + # Agora são criados dois Workers: um para train_model-queue e outro para cleanup-queue assert mock_worker_class.call_count == 2 # Verify cleanup was performed @@ -523,7 +533,7 @@ def test_worker_module_docstring(): @patch('model_manager.worker.worker.start_http_server') @patch('model_manager.worker.worker.metrics') def test_start_prometheus_server_prints_success( - mock_metrics, mock_start_http_server, capsys, mock_env_vars + mock_metrics, mock_start_http_server, capsys, mock_env_vars, mock_logger ): """Test that start_prometheus_server prints success message.""" from model_manager.worker.worker import start_prometheus_server @@ -531,25 +541,29 @@ def test_start_prometheus_server_prints_success( mock_app_up = Mock() mock_metrics.APP_UP.labels.return_value = mock_app_up - start_prometheus_server() + metadata = {'pod_id': 'test-pod-123', 'workflow_name': 'train_model'} - captured = capsys.readouterr() - assert 'Prometheus server started on port 9090' in captured.out + start_prometheus_server(mock_logger, metadata) + + # Agora a mensagem é enviada via logger + mock_logger.custom_info.assert_called_once() @patch('model_manager.worker.worker.start_http_server') @patch('model_manager.worker.worker.metrics') @patch('model_manager.worker.worker.os._exit') def test_start_prometheus_server_prints_failure( - mock_exit, mock_metrics, mock_start_http_server, capsys, mock_env_vars + mock_exit, mock_metrics, mock_start_http_server, capsys, mock_env_vars, mock_logger ): """Test that start_prometheus_server prints failure message.""" from model_manager.worker.worker import start_prometheus_server mock_start_http_server.side_effect = Exception('Test error') - start_prometheus_server() + metadata = {'pod_id': 'test-pod-123', 'workflow_name': 'train_model'} - captured = capsys.readouterr() - assert 'Failed to start Prometheus server' in captured.out - assert 'Test error' in captured.out + start_prometheus_server(mock_logger, metadata) + + # Agora o erro é logado via logger crítico + mock_logger.custom_critical.assert_called_once() + mock_exit.assert_called_once_with(1)