"""Worker para procesar jobs de restauración.""" import time import shutil from pathlib import Path from typing import Optional from datetime import datetime from PySide6.QtCore import QObject, Signal, QRunnable from ..constants import JobStatus, StepType from ..db.job_repository import JobRepository from ..db.job_step_repository import JobStepRepository from ..db.event_repository import EventRepository from ..db.node_repository import NodeRepository from ..extract.seven_zip import SevenZipExtractor from ..sql.sql_manager import SQLServerManager from ..panel import panel_client from ..transfer import sftp_copy from ..utils.logger import app_logger class RestoreDeferred(Exception): """ Señala que el job no puede ejecutarse ahora pero NO es un error: no hay servidor de restauración activo en el PANEL. El job se difiere (se borra para que el próximo escaneo lo reintente), no se marca como fallido (G8 del plan). """ class RestoreWorkerSignals(QObject): """Señales para comunicación con la UI.""" job_started = Signal(str) # job_id job_progress = Signal(str, str) # job_id, status job_completed = Signal(str, bool) # job_id, success error_occurred = Signal(str, str) # job_id, error class RestoreWorker(QRunnable): """Worker que procesa un job de restauración.""" def __init__( self, job_id: str, config: dict, dry_run: bool = False ): """ Inicializa el worker. Args: job_id: ID del job a procesar config: Configuración de la aplicación dry_run: Modo dry run (no ejecuta RESTORE) """ super().__init__() self.job_id = job_id self.config = config self.dry_run = dry_run self.signals = RestoreWorkerSignals() self._extract_dir: Optional[str] = None self._bak_path: Optional[str] = None # Servidor de restauración activo del PANEL (None = modo local sin PANEL). self._target: Optional[dict] = None def _panel_configured(self) -> bool: panel_cfg = self.config.get("panel", {}) return bool((panel_cfg.get("api_url") or "").strip()) def _legacy_panel_mode(self) -> Optional[str]: """Detecta mode legacy en config; solo para log de deprecación.""" mode = (self.config.get("panel", {}).get("mode") or "").strip() return mode if mode in ("orchestrator", "hub_restore", "colocated") else None def run(self): """Ejecuta el procesamiento del job.""" start_time = time.time() try: app_logger.info(f"Iniciando procesamiento de job {self.job_id}") self.signals.job_started.emit(self.job_id) # Obtener job job = JobRepository.get(self.job_id) if not job: raise ValueError(f"Job {self.job_id} no encontrado") if not self._panel_configured(): self._process_node_mapping(job) job = JobRepository.get(self.job_id) self._target = None app_logger.info("PANEL no configurado: usando configuración SQL local") else: legacy_mode = self._legacy_panel_mode() if legacy_mode and legacy_mode != "colocated": app_logger.warning( f"panel.mode={legacy_mode} está obsoleto; " "se usa enrutamiento automático (resolve-route)" ) route = self._resolve_route(job) job = JobRepository.get(self.job_id) if route["action"] == "forward": self._forward_zip(job, route, start_time) return self._target = route["target"] # Extraer y restaurar localmente. self._extract_backup(job) self._restore_database(job) self._cleanup(job) # Actualizar tiempos total_ms = int((time.time() - start_time) * 1000) JobRepository.update_timing(self.job_id, total_ms=total_ms) # Marcar como completado JobRepository.update_status(self.job_id, JobStatus.COMPLETED) app_logger.info( f"Job {self.job_id} completado exitosamente en {total_ms}ms" ) self._report_to_panel(job, "completed", duration_ms=total_ms) self.signals.job_completed.emit(self.job_id, True) except RestoreDeferred as e: # No hay servidor activo: borrar el job para que el próximo escaneo # vuelva a detectar el archivo y reintente. No cuenta como fallo. app_logger.warning(f"Job {self.job_id} diferido: {e}") EventRepository.create( "WARNING", f"Restauración diferida (sin servidor activo): {e}", self.job_id ) JobRepository.delete(self.job_id) self.signals.job_completed.emit(self.job_id, False) except Exception as e: error_msg = str(e) app_logger.error(f"Error procesando job {self.job_id}: {error_msg}", exc_info=True) JobRepository.update_status( self.job_id, JobStatus.FAILED, error=error_msg, increment_attempts=True ) EventRepository.create("ERROR", f"Job falló: {error_msg}", self.job_id) job = JobRepository.get(self.job_id) self._report_to_panel(job, "failed", error_message=error_msg) self.signals.error_occurred.emit(self.job_id, error_msg) self.signals.job_completed.emit(self.job_id, False) def _resolve_route(self, job) -> dict: """ Enrutamiento automático vía panel: restore_local o forward según nodo/asignación. Actualiza db_name del job desde el panel. """ panel_cfg = self.config.get("panel", {}) api_url = (panel_cfg.get("api_url") or "").strip() api_token = (panel_cfg.get("api_token") or "").strip() instance_key = (panel_cfg.get("instance_key") or "").strip() if not instance_key: raise RestoreDeferred( "Falta instance_key (instancia/servidor) en la configuración del PANEL" ) route = panel_client.resolve_route( api_url, api_token, job.source_name, instance_key=instance_key, ) if not route: raise RestoreDeferred( f"El PANEL no pudo resolver la ruta para '{job.source_name}'" ) node_key = (route.get("node_key") or Path(job.source_name).name).upper() JobRepository.update_node_and_db(self.job_id, node_key, route["db_name"]) app_logger.info( f"Ruta resuelta: action={route.get('action')}, db={route.get('db_name')}" ) return route def _collect_zip_paths(self, source_path: str) -> list[str]: """Rutas locales del ZIP (incluye todas las partes multipart).""" path = Path(source_path) if SevenZipExtractor.is_multipart(str(path)): base_name = path.stem.split(".zip")[0] parts = sorted(path.parent.glob(f"{base_name}.zip.*")) return [str(p) for p in parts] return [str(path)] def _move_zip_to_processed(self, job): """Mueve el ZIP (y partes multipart) a la carpeta processed.""" processed_folder = Path(self.config["paths"]["processed_folder"]) date_folder = processed_folder / datetime.now().strftime("%Y-%m-%d") date_folder.mkdir(parents=True, exist_ok=True) for zip_path in self._collect_zip_paths(job.source_path): src = Path(zip_path) dest = date_folder / src.name shutil.move(str(src), str(dest)) app_logger.info(f"Movido: {src.name} -> {dest}") def _forward_zip(self, job, route: dict, start_time: float): """Reenvía el ZIP al input_folder del servidor destino vía SFTP.""" target = route["target"] self._target = target step_id = JobStepRepository.create(self.job_id, StepType.FORWARD_ZIP) try: zip_paths = self._collect_zip_paths(job.source_path) remote_folder = target["input_folder"] uploaded = sftp_copy.upload_zip_parts(zip_paths, target, remote_folder) JobStepRepository.complete( step_id, exit_code=0, stdout=f"Destino: {target.get('name')} ({len(uploaded)} archivo(s))", ) except Exception as e: JobStepRepository.complete(step_id, exit_code=1, error=str(e)) raise self._move_zip_to_processed(job) total_ms = int((time.time() - start_time) * 1000) JobRepository.update_timing(self.job_id, total_ms=total_ms) JobRepository.update_status(self.job_id, JobStatus.COMPLETED) app_logger.info( f"Job {self.job_id} reenviado a '{target.get('name')}' en {total_ms}ms" ) EventRepository.create( "INFO", f"ZIP reenviado por SFTP a {target.get('name')}", self.job_id, ) self._report_to_panel(job, "forwarded", duration_ms=total_ms) self.signals.job_completed.emit(self.job_id, True) def _report_to_panel( self, job, status: str, duration_ms: Optional[int] = None, error_message: Optional[str] = None, ): """Reporta el resultado del job al PANEL (best-effort, solo en modo PANEL).""" if not self._target: return panel_cfg = self.config.get("panel", {}) panel_client.report_job_result( api_url=(panel_cfg.get("api_url") or "").strip(), api_token=(panel_cfg.get("api_token") or "").strip(), filename=job.source_name if job else "", status=status, restore_target_id=self._target.get("id"), db_name=job.db_name if job else None, duration_ms=duration_ms, error_message=error_message, ) def _process_node_mapping(self, job): """Procesa el mapeo de nodo a base de datos.""" step_id = JobStepRepository.create(self.job_id, StepType.NODE_MAPPING) try: # Obtener node name del archivo (con extensión, uppercase) node_name = Path(job.source_name).name.upper() # Buscar mapeo en la tabla nodes db_name = NodeRepository.get_db_for_node(node_name) if not db_name: raise ValueError( f"NODE_NOT_MAPPED: No existe mapeo activo para nodo '{node_name}'" ) # Actualizar job con node_name y db_name JobRepository.update_node_and_db(self.job_id, node_name, db_name) app_logger.info(f"Node '{node_name}' mapeado a DB '{db_name}'") JobStepRepository.complete(step_id, exit_code=0, stdout=f"DB: {db_name}") except Exception as e: JobStepRepository.complete(step_id, exit_code=1, error=str(e)) raise def _extract_backup(self, job): """Extrae el backup del archivo ZIP.""" JobRepository.update_status(self.job_id, JobStatus.EXTRACTING) self.signals.job_progress.emit(self.job_id, JobStatus.EXTRACTING) step_id = JobStepRepository.create(self.job_id, StepType.EXTRACT) extract_start = time.time() try: # Configuración seven_zip_path = self.config["paths"]["seven_zip_exe"] extract_base = self.config["paths"]["extract_folder"] timeout_minutes = self.config["timeouts"]["extract_minutes"] # Crear carpeta de extracción para este job self._extract_dir = str(Path(extract_base) / self.job_id) Path(self._extract_dir).mkdir(parents=True, exist_ok=True) # Extraer extractor = SevenZipExtractor(seven_zip_path) # Si es multipart, asegurar que usamos el primer archivo source_path = job.source_path if extractor.is_multipart(source_path): source_path = extractor.get_first_part(source_path) app_logger.info(f"Archivo multipart detectado, usando: {source_path}") exit_code, stdout, stderr = extractor.extract( source_path, self._extract_dir, timeout_minutes ) extract_ms = int((time.time() - extract_start) * 1000) JobRepository.update_timing(self.job_id, extract_ms=extract_ms) if exit_code != 0: raise RuntimeError(f"7-Zip falló con código {exit_code}: {stderr}") JobStepRepository.complete( step_id, exit_code=exit_code, stdout=stdout[:1000] if stdout else None, stderr=stderr[:1000] if stderr else None ) # Buscar archivo .bak locate_step_id = JobStepRepository.create(self.job_id, StepType.LOCATE_BAK) self._bak_path = SevenZipExtractor.find_bak_file(self._extract_dir) if not self._bak_path: raise FileNotFoundError("No se encontró archivo .bak en el archivo extraído") JobStepRepository.complete( locate_step_id, exit_code=0, stdout=f"BAK: {Path(self._bak_path).name}" ) except Exception as e: JobStepRepository.complete(step_id, exit_code=1, error=str(e)) raise def _restore_database(self, job): """Restaura la base de datos desde el backup.""" JobRepository.update_status(self.job_id, JobStatus.RESTORING) self.signals.job_progress.emit(self.job_id, JobStatus.RESTORING) # Recargar job para obtener db_name actualizado job = JobRepository.get(self.job_id) if not job.db_name: raise ValueError("DB name no está configurado en el job") # Servidor y credenciales: del servidor activo del PANEL, o config local. if self._target: server = self._target["server"] use_windows_auth = False username = self._target["username"] password = self._target["password"] data_folder = self._target["data_folder"] else: sql_config = self.config["sql"] server = sql_config["server"] use_windows_auth = sql_config["use_windows_auth"] username = sql_config.get("username") password = sql_config.get("password") data_folder = self.config["paths"]["data_sql_folder"] # Conectar a SQL connect_step_id = JobStepRepository.create(self.job_id, StepType.SQL_CONNECT) try: sql_manager = SQLServerManager( server=server, use_windows_auth=use_windows_auth, username=username, password=password ) # Test connection success, error = sql_manager.test_connection() if not success: raise RuntimeError(f"Conexión SQL falló: {error}") JobStepRepository.complete(connect_step_id, exit_code=0) except Exception as e: JobStepRepository.complete(connect_step_id, exit_code=1, error=str(e)) raise sql_backup_path = self._bak_path # Obtener FILELISTONLY filelist_step_id = JobStepRepository.create(self.job_id, StepType.FILELIST) filelist_start = time.time() try: logical_files, error = sql_manager.get_filelist_from_backup( sql_backup_path, timeout_minutes=self.config["timeouts"]["restore_minutes"] ) filelist_ms = int((time.time() - filelist_start) * 1000) JobRepository.update_timing(self.job_id, filelist_ms=filelist_ms) if error: raise RuntimeError(error) files_info = ", ".join([f"{lf.logical_name}({lf.type})" for lf in logical_files]) JobStepRepository.complete( filelist_step_id, exit_code=0, stdout=files_info[:1000] ) except Exception as e: JobStepRepository.complete(filelist_step_id, exit_code=1, error=str(e)) raise # RESTORE DATABASE restore_step_id = JobStepRepository.create(self.job_id, StepType.RESTORE) restore_start = time.time() try: success, stdout, error = sql_manager.restore_database( db_name=job.db_name, backup_path=sql_backup_path, data_folder=data_folder, logical_files=logical_files, timeout_minutes=self.config["timeouts"]["restore_minutes"], dry_run=self.dry_run ) restore_ms = int((time.time() - restore_start) * 1000) JobRepository.update_timing(self.job_id, restore_ms=restore_ms) if not success: raise RuntimeError(error or "RESTORE falló sin mensaje de error") JobStepRepository.complete( restore_step_id, exit_code=0, stdout=stdout[:1000] if stdout else None ) except Exception as e: JobStepRepository.complete(restore_step_id, exit_code=1, error=str(e)) raise def _cleanup(self, job): """Limpia archivos temporales y mueve el ZIP.""" JobRepository.update_status(self.job_id, JobStatus.CLEANING) self.signals.job_progress.emit(self.job_id, JobStatus.CLEANING) step_id = JobStepRepository.create(self.job_id, StepType.CLEANUP) try: # Eliminar carpeta de extracción if self._extract_dir and Path(self._extract_dir).exists(): shutil.rmtree(self._extract_dir) app_logger.info(f"Carpeta de extracción eliminada: {self._extract_dir}") # Mover ZIP a Processed processed_folder = Path(self.config["paths"]["processed_folder"]) date_folder = processed_folder / datetime.now().strftime("%Y-%m-%d") date_folder.mkdir(parents=True, exist_ok=True) source_path = Path(job.source_path) dest_path = date_folder / source_path.name # Si es multipart, mover todas las partes if SevenZipExtractor.is_multipart(str(source_path)): # Buscar todas las partes base_name = source_path.stem.split('.zip')[0] parts = list(source_path.parent.glob(f"{base_name}.zip.*")) for part in parts: part_dest = date_folder / part.name shutil.move(str(part), str(part_dest)) app_logger.info(f"Movido: {part.name} -> {part_dest}") else: shutil.move(str(source_path), str(dest_path)) app_logger.info(f"Movido: {source_path.name} -> {dest_path}") JobStepRepository.complete(step_id, exit_code=0) except Exception as e: # No fallar el job por errores de limpieza app_logger.warning(f"Error en limpieza (no crítico): {e}") JobStepRepository.complete(step_id, exit_code=1, error=str(e))