diff --git a/model_manager/activities/activities.py b/model_manager/activities/activities.py index 66d1bdd..96f8b20 100644 --- a/model_manager/activities/activities.py +++ b/model_manager/activities/activities.py @@ -123,7 +123,6 @@ class Activities(ExperimentTracking, Training, Cleanup): Cleanup.__init__( self, - minio_repository=self.minio_repository, logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller, diff --git a/model_manager/activities/cleanup.py b/model_manager/activities/cleanup.py index 81d85d8..3bc6bce 100644 --- a/model_manager/activities/cleanup.py +++ b/model_manager/activities/cleanup.py @@ -21,7 +21,6 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.observability.logger import Logger from sientia_do.observability.metrics_controller import MetricsController from sientia_do.observability.sientia_monitoring import SientiaMonitoring - from sientia_do.repository.minio_repository import MinioRepository from model_manager.metrics import ACTIVITY_EXECUTION_TOTAL, WORKFLOW_EXECUTION_TOTAL diff --git a/model_manager/activities/training.py b/model_manager/activities/training.py index 4bc771f..0d6f31c 100644 --- a/model_manager/activities/training.py +++ b/model_manager/activities/training.py @@ -12,7 +12,6 @@ with workflow.unsafe.imports_passed_through(): import traceback from typing import Any - import pandas as pd from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger @@ -21,9 +20,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.repository.minio_repository import MinioRepository from sientia_model.model_repository.mlflow_repository import SientiaMLflowRepository from sientia_model.model_repository.plugin_store import PluginStore - from sientia_model.wrappers.sientia_model import SientiaModel - from model_manager.metrics import ACTIVITY_EXECUTION_TOTAL, WORKFLOW_EXECUTION_TOTAL from model_manager.utils.models.train_model_params import TrainModelParams from model_manager.utils.repository.data_manager_repository import DataManagerRepository @@ -258,7 +255,6 @@ class Training(SientiaMonitoring): 'run_id': run_info.run_id, } except Exception as e: # noqa: BLE001 - metrics_status = 'error' error_msg = f'Error training model - error: {str(e)}' @@ -273,12 +269,9 @@ class Training(SientiaMonitoring): attachment_content=trace, ) - raise ModelTrainingError( - model_trained=model_trained, - model_saved=model_saved, - ) from e + raise e - @activity.defn(name='cleanup_resources') + activity.defn(name='cleanup_resources') async def cleanup_resources(self, input_data: dict[str, Any]) -> None: """ Cleanup temporary resources created during training. @@ -286,38 +279,22 @@ class Training(SientiaMonitoring): Args: input_data: Cleanup configuration containing: - metadata (dict): Workflow execution metadata. - - bucket_name (str): MinIO bucket of the uploaded file. - - file_name (str): MinIO object key to delete. + - run_dir (str): Temporary directory to remove. Raises: Exception: If cleanup fails (after sending notification). """ metadata = input_data.get('metadata', {}) - bucket_name = input_data.get('bucket_name', '') - file_name = input_data.get('file_name', '') - val_file_name = input_data.get('val_file_name') + run_dir = input_data.get('run_dir', '') try: - await self.minio_repository.delete_file( - object_name=file_name, - bucket=bucket_name, - metadata=metadata, - ) - if val_file_name: - await self.minio_repository.delete_file( - object_name=val_file_name, - bucket=bucket_name, - metadata=metadata, - ) + self.model_repository.cleanup_run_directory(run_dir) except Exception as e: # noqa: BLE001 - error_msg = ( - 'Error cleaning up resources - ' - f'File: {bucket_name}/{file_name}, Error: {str(e)}' - ) + error_msg = f'Error cleaning up resources - Run directory: {run_dir}, Error: {str(e)}' trace = traceback.format_exc() - await self.send_notification_async( + self.send_notification( metadata=metadata, notification_id='CLEANUP_RESOURCES_ERROR', message=error_msg, @@ -327,3 +304,4 @@ class Training(SientiaMonitoring): ) raise + diff --git a/model_manager/workflows/train_model.py b/model_manager/workflows/train_model.py index 9ff81d8..2960fc6 100644 --- a/model_manager/workflows/train_model.py +++ b/model_manager/workflows/train_model.py @@ -124,10 +124,7 @@ class TrainModel: finally: try: await self._cleanup_resources( - experiment_run_id=experiment_run_id, - bucket_name=train_params.bucket_name, - file_name=train_params.file_name, - val_file_name=train_params.val_file_name, + run_dir=train_result.get('run_dir'), metadata=metadata, ) except Exception: @@ -278,22 +275,13 @@ class TrainModel: return train_result except Exception as e: - # Mapear flags -> status - # False/False: erro no treino - # True/False: erro ao salvar (MLflow) - # False/True: estado inconsistente, tratar como erro de treino - # True/True: não deveria cair aqui; tratar como erro genérico de treino - status = ExperimentStatus.TRAINING_ERROR - - if isinstance(e, ModelTrainingError) and (e.model_trained and not e.model_saved): - status = ExperimentStatus.TRACKING_SEND_ERROR try: await self._update_experiment_run( metadata=metadata, experiment_run_id=experiment_run_id, update_type=UpdateType.STATUS_WITH_ERROR, - status=status, + status=ExperimentStatus.TRAINING_ERROR, error_message=self._extract_error_message(e), ) except Exception: @@ -302,54 +290,27 @@ class TrainModel: async def _cleanup_resources( self, - experiment_run_id: int, - bucket_name: str, - file_name: str, - val_file_name: str | None, + run_dir: str, metadata: dict[str, Any], ) -> None: """ Cleanup resources. - This method deletes the training file from MinIO. On success, updates - DB status to FILE_DELETED. - On error, updates DB status to FILE_DELETE_ERROR. + This method removes the temporary run directory via activity. Args: - experiment_run_id: Validated experiment run ID + run_dir: Temporary directory to remove metadata: Workflow execution metadata """ - try: - await workflow.execute_activity_method( - Activities.cleanup_resources, - { - **metadata, - 'bucket_name': bucket_name, - 'file_name': file_name, - 'val_file_name': val_file_name, - }, - retry_policy=network_retry_policy, - start_to_close_timeout=timedelta(seconds=TIMEOUT_DELETE_FILE), - ) - - await self._update_experiment_run( - metadata=metadata, - experiment_run_id=experiment_run_id, - update_type=UpdateType.STATUS, - status=ExperimentStatus.FILE_DELETED, - ) - except Exception as e: - try: - await self._update_experiment_run( - metadata=metadata, - experiment_run_id=experiment_run_id, - update_type=UpdateType.STATUS_WITH_ERROR, - status=ExperimentStatus.FILE_DELETE_ERROR, - error_message=self._extract_error_message(e), - ) - except Exception: - pass - raise + await workflow.execute_activity_method( + Activities.cleanup_resources, + { + **metadata, + 'run_dir': run_dir, + }, + retry_policy=network_retry_policy, + start_to_close_timeout=timedelta(seconds=TIMEOUT_DELETE_FILE), + ) async def _update_experiment_run( self,