From 94e11df80359579437c8a65b198ab9910a181a75 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Thu, 9 Oct 2025 17:48:43 -0300 Subject: [PATCH 1/5] SIENTIAPDE-1252: Remove experiment_description from TrainModelParams and related tests. --- model_manager/utils/models/train_model_params.py | 5 ----- tests/activities/test_training.py | 5 ----- tests/utils/models/test_train_model_params.py | 15 +++++---------- tests/utils/models/test_train_model_result.py | 1 - .../utils/repository/test_training_repository.py | 3 --- 5 files changed, 5 insertions(+), 24 deletions(-) diff --git a/model_manager/utils/models/train_model_params.py b/model_manager/utils/models/train_model_params.py index 251e252..af1b1c7 100644 --- a/model_manager/utils/models/train_model_params.py +++ b/model_manager/utils/models/train_model_params.py @@ -34,7 +34,6 @@ class TrainModelParams: shuffle (bool): Whether to shuffle the data during train/test split. experiment_run_id (int): Unique identifier for the experiment run. experiment_name (str): Name of the experiment for tracking. - experiment_description (str): Description of the experiment. removed_intervals (list): List of time intervals to remove from the data. """ @@ -56,7 +55,6 @@ class TrainModelParams: shuffle: bool experiment_run_id: int experiment_name: str - experiment_description: str removed_intervals: list @classmethod @@ -114,9 +112,6 @@ class TrainModelParams: data.get('experiment_run_id'), int, 'experiment_run_id' ), experiment_name=cls._check_none(data.get('experiment_name'), str, 'experiment_name'), - experiment_description=cls._check_none( - data.get('experiment_description'), str, 'experiment_description' - ), removed_intervals=cls._check_type( data.get('removed_intervals'), list, 'removed_intervals' ), diff --git a/tests/activities/test_training.py b/tests/activities/test_training.py index 4c51c3e..9616ee1 100644 --- a/tests/activities/test_training.py +++ b/tests/activities/test_training.py @@ -61,7 +61,6 @@ async def test_train_model_success(mock_training_repository_class): 'upp_lim': {'feature1': 100.0, 'feature2': 100.0}, 'window': 10, 'experiment_name': 'test_experiment', - 'experiment_description': 'Test experiment', 'removed_intervals': [], } @@ -118,7 +117,6 @@ async def test_train_model_invalid_file_type(mock_training_repository_class): 'upp_lim': {'feature1': 100.0}, 'window': 10, 'experiment_name': 'test_experiment', - 'experiment_description': 'Test experiment', 'removed_intervals': [], }, } @@ -169,7 +167,6 @@ async def test_train_model_training_error(mock_training_repository_class): 'upp_lim': {'feature1': 100.0}, 'window': 10, 'experiment_name': 'test_experiment', - 'experiment_description': 'Test experiment', 'removed_intervals': [], }, } @@ -219,7 +216,6 @@ async def test_train_model_sends_notification_on_error(mock_training_repository_ 'upp_lim': {'feature1': 100.0}, 'window': 10, 'experiment_name': 'test_experiment', - 'experiment_description': 'Test experiment', 'removed_intervals': [], }, } @@ -274,7 +270,6 @@ async def test_train_model_after_calculation_error(mock_training_repository_clas 'upp_lim': {'feature1': 100.0}, 'window': 10, 'experiment_name': 'test_experiment', - 'experiment_description': 'Test experiment', 'removed_intervals': [], }, } diff --git a/tests/utils/models/test_train_model_params.py b/tests/utils/models/test_train_model_params.py index 9622fc5..f71a4af 100644 --- a/tests/utils/models/test_train_model_params.py +++ b/tests/utils/models/test_train_model_params.py @@ -26,8 +26,7 @@ def valid_params_dict(): 'train_size': 80, 'shuffle': True, 'experiment_run_id': 123, - 'experiment_name': 'test-experiment', - 'experiment_description': 'Test experiment description', + 'experiment_name': 'Test experiment name', 'removed_intervals': [], } @@ -53,8 +52,7 @@ def test_train_model_params_creation_with_valid_params(valid_params_dict): assert params.train_size == 80 assert params.shuffle is True assert params.experiment_run_id == 123 - assert params.experiment_name == 'test-experiment' - assert params.experiment_description == 'Test experiment description' + assert params.experiment_name == 'Test experiment name' assert params.removed_intervals == [] @@ -201,13 +199,13 @@ def test_train_model_params_removed_intervals_with_values(valid_params_dict): def test_train_model_params_all_fields_count(): - """Test that TrainModelParams has exactly 20 required fields.""" + """Test that TrainModelParams has exactly 19 required fields.""" import inspect sig = inspect.signature(TrainModelParams.__init__) # Subtract 1 for 'self' param_count = len(sig.parameters) - 1 - assert param_count == 20 + assert param_count == 19 def test_train_model_params_with_minimal_valid_data(): @@ -230,8 +228,7 @@ def test_train_model_params_with_minimal_valid_data(): train_size=50, shuffle=False, experiment_run_id=1, - experiment_name='exp', - experiment_description='desc', + experiment_name='name', removed_intervals=[], ) @@ -261,7 +258,6 @@ def test_train_model_params_check_none_method(): 'shuffle': True, 'experiment_run_id': 123, 'experiment_name': 'exp', - 'experiment_description': 'desc', 'removed_intervals': [], } @@ -293,7 +289,6 @@ def test_train_model_params_check_type_method(): 'shuffle': True, 'experiment_run_id': 123, 'experiment_name': 'exp', - 'experiment_description': 'desc', 'removed_intervals': [], } diff --git a/tests/utils/models/test_train_model_result.py b/tests/utils/models/test_train_model_result.py index 8ec77b4..0ba7912 100644 --- a/tests/utils/models/test_train_model_result.py +++ b/tests/utils/models/test_train_model_result.py @@ -31,7 +31,6 @@ def sample_params(): shuffle=True, experiment_run_id=123, experiment_name='test-experiment', - experiment_description='Test experiment description', removed_intervals=[], ) diff --git a/tests/utils/repository/test_training_repository.py b/tests/utils/repository/test_training_repository.py index fcd90f8..f60e826 100644 --- a/tests/utils/repository/test_training_repository.py +++ b/tests/utils/repository/test_training_repository.py @@ -46,7 +46,6 @@ def train_params(): shuffle=True, experiment_run_id=123, experiment_name='test_experiment', - experiment_description='Test experiment', removed_intervals=[], ) @@ -200,7 +199,6 @@ def test_init_scaler_dict_without_scaler(training_repository): shuffle=True, experiment_run_id=123, experiment_name='test', - experiment_description='test', removed_intervals=[], ) @@ -281,7 +279,6 @@ def test_after_train_calculation_without_scaler(mock_r2, mock_mae, mock_mse, tra shuffle=True, experiment_run_id=123, experiment_name='test', - experiment_description='test', removed_intervals=[], ) From 52b7b2326173b8634aba8b6f4cf5a55d98be8648 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Thu, 9 Oct 2025 17:52:35 -0300 Subject: [PATCH 2/5] SIENTIAPDE-1252: Update Python setup action and simplify dependency caching in quality gate workflow --- .github/workflows/quality-gate.yml | 14 +++++--------- 1 file changed, 5 insertions(+), 9 deletions(-) diff --git a/.github/workflows/quality-gate.yml b/.github/workflows/quality-gate.yml index 36ce79b..1708d52 100644 --- a/.github/workflows/quality-gate.yml +++ b/.github/workflows/quality-gate.yml @@ -193,17 +193,13 @@ jobs: df -h - name: 🔧 Setup Python - uses: actions/setup-python@v4 + uses: actions/setup-python@v5 with: python-version: "3.11" - - - name: 🗄️ Cache Python dependencies - uses: actions/cache@v3 - with: - path: ~/.cache/pip - key: ${{ runner.os }}-pip-${{ hashFiles('requirements.txt', 'requirements-dev.txt') }} - restore-keys: | - ${{ runner.os }}-pip- + cache: 'pip' + cache-dependency-path: | + requirements.txt + requirements-dev.txt - name: 📦 Install Dependencies run: | From 3fd5ab79feb8cfa98a7ba8ee271439ebd5003c18 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Fri, 10 Oct 2025 17:56:09 -0300 Subject: [PATCH 3/5] SIENTIAPDE-1252: Implement artifact generation and MLflow logging for model training results This commit introduces artifact generation and MLflow logging capabilities to the model training process. It includes the following changes: - Added methods to generate reports, save data files, and log model parameters, metrics, models, and artifacts to MLflow. - Implemented error handling for various scenarios, such as missing files, invalid data, and MLflow connection errors. - Created a new 'header.html' file for report styling and navigation. - Modified the 'model_repository.py' file to include the new artifact generation and MLflow logging methods. - Added comprehensive unit tests to ensure the functionality and robustness of the new features. --- model_manager/reports/header.html | 166 +++++ .../utils/repository/model_repository.py | 377 ++++++++++ .../utils/repository/test_model_repository.py | 677 ++++++++++++++++++ 3 files changed, 1220 insertions(+) create mode 100644 model_manager/reports/header.html diff --git a/model_manager/reports/header.html b/model_manager/reports/header.html new file mode 100644 index 0000000..fd646a2 --- /dev/null +++ b/model_manager/reports/header.html @@ -0,0 +1,166 @@ + + + + + + + + + Report + + + + + + +
+
+ +

Report

+ +
+ info +

+ Note that "current"
+ is related to the test
+ set while "reference"
+ refers to the training
+ set +

+
+
+
+
+
+
+
+
+
+
+
+
+ + + + diff --git a/model_manager/utils/repository/model_repository.py b/model_manager/utils/repository/model_repository.py index 5dafc76..ac6ad36 100644 --- a/model_manager/utils/repository/model_repository.py +++ b/model_manager/utils/repository/model_repository.py @@ -11,16 +11,21 @@ By Monitoring we mean the evaluation of the performance of models, the generatio """ +import shutil import traceback from datetime import datetime from os import makedirs, path, remove import mlflow +import numpy as np import pandas as pd from sientia.ModelServing import ModelServing # type: ignore[import-untyped] +from sientia.reports import Reports # type: ignore[import-untyped] from sientia_do.observability.logger import Logger from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ +from model_manager.utils.models.train_model_result import TrainModelResult + class MLFlowRepository: def __init__(self, host, username, password, logger: Logger): @@ -455,3 +460,375 @@ class MLFlowRepository: metadata['mlflow_experiment_id'] = experiment_id return metadata + + def get_next_run_name_new(self, experiment_name: str) -> str: + """ + Generates the next run name for a given experiment. + + Args: + experiment_name (str): The name of the experiment for which the next run name is being generated. + + Returns: + str: A unique run name in the format "-". + """ + runs = self.model_serving.search_runs_by_name( + experiment_names=[experiment_name], order_by=['start_time desc'] + ) + + next_run_number = len(runs) + 1 + return f'{experiment_name}-{next_run_number}' + + def generate_artifacts(self, data: TrainModelResult) -> TrainModelResult: + """ + Generates and organizes artifacts related to the training process, such as reports and data files. + + Args: + data: The training model result containing the datasets, model, and parameters. + + Returns: + The updated result object with paths to the generated artifacts. + + Raises: + FileNotFoundError: If the reports directory or header.html file does not exist. + ValueError: If run_name is not set. + """ + # Validate that run_name is set + if not data.run_name: + error_msg = 'run_name must be set before generating artifacts' + self.logger.error(error_msg) + raise ValueError(error_msg) + + reference_data, current_data = self._init_artifacts_data(data) + base_path = self._get_reports_directory() + + # Validate that reports directory exists + if not path.exists(base_path): + error_msg = f'Reports directory does not exist: {base_path}' + self.logger.error(error_msg) + raise FileNotFoundError(error_msg) + + data.run_dir = self._create_run_directory(base_path, data.run_name) + header_file_path = path.join(base_path, 'header.html') + + # Validate that header.html exists + if not path.exists(header_file_path): + error_msg = f'Header file does not exist: {header_file_path}' + self.logger.error(error_msg) + raise FileNotFoundError(error_msg) + + self._setup_run_directory(data.run_dir, header_file_path) + return self._generate_report(reference_data, current_data, data) + + def save_run(self, data: TrainModelResult): + """ + Logs the details of a machine learning run, including parameters, metrics, models, and artifacts, + to the Sientia tracking system. + + Args: + data: The training model result containing the datasets, model, parameters, + and evaluation metrics. + + Raises: + ValueError: If required metrics or artifacts are missing. + Exception: If MLflow logging fails for any reason. + """ + # Validate that required artifacts exist before attempting to log + if not data.report_path or not path.exists(data.report_path): + error_msg = f'Report file does not exist: {data.report_path}' + self.logger.error(error_msg) + raise ValueError(error_msg) + + if not data.train_data_path or not path.exists(data.train_data_path): + error_msg = f'Training data file does not exist: {data.train_data_path}' + self.logger.error(error_msg) + raise ValueError(error_msg) + + if not data.test_data_path or not path.exists(data.test_data_path): + error_msg = f'Test data file does not exist: {data.test_data_path}' + self.logger.error(error_msg) + raise ValueError(error_msg) + + # Validate that metrics are present + if data.mse_val is None or data.r2_val is None or data.mae_val is None: + error_msg = 'One or more metrics (MSE, R2, MAE) are None' + self.logger.error(error_msg) + raise ValueError(error_msg) + + try: + # Prepare parameters + train_test_split = f'{data.params.train_size}-{100 - data.params.train_size}' + interval_strs = [ + (str(interval[0]), str(interval[1])) + for interval in (data.params.removed_intervals or []) + ] + + # Set experiment and create run + self.model_serving.set_experiment(data.params.experiment_name) + self.logger.info( + f"Logging run '{data.run_name}' to experiment '{data.params.experiment_name}'" + ) + + with self.model_serving.save_experiment( + run_name=data.run_name, description=data.params.experiment_name + ): + # Log model parameters + self.model_serving.log_param('model_type', 'Linear Regression') + self.model_serving.log_param('target_variable', data.params.target_variable) + self.model_serving.log_param('input_variables', data.params.variable_columns) + self.model_serving.log_param('lag_train', data.params.lag_train) + self.model_serving.log_param('lag_val', data.params.lag_val) + self.model_serving.log_param('ma', data.params.window) + self.model_serving.log_param('low_lim', data.params.low_lim) + self.model_serving.log_param('upp_lim', data.params.upp_lim) + self.model_serving.log_param('normalized', data.scaler_dict) + self.model_serving.log_param('ar', data.params.include_ar) + self.model_serving.log_param('Train_test_split', train_test_split) + self.model_serving.log_param('Removed_intervals', interval_strs) + self.model_serving.log_param('Retrain', False) + + # Log evaluation metrics + self.model_serving.log_metric('MSE', data.mse_val) + self.model_serving.log_metric('R2', data.r2_val) + self.model_serving.log_metric('MAE', data.mae_val) + + # Log models + self.model_serving.log_model(data.process_data, 'data_model') + self.model_serving.log_model(data.regr, 'prediction_model') + + # Log artifacts + self.model_serving.log_artifact(data.report_path) + self.model_serving.log_artifact(data.train_data_path) + self.model_serving.log_artifact(data.test_data_path) + + self.logger.info( + f"Successfully logged run '{data.run_name}' with metrics: MSE={data.mse_val:.4f}, R2={data.r2_val:.4f}, MAE={data.mae_val:.4f}" + ) + + except Exception as e: + error_msg = f"Failed to save run '{data.run_name}' to MLflow: {str(e)}" + self.logger.error(error_msg) + raise Exception(error_msg) from e + + def _init_artifacts_data(self, data: TrainModelResult) -> tuple[pd.DataFrame, pd.DataFrame]: + """ + Prepares the reference and current datasets for artifact generation. + + Args: + data: The training model result containing the datasets and model. + + Returns: + tuple: A tuple containing: + - reference_data: The training dataset with predictions added. + - current_data: The testing dataset with predictions added. + + Raises: + ValueError: If training or test datasets are empty or invalid. + AttributeError: If required attributes are missing from the data object. + """ + # Validate that required DataFrames are not empty + # Note: x_train, y_train, x_test, y_test, and regr are required fields in TrainModelResult + # so we only check if they are empty, not None + if data.x_train.empty: + error_msg = 'Training features (x_train) are empty' + self.logger.error(error_msg) + raise ValueError(error_msg) + + if data.y_train.empty: + error_msg = 'Training target (y_train) is empty' + self.logger.error(error_msg) + raise ValueError(error_msg) + + if data.x_test.empty: + error_msg = 'Test features (x_test) are empty' + self.logger.error(error_msg) + raise ValueError(error_msg) + + if data.y_test.empty: + error_msg = 'Test target (y_test) is empty' + self.logger.error(error_msg) + raise ValueError(error_msg) + + # Validate that predictions exist (y_pred is optional, so check for None) + if data.y_pred is None: + error_msg = 'Test predictions (y_pred) are None' + self.logger.error(error_msg) + raise ValueError(error_msg) + + # Prepare reference data (training set) + reference_data = pd.concat([data.x_train, data.y_train], axis=1) + reference_data = reference_data.rename(columns={data.params.target_variable: 'target'}) + reference_data['prediction'] = data.regr.predict(data.x_train) + + # Prepare current data (test set) + current_data = pd.concat([data.x_test, data.y_test], axis=1) + current_data = current_data.rename(columns={data.params.target_variable: 'target'}) + current_data['prediction'] = data.y_pred + + return reference_data, current_data + + def _create_run_directory(self, base_path: str, run_name: str) -> str: + """ + Creates a directory inside the 'reports' folder with the run name and a timestamp. + + Uses microsecond precision in timestamp to minimize collision probability + in high-concurrency scenarios. + + Args: + base_path (str): The path to the 'reports' folder. + run_name (str): The name of the run. + + Returns: + str: The path to the created directory. + + Raises: + PermissionError: If there are insufficient permissions to create the directory. + OSError: If directory creation fails for any other reason. + """ + # Use microsecond precision to reduce collision probability + timestamp = datetime.now().strftime('%Y%m%d_%H%M%S_%f') + run_dir = path.join(base_path, f'{run_name}_{timestamp}') + + try: + makedirs(run_dir, exist_ok=True) + self.logger.info(f'Created run directory: {run_dir}') + return run_dir + except PermissionError as e: + error_msg = f'Permission denied when creating directory: {run_dir}' + self.logger.error(error_msg) + raise PermissionError(error_msg) from e + except OSError as e: + error_msg = f'Failed to create directory {run_dir}: {str(e)}' + self.logger.error(error_msg) + raise OSError(error_msg) from e + + def _setup_run_directory(self, run_dir: str, header_file_path: str): + """ + Creates empty files and copies a header file into the specified run directory. + + Note: Lock removed as each run has its own unique directory, so no synchronization + is needed between different runs. File operations within the same directory are + atomic at the OS level. + + Args: + run_dir (str): The path to the run directory where the files will be created. + header_file_path (str): The path to the header.html file to be copied. + + Raises: + FileNotFoundError: If the header file does not exist. + PermissionError: If there are insufficient permissions to create files. + OSError: If file creation or copying fails for any other reason. + """ + empty_files = ['data_drift.html', 'data_quality.html', 'regression.html'] + + try: + # Create empty placeholder files + for file_name in empty_files: + file_path = path.join(run_dir, file_name) + with open(file_path, 'w'): + pass # Create empty file + + # Copy header file to run directory + header_dest = path.join(run_dir, 'header.html') + shutil.copy(header_file_path, header_dest) + + self.logger.info(f'Run directory setup completed successfully in: {run_dir}') + + except FileNotFoundError as e: + error_msg = f'Header file not found: {header_file_path}' + self.logger.error(error_msg) + raise FileNotFoundError(error_msg) from e + except PermissionError as e: + error_msg = f'Permission denied when setting up directory: {run_dir}' + self.logger.error(error_msg) + raise PermissionError(error_msg) from e + except OSError as e: + error_msg = f'Failed to setup run directory {run_dir}: {str(e)}' + self.logger.error(error_msg) + raise OSError(error_msg) from e + + def _generate_report( + self, reference_data: pd.DataFrame, current_data: pd.DataFrame, data: TrainModelResult + ) -> TrainModelResult: + """ + Generates a comprehensive report summarizing data quality, data drift, and regression analysis. + + Args: + reference_data (pd.DataFrame): The training dataset with predictions added. + current_data (pd.DataFrame): The testing dataset with predictions added. + data: The training model result containing the datasets, model, and parameters. + + Returns: + The updated result object with paths to the generated report and data files. + + Raises: + ValueError: If data conversion to float64 fails or DataFrames are invalid. + PermissionError: If there are insufficient permissions to write files. + OSError: If file writing fails for any other reason. + """ + try: + # Convert data to float64 for report generation + # This may raise ValueError if data contains non-numeric values + reference_data_float = reference_data.astype(np.float64) + current_data_float = current_data.astype(np.float64) + + # Initialize report generator + report = Reports( + reference_data=reference_data_float, + current_data=current_data_float, + base_path=data.run_dir, + ) + + # Generate report sections + report.add_data_quality_section(columns=data.params.variable_columns + ['target']) + report.add_data_drift_section(columns=data.params.variable_columns + ['target']) + report.add_regression_section() + + # Validate that run_dir is set (should be set by _create_run_directory) + if not data.run_dir: + error_msg = 'run_dir is not set after directory creation' + self.logger.error(error_msg) + raise ValueError(error_msg) + + # Save HTML report + data.report_path = path.join(data.run_dir, 'report.html') + report.save_all_sections_html(data.report_path) + self.logger.info(f'Generated HTML report: {data.report_path}') + + # Save training data CSV + data.train_data_path = path.join(data.run_dir, 'train_data.csv') + reference_data.to_csv(data.train_data_path, index=False) + self.logger.info(f'Saved training data: {data.train_data_path}') + + # Save test data CSV + data.test_data_path = path.join(data.run_dir, 'test_data.csv') + current_data.to_csv(data.test_data_path, index=False) + self.logger.info(f'Saved test data: {data.test_data_path}') + + return data + + except ValueError as e: + error_msg = f'Failed to convert data to float64 for report generation: {str(e)}' + self.logger.error(error_msg) + raise ValueError(error_msg) from e + except PermissionError as e: + error_msg = f'Permission denied when writing report files to: {data.run_dir}' + self.logger.error(error_msg) + raise PermissionError(error_msg) from e + except OSError as e: + error_msg = f'Failed to generate report in {data.run_dir}: {str(e)}' + self.logger.error(error_msg) + raise OSError(error_msg) from e + + def _get_reports_directory(self) -> str: + """ + Get the absolute path to the reports directory. + + Returns: + str: Absolute path to model_manager/reports directory. + """ + # Get the directory where this file is located (model_manager/utils/repository/) + current_file_dir = path.dirname(path.abspath(__file__)) + # Navigate up to model_manager/ and then to reports/ + model_manager_dir = path.dirname(path.dirname(current_file_dir)) + reports_dir = path.join(model_manager_dir, 'reports') + return reports_dir diff --git a/tests/utils/repository/test_model_repository.py b/tests/utils/repository/test_model_repository.py index a0c013f..c1e0197 100644 --- a/tests/utils/repository/test_model_repository.py +++ b/tests/utils/repository/test_model_repository.py @@ -497,3 +497,680 @@ def test_update_production_model(mlflow_repository): 'mlflow_run_id': '0', 'mlflow_experiment_id': '0', } + + +# ========== Tests for Model Artifact Generation Methods ========== + + +def test_get_next_run_name_new(mlflow_repository): + """Test get_next_run_name generates correct run name based on existing runs.""" + mlflow_repository.model_serving.search_runs_by_name.return_value = [ + MagicMock(), + MagicMock(), + MagicMock(), + ] + + result = mlflow_repository.get_next_run_name_new('test_experiment') + + mlflow_repository.model_serving.search_runs_by_name.assert_called_once_with( + experiment_names=['test_experiment'], order_by=['start_time desc'] + ) + assert result == 'test_experiment-4' + + +def test_get_next_run_name_new_first_run(mlflow_repository): + """Test get_next_run_name for first run (no existing runs).""" + mlflow_repository.model_serving.search_runs_by_name.return_value = [] + + result = mlflow_repository.get_next_run_name_new('new_experiment') + + assert result == 'new_experiment-1' + + +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_artifacts_success(mock_path, mlflow_repository): + """Test generate_artifacts successfully creates all artifacts.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + # Mock data + params = MagicMock(spec=TrainModelParams) + params.target_variable = 'target' + params.variable_columns = ['feat1', 'feat2'] + params.experiment_name = 'test_exp' + + data = MagicMock(spec=TrainModelResult) + data.run_name = 'test_run-1' + data.params = params + data.x_train = DataFrame({'feat1': [1, 2], 'feat2': [3, 4]}) + data.y_train = DataFrame({'target': [5, 6]}) + data.x_test = DataFrame({'feat1': [7, 8], 'feat2': [9, 10]}) + data.y_test = DataFrame({'target': [11, 12]}) + data.regr = MagicMock() + data.regr.predict = MagicMock(return_value=np.array([5.1, 6.1])) + data.y_pred = np.array([11.1, 12.1]) + + # Mock path operations + mock_path.exists.return_value = True + mock_path.join.side_effect = lambda *args: '/'.join(args) + + # Mock private methods + mlflow_repository._get_reports_directory = MagicMock(return_value='/reports') + mlflow_repository._create_run_directory = MagicMock(return_value='/reports/test_run-1_20231010') + mlflow_repository._setup_run_directory = MagicMock() + mlflow_repository._generate_report = MagicMock(return_value=data) + + result = mlflow_repository.generate_artifacts(data) + + # Assertions + mlflow_repository._get_reports_directory.assert_called_once() + mlflow_repository._create_run_directory.assert_called_once_with('/reports', 'test_run-1') + mlflow_repository._setup_run_directory.assert_called_once() + mlflow_repository._generate_report.assert_called_once() + assert result == data + + +def test_generate_artifacts_missing_run_name(mlflow_repository): + """Test generate_artifacts raises ValueError when run_name is not set.""" + from model_manager.utils.models.train_model_result import TrainModelResult + + data = MagicMock(spec=TrainModelResult) + data.run_name = None + + with pytest.raises(ValueError) as exc_info: + mlflow_repository.generate_artifacts(data) + + assert 'run_name must be set before generating artifacts' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_artifacts_reports_directory_not_exists(mock_path, mlflow_repository): + """Test generate_artifacts raises FileNotFoundError when reports directory doesn't exist.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.target_variable = 'target' + + data = MagicMock(spec=TrainModelResult) + data.run_name = 'test_run-1' + data.params = params + data.x_train = DataFrame({'feat1': [1]}) + data.y_train = DataFrame({'target': [2]}) + data.x_test = DataFrame({'feat1': [3]}) + data.y_test = DataFrame({'target': [4]}) + data.regr = MagicMock() + data.y_pred = np.array([4.1]) + + mlflow_repository._get_reports_directory = MagicMock(return_value='/reports') + mock_path.exists.return_value = False + + with pytest.raises(FileNotFoundError) as exc_info: + mlflow_repository.generate_artifacts(data) + + assert 'Reports directory does not exist' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_artifacts_header_file_not_exists(mock_path, mlflow_repository): + """Test generate_artifacts raises FileNotFoundError when header.html doesn't exist.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.target_variable = 'target' + + data = MagicMock(spec=TrainModelResult) + data.run_name = 'test_run-1' + data.params = params + data.x_train = DataFrame({'feat1': [1]}) + data.y_train = DataFrame({'target': [2]}) + data.x_test = DataFrame({'feat1': [3]}) + data.y_test = DataFrame({'target': [4]}) + data.regr = MagicMock() + data.regr.predict = MagicMock(return_value=np.array([2.1])) + data.y_pred = np.array([4.1]) + + mlflow_repository._get_reports_directory = MagicMock(return_value='/reports') + mlflow_repository._create_run_directory = MagicMock(return_value='/reports/test_run-1_20231010') + + # First call returns True (reports dir exists), second returns False (header.html doesn't exist) + mock_path.exists.side_effect = [True, False] + mock_path.join.side_effect = lambda *args: '/'.join(args) + + with pytest.raises(FileNotFoundError) as exc_info: + mlflow_repository.generate_artifacts(data) + + assert 'Header file does not exist' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.path') +def test_save_run_success(mock_path, mlflow_repository): + """Test save_run successfully logs all parameters, metrics, models, and artifacts.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.train_size = 80 + params.removed_intervals = [(1, 10), (20, 30)] + params.experiment_name = 'test_exp' + params.target_variable = 'target' + params.variable_columns = ['feat1', 'feat2'] + params.lag_train = 5 + params.lag_val = 3 + params.window = 10 + params.low_lim = 0.0 + params.upp_lim = 1.0 + params.include_ar = True + + data = MagicMock(spec=TrainModelResult) + data.run_name = 'test_run-1' + data.params = params + data.report_path = '/reports/report.html' + data.train_data_path = '/reports/train.csv' + data.test_data_path = '/reports/test.csv' + data.mse_val = 0.123 + data.r2_val = 0.987 + data.mae_val = 0.456 + data.scaler_dict = {'scaler': 'minmax'} + data.process_data = MagicMock() + data.regr = MagicMock() + + mock_path.exists.return_value = True + + mlflow_repository.save_run(data) + + # Verify experiment was set + mlflow_repository.model_serving.set_experiment.assert_called_once_with('test_exp') + + # Verify parameters were logged + assert mlflow_repository.model_serving.log_param.call_count == 13 + + # Verify metrics were logged + mlflow_repository.model_serving.log_metric.assert_any_call('MSE', 0.123) + mlflow_repository.model_serving.log_metric.assert_any_call('R2', 0.987) + mlflow_repository.model_serving.log_metric.assert_any_call('MAE', 0.456) + + # Verify models were logged + mlflow_repository.model_serving.log_model.assert_any_call(data.process_data, 'data_model') + mlflow_repository.model_serving.log_model.assert_any_call(data.regr, 'prediction_model') + + # Verify artifacts were logged + mlflow_repository.model_serving.log_artifact.assert_any_call('/reports/report.html') + mlflow_repository.model_serving.log_artifact.assert_any_call('/reports/train.csv') + mlflow_repository.model_serving.log_artifact.assert_any_call('/reports/test.csv') + + +@patch('model_manager.utils.repository.model_repository.path') +def test_save_run_missing_report_path(mock_path, mlflow_repository): + """Test save_run raises ValueError when report_path is missing.""" + from model_manager.utils.models.train_model_result import TrainModelResult + + data = MagicMock(spec=TrainModelResult) + data.report_path = None + + with pytest.raises(ValueError) as exc_info: + mlflow_repository.save_run(data) + + assert 'Report file does not exist' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.path') +def test_save_run_missing_metrics(mock_path, mlflow_repository): + """Test save_run raises ValueError when metrics are None.""" + from model_manager.utils.models.train_model_result import TrainModelResult + + data = MagicMock(spec=TrainModelResult) + data.report_path = '/reports/report.html' + data.train_data_path = '/reports/train.csv' + data.test_data_path = '/reports/test.csv' + data.mse_val = None + data.r2_val = 0.987 + data.mae_val = 0.456 + + mock_path.exists.return_value = True + + with pytest.raises(ValueError) as exc_info: + mlflow_repository.save_run(data) + + assert 'One or more metrics (MSE, R2, MAE) are None' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.path') +def test_save_run_mlflow_error(mock_path, mlflow_repository): + """Test save_run handles MLflow errors gracefully.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.train_size = 80 + params.removed_intervals = [] + params.experiment_name = 'test_exp' + + data = MagicMock(spec=TrainModelResult) + data.run_name = 'test_run-1' + data.params = params + data.report_path = '/reports/report.html' + data.train_data_path = '/reports/train.csv' + data.test_data_path = '/reports/test.csv' + data.mse_val = 0.123 + data.r2_val = 0.987 + data.mae_val = 0.456 + + mock_path.exists.return_value = True + mlflow_repository.model_serving.set_experiment.side_effect = Exception( + 'MLflow connection error' + ) + + with pytest.raises(Exception) as exc_info: + mlflow_repository.save_run(data) + + assert 'Failed to save run' in str(exc_info.value) + assert 'MLflow connection error' in str(exc_info.value) + + +# ========== Additional Tests for 100% Coverage ========== + + +def test_init_artifacts_data_empty_x_train(mlflow_repository): + """Test _init_artifacts_data raises ValueError when x_train is empty.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + data = MagicMock(spec=TrainModelResult) + data.params = params + data.x_train = DataFrame() # Empty DataFrame + data.y_train = DataFrame({'target': [1]}) + data.x_test = DataFrame({'feat1': [1]}) + data.y_test = DataFrame({'target': [1]}) + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._init_artifacts_data(data) + + assert 'Training features (x_train) are empty' in str(exc_info.value) + + +def test_init_artifacts_data_empty_y_train(mlflow_repository): + """Test _init_artifacts_data raises ValueError when y_train is empty.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + data = MagicMock(spec=TrainModelResult) + data.params = params + data.x_train = DataFrame({'feat1': [1]}) + data.y_train = DataFrame() # Empty DataFrame + data.x_test = DataFrame({'feat1': [1]}) + data.y_test = DataFrame({'target': [1]}) + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._init_artifacts_data(data) + + assert 'Training target (y_train) is empty' in str(exc_info.value) + + +def test_init_artifacts_data_empty_x_test(mlflow_repository): + """Test _init_artifacts_data raises ValueError when x_test is empty.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + data = MagicMock(spec=TrainModelResult) + data.params = params + data.x_train = DataFrame({'feat1': [1]}) + data.y_train = DataFrame({'target': [1]}) + data.x_test = DataFrame() # Empty DataFrame + data.y_test = DataFrame({'target': [1]}) + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._init_artifacts_data(data) + + assert 'Test features (x_test) are empty' in str(exc_info.value) + + +def test_init_artifacts_data_empty_y_test(mlflow_repository): + """Test _init_artifacts_data raises ValueError when y_test is empty.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + data = MagicMock(spec=TrainModelResult) + data.params = params + data.x_train = DataFrame({'feat1': [1]}) + data.y_train = DataFrame({'target': [1]}) + data.x_test = DataFrame({'feat1': [1]}) + data.y_test = DataFrame() # Empty DataFrame + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._init_artifacts_data(data) + + assert 'Test target (y_test) is empty' in str(exc_info.value) + + +def test_init_artifacts_data_none_y_pred(mlflow_repository): + """Test _init_artifacts_data raises ValueError when y_pred is None.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + data = MagicMock(spec=TrainModelResult) + data.params = params + data.x_train = DataFrame({'feat1': [1]}) + data.y_train = DataFrame({'target': [1]}) + data.x_test = DataFrame({'feat1': [1]}) + data.y_test = DataFrame({'target': [1]}) + data.y_pred = None + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._init_artifacts_data(data) + + assert 'Test predictions (y_pred) are None' in str(exc_info.value) + + +def test_init_artifacts_data_success(mlflow_repository): + """Test _init_artifacts_data successfully prepares data.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.target_variable = 'target' + + data = MagicMock(spec=TrainModelResult) + data.params = params + data.x_train = DataFrame({'feat1': [1, 2]}) + data.y_train = DataFrame({'target': [3, 4]}) + data.x_test = DataFrame({'feat1': [5, 6]}) + data.y_test = DataFrame({'target': [7, 8]}) + data.regr = MagicMock() + data.regr.predict = MagicMock(return_value=np.array([3.1, 4.1])) + data.y_pred = np.array([7.1, 8.1]) + + reference_data, current_data = mlflow_repository._init_artifacts_data(data) + + assert 'target' in reference_data.columns + assert 'prediction' in reference_data.columns + assert 'target' in current_data.columns + assert 'prediction' in current_data.columns + assert len(reference_data) == 2 + assert len(current_data) == 2 + + +@patch('model_manager.utils.repository.model_repository.makedirs') +@patch('model_manager.utils.repository.model_repository.path') +def test_create_run_directory_success(mock_path, mock_makedirs, mlflow_repository): + """Test _create_run_directory successfully creates directory.""" + mock_path.join.return_value = '/reports/test_run_20231010_123456_123456' + + result = mlflow_repository._create_run_directory('/reports', 'test_run') + + mock_makedirs.assert_called_once_with('/reports/test_run_20231010_123456_123456', exist_ok=True) + assert result == '/reports/test_run_20231010_123456_123456' + + +@patch('model_manager.utils.repository.model_repository.makedirs') +@patch('model_manager.utils.repository.model_repository.path') +def test_create_run_directory_permission_error(mock_path, mock_makedirs, mlflow_repository): + """Test _create_run_directory handles PermissionError.""" + mock_path.join.return_value = '/reports/test_run_20231010' + mock_makedirs.side_effect = PermissionError('Permission denied') + + with pytest.raises(PermissionError) as exc_info: + mlflow_repository._create_run_directory('/reports', 'test_run') + + assert 'Permission denied when creating directory' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.makedirs') +@patch('model_manager.utils.repository.model_repository.path') +def test_create_run_directory_os_error(mock_path, mock_makedirs, mlflow_repository): + """Test _create_run_directory handles OSError.""" + mock_path.join.return_value = '/reports/test_run_20231010' + mock_makedirs.side_effect = OSError('Disk full') + + with pytest.raises(OSError) as exc_info: + mlflow_repository._create_run_directory('/reports', 'test_run') + + assert 'Failed to create directory' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.shutil') +@patch('model_manager.utils.repository.model_repository.path') +def test_setup_run_directory_success(mock_path, mock_shutil, mlflow_repository): + """Test _setup_run_directory successfully sets up directory.""" + mock_path.join.side_effect = lambda *args: '/'.join(args) + mock_open = MagicMock() + + with patch('builtins.open', mock_open): + mlflow_repository._setup_run_directory('/run_dir', '/reports/header.html') + + assert mock_open.call_count == 3 # 3 empty files + mock_shutil.copy.assert_called_once_with('/reports/header.html', '/run_dir/header.html') + + +@patch('model_manager.utils.repository.model_repository.shutil') +@patch('model_manager.utils.repository.model_repository.path') +def test_setup_run_directory_file_not_found(mock_path, mock_shutil, mlflow_repository): + """Test _setup_run_directory handles FileNotFoundError.""" + mock_path.join.side_effect = lambda *args: '/'.join(args) + mock_shutil.copy.side_effect = FileNotFoundError('Header not found') + + mock_open = MagicMock() + with patch('builtins.open', mock_open): + with pytest.raises(FileNotFoundError) as exc_info: + mlflow_repository._setup_run_directory('/run_dir', '/reports/header.html') + + assert 'Header file not found' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.shutil') +@patch('model_manager.utils.repository.model_repository.path') +def test_setup_run_directory_permission_error(mock_path, mock_shutil, mlflow_repository): + """Test _setup_run_directory handles PermissionError.""" + mock_path.join.side_effect = lambda *args: '/'.join(args) + + mock_open = MagicMock() + mock_open.side_effect = PermissionError('Permission denied') + + with patch('builtins.open', mock_open): + with pytest.raises(PermissionError) as exc_info: + mlflow_repository._setup_run_directory('/run_dir', '/reports/header.html') + + assert 'Permission denied when setting up directory' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.shutil') +@patch('model_manager.utils.repository.model_repository.path') +def test_setup_run_directory_os_error(mock_path, mock_shutil, mlflow_repository): + """Test _setup_run_directory handles OSError.""" + mock_path.join.side_effect = lambda *args: '/'.join(args) + + mock_open = MagicMock() + mock_open.side_effect = OSError('Disk error') + + with patch('builtins.open', mock_open): + with pytest.raises(OSError) as exc_info: + mlflow_repository._setup_run_directory('/run_dir', '/reports/header.html') + + assert 'Failed to setup run directory' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.Reports') +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_report_success(mock_path, mock_reports, mlflow_repository): + """Test _generate_report successfully generates all reports.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.variable_columns = ['feat1', 'feat2'] + + data = MagicMock(spec=TrainModelResult) + data.params = params + data.run_dir = '/run_dir' + + reference_data = DataFrame( + {'feat1': [1.0], 'feat2': [2.0], 'target': [3.0], 'prediction': [3.1]} + ) + current_data = DataFrame({'feat1': [4.0], 'feat2': [5.0], 'target': [6.0], 'prediction': [6.1]}) + + mock_path.join.side_effect = lambda *args: '/'.join(args) + mock_report_instance = MagicMock() + mock_reports.return_value = mock_report_instance + + # Mock DataFrame.to_csv to avoid actual file writing + with patch.object(DataFrame, 'to_csv'): + result = mlflow_repository._generate_report(reference_data, current_data, data) + + mock_reports.assert_called_once() + mock_report_instance.add_data_quality_section.assert_called_once() + mock_report_instance.add_data_drift_section.assert_called_once() + mock_report_instance.add_regression_section.assert_called_once() + mock_report_instance.save_all_sections_html.assert_called_once() + assert result == data + + +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_report_value_error(mock_path, mlflow_repository): + """Test _generate_report handles ValueError from data conversion.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + data = MagicMock(spec=TrainModelResult) + data.params = params + data.run_dir = '/run_dir' + + # DataFrame with non-numeric data + reference_data = DataFrame({'feat1': ['a', 'b']}) + current_data = DataFrame({'feat1': ['c', 'd']}) + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._generate_report(reference_data, current_data, data) + + assert 'Failed to convert data to float64' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.Reports') +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_report_permission_error(mock_path, mock_reports, mlflow_repository): + """Test _generate_report handles PermissionError.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.variable_columns = ['feat1'] + + data = MagicMock(spec=TrainModelResult) + data.params = params + data.run_dir = '/run_dir' + + reference_data = DataFrame({'feat1': [1.0]}) + current_data = DataFrame({'feat1': [2.0]}) + + mock_path.join.side_effect = lambda *args: '/'.join(args) + mock_report_instance = MagicMock() + mock_reports.return_value = mock_report_instance + mock_report_instance.save_all_sections_html.side_effect = PermissionError('Permission denied') + + with pytest.raises(PermissionError) as exc_info: + mlflow_repository._generate_report(reference_data, current_data, data) + + assert 'Permission denied when writing report files' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.Reports') +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_report_os_error(mock_path, mock_reports, mlflow_repository): + """Test _generate_report handles OSError.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.variable_columns = ['feat1'] + + data = MagicMock(spec=TrainModelResult) + data.params = params + data.run_dir = '/run_dir' + + reference_data = DataFrame({'feat1': [1.0]}) + current_data = DataFrame({'feat1': [2.0]}) + + mock_path.join.side_effect = lambda *args: '/'.join(args) + mock_report_instance = MagicMock() + mock_reports.return_value = mock_report_instance + mock_report_instance.save_all_sections_html.side_effect = OSError('Disk error') + + with pytest.raises(OSError) as exc_info: + mlflow_repository._generate_report(reference_data, current_data, data) + + assert 'Failed to generate report' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.Reports') +@patch('model_manager.utils.repository.model_repository.path') +def test_generate_report_run_dir_none(mock_path, mock_reports, mlflow_repository): + """Test _generate_report raises ValueError when run_dir is None.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.variable_columns = ['feat1'] + + data = MagicMock(spec=TrainModelResult) + data.params = params + data.run_dir = None # Not set + + reference_data = DataFrame({'feat1': [1.0]}) + current_data = DataFrame({'feat1': [2.0]}) + + mock_report_instance = MagicMock() + mock_reports.return_value = mock_report_instance + + with pytest.raises(ValueError) as exc_info: + mlflow_repository._generate_report(reference_data, current_data, data) + + assert 'run_dir is not set after directory creation' in str(exc_info.value) + + +def test_get_reports_directory(mlflow_repository): + """Test _get_reports_directory returns correct path.""" + result = mlflow_repository._get_reports_directory() + + assert result.endswith('model_manager/reports') + assert 'model_manager' in result + + +@patch('model_manager.utils.repository.model_repository.path') +def test_save_run_missing_train_data_path(mock_path, mlflow_repository): + """Test save_run raises ValueError when train_data_path is missing.""" + from model_manager.utils.models.train_model_result import TrainModelResult + + data = MagicMock(spec=TrainModelResult) + data.report_path = '/reports/report.html' + data.train_data_path = None + + mock_path.exists.return_value = True + + with pytest.raises(ValueError) as exc_info: + mlflow_repository.save_run(data) + + assert 'Training data file does not exist' in str(exc_info.value) + + +@patch('model_manager.utils.repository.model_repository.path') +def test_save_run_missing_test_data_path(mock_path, mlflow_repository): + """Test save_run raises ValueError when test_data_path is missing.""" + from model_manager.utils.models.train_model_result import TrainModelResult + + data = MagicMock(spec=TrainModelResult) + data.report_path = '/reports/report.html' + data.train_data_path = '/reports/train.csv' + data.test_data_path = None + + mock_path.exists.return_value = True + + with pytest.raises(ValueError) as exc_info: + mlflow_repository.save_run(data) + + assert 'Test data file does not exist' in str(exc_info.value) From b3c749872c84221121e3f5977a0f04682a787ee0 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Mon, 13 Oct 2025 10:18:45 -0300 Subject: [PATCH 4/5] SIENTIAPDE-1252: Implement save_model activity to save trained models to MLflow with comprehensive error handling and add corresponding unit tests. --- model_manager/activities/mlflow.py | 104 ++++++++++++ tests/activities/test_mlflow.py | 248 +++++++++++++++++++++++++++++ 2 files changed, 352 insertions(+) diff --git a/model_manager/activities/mlflow.py b/model_manager/activities/mlflow.py index 9e5a900..0458c05 100644 --- a/model_manager/activities/mlflow.py +++ b/model_manager/activities/mlflow.py @@ -344,3 +344,107 @@ class MLFlow(BaseActivity): ) self.error(trace, metadata=metadata) raise e + + @activity.defn(name='save_model') + async def save_model(self, input_data: dict[str, Any]) -> dict[str, Any]: + """ + Save a trained ML model and its artifacts to MLflow with comprehensive error handling. + + This activity orchestrates the complete model saving pipeline: + 1. Generates the next run name for the experiment + 2. Creates and organizes artifacts (reports, data files) + 3. Logs model, parameters, metrics, and artifacts to MLflow + 4. Returns success/failure status with results or error message + + The activity does NOT raise exceptions on failure - it catches all errors, + sends notifications, and returns a failure status. This allows the workflow + to handle the error gracefully and update the database accordingly. + + Args: + input_data: Configuration for model saving operation + Required keys: + - metadata (dict): Workflow execution metadata + - train_result (TrainModelResult): Training result with model and metrics + + Returns: + dict: Save result with the following structure: + { + 'success': bool, # True if saving succeeded, False otherwise + 'result': TrainModelResult | None, # Updated result if success=True + 'error_message': str | None # Error message if success=False + } + + Example: + # Successful save + result = await save_model({ + 'metadata': {'workflow_id': 'save-123', 'experiment_run_id': 456}, + 'train_result': TrainModelResult(...) + }) + # Returns: {'success': True, 'result': TrainModelResult(...), 'error_message': None} + + # Failed save + # Returns: {'success': False, 'result': None, 'error_message': 'Error details...'} + """ + metadata = input_data.get('metadata', {}) + train_result = input_data['train_result'] + + try: + experiment_name = train_result.params.experiment_name + + self.info( + f'Starting model save for experiment: {experiment_name}', + metadata, + ) + + # Step 1: Generate next run name + self.info('Generating run name', metadata) + train_result.run_name = self.model_monitoring_repository.get_next_run_name( + experiment_name + ) + self.info(f'Generated run name: {train_result.run_name}', metadata) + + # Step 2: Generate artifacts (reports, CSV files) + self.info('Generating artifacts', metadata) + train_result = self.model_monitoring_repository.generate_artifacts(train_result) + self.info('Artifacts generated successfully', metadata) + + # Step 3: Save run to MLflow + self.info('Saving run to MLflow', metadata) + self.model_monitoring_repository.save_run(train_result) + + self.info( + f'Model saved successfully - Run: {train_result.run_name}, ' + f'Experiment: {experiment_name}', + metadata, + ) + + return { + 'success': True, + 'result': train_result, + 'error_message': None, + } + + except Exception as e: # noqa: BLE001 + error_msg = f'Error saving model - Experiment: {train_result.params.experiment_name if train_result and train_result.params else "unknown"}, Error: {str(e)}' + trace = traceback.format_exc() + + # Send notification (MongoDB) + self.send_notification( + metadata=metadata, + notification_id='SAVE_MODEL_ERROR', + message=error_msg, + block='save_model', + level=NotificationLevel.ERROR, + attachment_content=trace, + ) + + # Log error with metadata + self.error(trace, metadata=metadata) + + # Return failure result (do NOT raise exception) + # This allows workflow to update database with error status + return { + 'success': False, + 'result': None, + 'error_message': str(e), + } diff --git a/tests/activities/test_mlflow.py b/tests/activities/test_mlflow.py index 80b7e85..779f9d4 100644 --- a/tests/activities/test_mlflow.py +++ b/tests/activities/test_mlflow.py @@ -299,3 +299,251 @@ async def test_update_production_model_error(mlflow): ) else: raise AssertionError('No exception raised') + + +@mark.asyncio +async def test_save_model_success(mlflow): + """Test save_model successfully saves model and artifacts to MLflow.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + # Mock train result + params = MagicMock(spec=TrainModelParams) + params.experiment_name = 'test_experiment' + + train_result = MagicMock(spec=TrainModelResult) + train_result.params = params + train_result.run_name = None # Will be set by get_next_run_name + + # Mock repository methods + mlflow.model_monitoring_repository.get_next_run_name.return_value = 'test_experiment-1' + mlflow.model_monitoring_repository.generate_artifacts.return_value = train_result + mlflow.model_monitoring_repository.save_run.return_value = None + + input_data = { + **metadata, + 'train_result': train_result, + } + + # Call the method + response = await mlflow.save_model(input_data) + + # Verify repository methods were called + mlflow.model_monitoring_repository.get_next_run_name.assert_called_once_with('test_experiment') + mlflow.model_monitoring_repository.generate_artifacts.assert_called_once_with(train_result) + mlflow.model_monitoring_repository.save_run.assert_called_once_with(train_result) + + # Verify response + assert response['success'] is True + assert response['result'] == train_result + assert response['error_message'] is None + assert train_result.run_name == 'test_experiment-1' + + +@mark.asyncio +async def test_save_model_get_next_run_name_error(mlflow): + """Test save_model handles error during get_next_run_name.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.experiment_name = 'test_experiment' + + train_result = MagicMock(spec=TrainModelResult) + train_result.params = params + + # Mock error in get_next_run_name + mlflow.model_monitoring_repository.get_next_run_name.side_effect = Exception( + 'MLflow connection error' + ) + + input_data = { + **metadata, + 'train_result': train_result, + } + + # Call the method + response = await mlflow.save_model(input_data) + + # Verify error handling + assert response['success'] is False + assert response['result'] is None + assert 'MLflow connection error' in response['error_message'] + + # Verify notification was sent + mlflow.send_notification.assert_called_once_with( + metadata=metadata['metadata'], + notification_id='SAVE_MODEL_ERROR', + message=ANY, + block='save_model', + level=NotificationLevel.ERROR, + attachment_content=ANY, + ) + + +@mark.asyncio +async def test_save_model_generate_artifacts_error(mlflow): + """Test save_model handles error during generate_artifacts.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.experiment_name = 'test_experiment' + + train_result = MagicMock(spec=TrainModelResult) + train_result.params = params + + # Mock successful get_next_run_name but error in generate_artifacts + mlflow.model_monitoring_repository.get_next_run_name.return_value = 'test_experiment-1' + mlflow.model_monitoring_repository.generate_artifacts.side_effect = FileNotFoundError( + 'Reports directory does not exist' + ) + + input_data = { + **metadata, + 'train_result': train_result, + } + + # Call the method + response = await mlflow.save_model(input_data) + + # Verify error handling + assert response['success'] is False + assert response['result'] is None + assert 'Reports directory does not exist' in response['error_message'] + + # Verify notification was sent + mlflow.send_notification.assert_called_once_with( + metadata=metadata['metadata'], + notification_id='SAVE_MODEL_ERROR', + message=ANY, + block='save_model', + level=NotificationLevel.ERROR, + attachment_content=ANY, + ) + + +@mark.asyncio +async def test_save_model_save_run_error(mlflow): + """Test save_model handles error during save_run.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.experiment_name = 'test_experiment' + + train_result = MagicMock(spec=TrainModelResult) + train_result.params = params + + # Mock successful get_next_run_name and generate_artifacts but error in save_run + mlflow.model_monitoring_repository.get_next_run_name.return_value = 'test_experiment-1' + mlflow.model_monitoring_repository.generate_artifacts.return_value = train_result + mlflow.model_monitoring_repository.save_run.side_effect = ValueError( + 'One or more metrics (MSE, R2, MAE) are None' + ) + + input_data = { + **metadata, + 'train_result': train_result, + } + + # Call the method + response = await mlflow.save_model(input_data) + + # Verify error handling + assert response['success'] is False + assert response['result'] is None + assert 'One or more metrics (MSE, R2, MAE) are None' in response['error_message'] + + # Verify notification was sent + mlflow.send_notification.assert_called_once_with( + metadata=metadata['metadata'], + notification_id='SAVE_MODEL_ERROR', + message=ANY, + block='save_model', + level=NotificationLevel.ERROR, + attachment_content=ANY, + ) + + +@mark.asyncio +async def test_save_model_missing_metadata(mlflow): + """Test save_model handles missing metadata gracefully.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.experiment_name = 'test_experiment' + + train_result = MagicMock(spec=TrainModelResult) + train_result.params = params + + # Mock repository methods + mlflow.model_monitoring_repository.get_next_run_name.return_value = 'test_experiment-1' + mlflow.model_monitoring_repository.generate_artifacts.return_value = train_result + mlflow.model_monitoring_repository.save_run.return_value = None + + # Input data without metadata + input_data = { + 'train_result': train_result, + } + + # Call the method + response = await mlflow.save_model(input_data) + + # Verify it still works (metadata defaults to {}) + assert response['success'] is True + assert response['result'] == train_result + assert response['error_message'] is None + + +@mark.asyncio +async def test_save_model_complete_flow(mlflow): + """Test save_model complete flow with all steps.""" + from model_manager.utils.models.train_model_params import TrainModelParams + from model_manager.utils.models.train_model_result import TrainModelResult + + params = MagicMock(spec=TrainModelParams) + params.experiment_name = 'production_model' + + train_result = MagicMock(spec=TrainModelResult) + train_result.params = params + train_result.run_name = None + train_result.run_dir = None + train_result.report_path = None + + # Mock complete flow + mlflow.model_monitoring_repository.get_next_run_name.return_value = 'production_model-5' + + # After generate_artifacts, paths should be set + updated_result = MagicMock(spec=TrainModelResult) + updated_result.params = params + updated_result.run_name = 'production_model-5' + updated_result.run_dir = '/reports/production_model-5_20231010' + updated_result.report_path = '/reports/production_model-5_20231010/report.html' + updated_result.train_data_path = '/reports/production_model-5_20231010/train_data.csv' + updated_result.test_data_path = '/reports/production_model-5_20231010/test_data.csv' + + mlflow.model_monitoring_repository.generate_artifacts.return_value = updated_result + mlflow.model_monitoring_repository.save_run.return_value = None + + input_data = { + **metadata, + 'train_result': train_result, + } + + # Call the method + response = await mlflow.save_model(input_data) + + # Verify complete flow + mlflow.model_monitoring_repository.get_next_run_name.assert_called_once_with('production_model') + mlflow.model_monitoring_repository.generate_artifacts.assert_called_once() + mlflow.model_monitoring_repository.save_run.assert_called_once_with(updated_result) + + # Verify response + assert response['success'] is True + assert response['result'] == updated_result + assert response['error_message'] is None + assert updated_result.run_name == 'production_model-5' + assert updated_result.run_dir is not None + assert updated_result.report_path is not None From 9d4ec195863dd6490f36dfbb8ec724eca0f4f559 Mon Sep 17 00:00:00 2001 From: Bruno Domingues Date: Mon, 13 Oct 2025 10:54:38 -0300 Subject: [PATCH 5/5] SIENTIAPDE-1252: Improve caching in quality gate, add alt text to logo, and change exception type in MLflow repository. --- .github/workflows/quality-gate.yml | 12 ++++++++---- model_manager/reports/header.html | 1 + model_manager/utils/repository/model_repository.py | 2 +- tests/utils/repository/test_model_repository.py | 2 +- 4 files changed, 11 insertions(+), 6 deletions(-) diff --git a/.github/workflows/quality-gate.yml b/.github/workflows/quality-gate.yml index 1708d52..7add2b0 100644 --- a/.github/workflows/quality-gate.yml +++ b/.github/workflows/quality-gate.yml @@ -196,10 +196,14 @@ jobs: uses: actions/setup-python@v5 with: python-version: "3.11" - cache: 'pip' - cache-dependency-path: | - requirements.txt - requirements-dev.txt + + - name: 💾 Cache pip packages + uses: actions/cache@v4 + with: + path: ~/.cache/pip + key: ${{ runner.os }}-pip-${{ hashFiles('requirements.txt', 'requirements-dev.txt') }} + restore-keys: | + ${{ runner.os }}-pip- - name: 📦 Install Dependencies run: | diff --git a/model_manager/reports/header.html b/model_manager/reports/header.html index fd646a2..cd73a20 100644 --- a/model_manager/reports/header.html +++ b/model_manager/reports/header.html @@ -110,6 +110,7 @@