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.
This commit is contained in:
@@ -17,6 +17,8 @@ PROJECT_NAME=sientia-model-manager
|
|||||||
|
|
||||||
TEMPORAL_HOST=temporal-frontend.temporal.svc.cluster.local:7233
|
TEMPORAL_HOST=temporal-frontend.temporal.svc.cluster.local:7233
|
||||||
TEMPORAL_NAMESPACE=model-manager
|
TEMPORAL_NAMESPACE=model-manager
|
||||||
|
TRAIN_TASK_QUEUE=train_model-queue
|
||||||
|
CLEANUP_TASK_QUEUE=cleanup-queue
|
||||||
|
|
||||||
MONGODB_USERNAME=mongo_user
|
MONGODB_USERNAME=mongo_user
|
||||||
MONGODB_PASSWORD=mongo_db_password
|
MONGODB_PASSWORD=mongo_db_password
|
||||||
|
|||||||
@@ -33,6 +33,7 @@ with workflow.unsafe.imports_passed_through():
|
|||||||
|
|
||||||
from prometheus_client import start_http_server
|
from prometheus_client import start_http_server
|
||||||
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
|
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 import metrics
|
||||||
from model_manager.activities.activities import Activities
|
from model_manager.activities.activities import Activities
|
||||||
@@ -48,6 +49,8 @@ with workflow.unsafe.imports_passed_through():
|
|||||||
|
|
||||||
POD_ID = os.getenv('POD_ID')
|
POD_ID = os.getenv('POD_ID')
|
||||||
SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091'))
|
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():
|
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(f'Starting Worker with POD_ID: {POD_ID}', metadata)
|
||||||
|
start_prometheus_server(logger, metadata)
|
||||||
logger.custom_info('Starting prometheus client...', metadata)
|
|
||||||
start_prometheus_server()
|
|
||||||
|
|
||||||
logger.custom_info('Starting Notification Handler...', metadata)
|
logger.custom_info('Starting Notification Handler...', metadata)
|
||||||
|
|
||||||
mongo_config = build_mongodb_config()
|
mongo_config = build_mongodb_config()
|
||||||
|
|
||||||
notification_handler = NotificationHandler(
|
notification_handler = NotificationHandler(
|
||||||
connection_string=mongo_config['connection_string'],
|
connection_string=mongo_config['connection_string'],
|
||||||
database=mongo_config['database_name'],
|
database=mongo_config['database_name'],
|
||||||
@@ -122,7 +122,7 @@ async def main():
|
|||||||
workers = [
|
workers = [
|
||||||
Worker(
|
Worker(
|
||||||
temporal_client,
|
temporal_client,
|
||||||
task_queue='train_model-queue',
|
task_queue=TRAIN_TASK_QUEUE,
|
||||||
workflows=[TrainModel],
|
workflows=[TrainModel],
|
||||||
activities=[
|
activities=[
|
||||||
activities.update_experiment_run,
|
activities.update_experiment_run,
|
||||||
@@ -139,7 +139,7 @@ async def main():
|
|||||||
),
|
),
|
||||||
Worker(
|
Worker(
|
||||||
temporal_client,
|
temporal_client,
|
||||||
task_queue='cleanup-queue',
|
task_queue=CLEANUP_TASK_QUEUE,
|
||||||
workflows=[CleanupFiles],
|
workflows=[CleanupFiles],
|
||||||
activities=[
|
activities=[
|
||||||
activities.cleanup_minio_files,
|
activities.cleanup_minio_files,
|
||||||
@@ -155,6 +155,7 @@ async def main():
|
|||||||
]
|
]
|
||||||
|
|
||||||
handlers = []
|
handlers = []
|
||||||
|
|
||||||
for w in workers:
|
for w in workers:
|
||||||
handlers.append(w.run())
|
handlers.append(w.run())
|
||||||
|
|
||||||
@@ -174,7 +175,7 @@ async def main():
|
|||||||
sys.exit(1)
|
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.
|
Starts the Prometheus metrics server for monitoring and observability.
|
||||||
|
|
||||||
@@ -194,10 +195,10 @@ def start_prometheus_server():
|
|||||||
try:
|
try:
|
||||||
port = int(os.getenv('HTTP_METRICS_PORT', 9090))
|
port = int(os.getenv('HTTP_METRICS_PORT', 9090))
|
||||||
start_http_server(port)
|
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
|
metrics.APP_UP.labels(pod_id=POD_ID).set(1) # Mark app as UP
|
||||||
except Exception as e: # noqa: BLE001
|
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)
|
os._exit(1)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -3,7 +3,7 @@ psycopg2-binary==2.9.11
|
|||||||
sqlalchemy==2.0.44
|
sqlalchemy==2.0.44
|
||||||
boto3==1.40.55
|
boto3==1.40.55
|
||||||
botocore==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
|
prometheus-client==0.23.1
|
||||||
mlflow==2.10.1
|
mlflow==2.10.1
|
||||||
evidently==0.4.21
|
evidently==0.4.21
|
||||||
|
|||||||
@@ -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.POD_ID', 'test-pod-123')
|
||||||
@patch('model_manager.worker.worker.start_http_server')
|
@patch('model_manager.worker.worker.start_http_server')
|
||||||
@patch('model_manager.worker.worker.metrics')
|
@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."""
|
"""Test successful Prometheus server startup."""
|
||||||
from model_manager.worker.worker import start_prometheus_server
|
from model_manager.worker.worker import start_prometheus_server
|
||||||
|
|
||||||
mock_app_up = Mock()
|
mock_app_up = Mock()
|
||||||
mock_metrics.APP_UP.labels.return_value = mock_app_up
|
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
|
# Verify HTTP server started
|
||||||
mock_start_http_server.assert_called_once_with(9090)
|
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
|
# Verify APP_UP metric was set to 1
|
||||||
mock_metrics.APP_UP.labels.assert_called_once_with(pod_id='test-pod-123')
|
mock_metrics.APP_UP.labels.assert_called_once_with(pod_id='test-pod-123')
|
||||||
mock_app_up.set.assert_called_once_with(1)
|
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.start_http_server')
|
||||||
@patch('model_manager.worker.worker.metrics')
|
@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."""
|
"""Test Prometheus server startup with custom port."""
|
||||||
from model_manager.worker.worker import start_prometheus_server
|
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_app_up = Mock()
|
||||||
mock_metrics.APP_UP.labels.return_value = mock_app_up
|
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)
|
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.metrics')
|
||||||
@patch('model_manager.worker.worker.os._exit')
|
@patch('model_manager.worker.worker.os._exit')
|
||||||
def test_start_prometheus_server_failure(
|
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."""
|
"""Test Prometheus server startup failure."""
|
||||||
from model_manager.worker.worker import start_prometheus_server
|
from model_manager.worker.worker import start_prometheus_server
|
||||||
|
|
||||||
mock_start_http_server.side_effect = OSError('Port already in use')
|
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_exit.assert_called_once_with(1)
|
||||||
|
mock_logger.custom_critical.assert_called_once()
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
@@ -234,7 +244,7 @@ async def test_main_successful_startup(
|
|||||||
mock_notification_handler_class.assert_called_once()
|
mock_notification_handler_class.assert_called_once()
|
||||||
mock_activities_class.assert_called_once()
|
mock_activities_class.assert_called_once()
|
||||||
mock_client_class.connect.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
|
assert mock_worker_class.call_count == 2
|
||||||
|
|
||||||
# Verify cleanup was performed
|
# 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.start_http_server')
|
||||||
@patch('model_manager.worker.worker.metrics')
|
@patch('model_manager.worker.worker.metrics')
|
||||||
def test_start_prometheus_server_prints_success(
|
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."""
|
"""Test that start_prometheus_server prints success message."""
|
||||||
from model_manager.worker.worker import start_prometheus_server
|
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_app_up = Mock()
|
||||||
mock_metrics.APP_UP.labels.return_value = mock_app_up
|
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()
|
start_prometheus_server(mock_logger, metadata)
|
||||||
assert 'Prometheus server started on port 9090' in captured.out
|
|
||||||
|
# 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.start_http_server')
|
||||||
@patch('model_manager.worker.worker.metrics')
|
@patch('model_manager.worker.worker.metrics')
|
||||||
@patch('model_manager.worker.worker.os._exit')
|
@patch('model_manager.worker.worker.os._exit')
|
||||||
def test_start_prometheus_server_prints_failure(
|
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."""
|
"""Test that start_prometheus_server prints failure message."""
|
||||||
from model_manager.worker.worker import start_prometheus_server
|
from model_manager.worker.worker import start_prometheus_server
|
||||||
|
|
||||||
mock_start_http_server.side_effect = Exception('Test error')
|
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()
|
start_prometheus_server(mock_logger, metadata)
|
||||||
assert 'Failed to start Prometheus server' in captured.out
|
|
||||||
assert 'Test error' in captured.out
|
# Agora o erro é logado via logger crítico
|
||||||
|
mock_logger.custom_critical.assert_called_once()
|
||||||
|
mock_exit.assert_called_once_with(1)
|
||||||
|
|||||||
Reference in New Issue
Block a user