"""Motor principal de la aplicación.""" import platform import socket from threading import Thread from typing import Optional from pathlib import Path from PySide6.QtCore import QObject, Signal, QThreadPool from .file_watcher import FileWatcher, FileStabilityChecker, calculate_file_hash from .maintenance_scheduler import DailyMaintenanceScheduler from .restore_worker import RestoreWorker from .retention import RetentionCleaner from .. import __version__ from ..db.job_repository import JobRepository from ..db.event_repository import EventRepository from ..db.config_repository import ConfigRepository from ..config.env_loader import apply_env_overrides, load_env_file from ..constants import APP_ARCH, APP_DIR, APP_PLATFORM, JobStatus, DEFAULT_CONFIG from ..panel import panel_client from ..utils.logger import app_logger class EngineSignals(QObject): """Señales del motor.""" job_created = Signal(str) # job_id stats_updated = Signal(dict) # stats config_loaded = Signal(dict) # config class RestoreEngine(QObject): """Motor principal de restauración.""" def __init__(self): """Inicializa el motor.""" super().__init__() self.signals = EngineSignals() # Thread pool para workers self._thread_pool = QThreadPool.globalInstance() # File watcher self._file_watcher: Optional[FileWatcher] = None # Mantenimiento diario (retención de respaldos obsoletos) self._maintenance: Optional[DailyMaintenanceScheduler] = None # Estado self._running = False self._paused = False # Configuración self._config = self._load_config() self._report_instance_config_to_panel() app_logger.info("RestoreEngine inicializado") def _load_config(self) -> dict: """Carga la configuración desde la base de datos y variables de entorno.""" load_env_file() config = ConfigRepository.get("app_config", DEFAULT_CONFIG.copy()) for key, value in DEFAULT_CONFIG.items(): if key not in config: config[key] = value config = apply_env_overrides(config) self.signals.config_loaded.emit(config) return config def save_config(self, config: dict): """ Guarda la configuración. Args: config: Configuración a guardar """ ConfigRepository.set("app_config", config) self._config = config app_logger.info("Configuración guardada") self._report_instance_config_to_panel() # Reconfigurar file watcher y mantenimiento si está corriendo if self._running: self._restart_file_watcher() self._restart_maintenance() def get_config(self) -> dict: """Obtiene la configuración actual.""" return self._config.copy() def start(self): """Inicia el motor.""" if self._running: app_logger.warning("Motor ya está corriendo") return # Validar configuración if not self._validate_config(): app_logger.error("Configuración inválida, no se puede iniciar el motor") return # Barrido de Temp: al arrancar no hay jobs corriendo, así que cualquier subcarpeta en # extract_folder es un remanente huérfano de una corrida previa (los caminos de # fallo/diferido no siempre limpiaban). Se elimina para que Temp no se acumule. self._purge_stale_temp() # Configurar thread pool extract_workers = self._config["concurrency"]["extract_workers"] restore_workers = self._config["concurrency"]["restore_workers"] max_threads = extract_workers + restore_workers self._thread_pool.setMaxThreadCount(max_threads) app_logger.info( f"Thread pool configurado: {max_threads} threads " f"(extract: {extract_workers}, restore: {restore_workers})" ) # Iniciar file watcher if self._config["features"]["auto_scan_enabled"]: self._start_file_watcher() # Iniciar el mantenimiento diario (limpieza de respaldos obsoletos) self._start_maintenance() self._running = True self._paused = False self._report_instance_config_to_panel() EventRepository.create("INFO", "Motor iniciado") app_logger.info("Motor iniciado") def _purge_stale_temp(self): """Elimina subcarpetas huérfanas en la carpeta temporal de extracción (al arrancar).""" import shutil extract_folder = (self._config.get("paths") or {}).get("extract_folder") if not extract_folder: return base = Path(extract_folder) if not base.is_dir(): return removed = 0 for child in base.iterdir(): if not child.is_dir(): continue try: shutil.rmtree(child, ignore_errors=True) removed += 1 except Exception as e: app_logger.warning(f"No se pudo limpiar Temp huérfano {child}: {e}") if removed: app_logger.info(f"Temp: {removed} carpeta(s) huérfana(s) eliminada(s) al iniciar") def stop(self): """Detiene el motor.""" if not self._running: return # Detener file watcher if self._file_watcher: self._file_watcher.stop() self._file_watcher = None # Detener el mantenimiento diario self._stop_maintenance() # Esperar a que terminen los workers self._thread_pool.waitForDone(msecs=30000) # 30s timeout self._running = False EventRepository.create("INFO", "Motor detenido") app_logger.info("Motor detenido") def pause(self): """Pausa el procesamiento (no toma nuevos jobs).""" self._paused = True if self._file_watcher: self._file_watcher.stop() EventRepository.create("INFO", "Motor pausado") app_logger.info("Motor pausado") def resume(self): """Reanuda el procesamiento.""" if not self._running: return self._paused = False if self._config["features"]["auto_scan_enabled"]: self._start_file_watcher() EventRepository.create("INFO", "Motor reanudado") app_logger.info("Motor reanudado") def is_running(self) -> bool: """Verifica si el motor está corriendo.""" return self._running def is_paused(self) -> bool: """Verifica si el motor está pausado.""" return self._paused def scan_now(self): """Fuerza un escaneo manual de la carpeta de entrada.""" if self._file_watcher: self._file_watcher.scan_now() else: # Crear watcher temporal para un escaneo self._start_file_watcher() if self._file_watcher: self._file_watcher.scan_now() def process_file(self, file_path: str): """ Procesa un archivo manualmente. Args: file_path: Ruta al archivo ZIP """ try: self._on_file_ready(file_path) except Exception as e: app_logger.error(f"Error procesando archivo {file_path}: {e}") def get_stats(self) -> dict: """Obtiene estadísticas del motor.""" stats = JobRepository.get_stats() stats["active_threads"] = self._thread_pool.activeThreadCount() stats["max_threads"] = self._thread_pool.maxThreadCount() stats["running"] = self._running stats["paused"] = self._paused # Tiempos promedio avg_times = JobRepository.get_average_times() stats.update(avg_times) return stats def _report_instance_config_to_panel(self) -> None: """Reporta input_folder al PANEL (best-effort, no bloquea el flujo).""" panel_cfg = self._config.get("panel", {}) api_url = (panel_cfg.get("api_url") or "").strip() api_token = (panel_cfg.get("api_token") or "").strip() if not api_url or not api_token: return paths = self._config.get("paths") or {} input_folder = paths.get("input_folder") or "" processed_folder = paths.get("processed_folder") or "" try: host_name = socket.gethostname() or platform.node() except Exception: host_name = platform.node() instance_key = (panel_cfg.get("instance_key") or "").strip() or None # platform/arch le dicen al PANEL qué artefacto le toca a este servidor cuando # instala o actualiza (a24c.cras_releases se llavea por version+platform+arch). panel_client.report_instance_config( api_url=api_url, api_token=api_token, input_folder=input_folder, processed_folder=processed_folder, host_name=host_name, app_version=__version__, instance_key=instance_key, platform_name=APP_PLATFORM, arch=APP_ARCH, # Dónde vive el ejecutable. El PANEL actualiza en esta ruta; si instalara en la # default crearía una segunda instalación y dejaría huérfano este config/.env. install_path=str(APP_DIR), ) def _validate_config(self) -> bool: """Valida que la configuración sea correcta.""" paths = self._config["paths"] # Validar carpetas requeridas. processed_folder/failed_folder se incluyen para que la # reubicación de ZIP (éxito/fallo) nunca falle por carpeta inexistente. required_paths = [ "input_folder", "extract_folder", "data_sql_folder", "processed_folder", "failed_folder", ] for key in required_paths: if not paths.get(key): app_logger.error(f"Falta configurar: {key}") return False path = Path(paths[key]) if not path.exists(): try: path.mkdir(parents=True, exist_ok=True) except Exception as e: app_logger.error(f"No se pudo crear carpeta {key}: {e}") return False # Validar 7-Zip if not paths.get("seven_zip_exe"): app_logger.error("Falta configurar ruta a 7-Zip") return False if not Path(paths["seven_zip_exe"]).exists(): app_logger.error(f"7-Zip no encontrado en: {paths['seven_zip_exe']}") return False return True def _start_file_watcher(self): """Inicia el file watcher.""" if self._file_watcher: self._file_watcher.stop() # Configurar stability checker stability_config = self._config["stability"] stability_checker = FileStabilityChecker( check_interval=stability_config["check_interval_seconds"], stable_duration=stability_config["stable_duration_seconds"], use_ready_marker=stability_config["use_ready_marker"] ) # Crear file watcher self._file_watcher = FileWatcher( watch_folder=self._config["paths"]["input_folder"], on_file_ready=self._on_file_ready, scan_interval=self._config["features"]["scan_interval_seconds"], stability_checker=stability_checker ) self._file_watcher.start() def _restart_file_watcher(self): """Reinicia el file watcher con nueva configuración.""" if self._running and not self._paused: self._start_file_watcher() def _start_maintenance(self): """Inicia el mantenimiento diario si la retención está habilitada.""" retention_cfg = self._config.get("retention", {}) if not retention_cfg.get("enabled", True): app_logger.info("Retención deshabilitada; no se inicia el mantenimiento diario") return self._maintenance = DailyMaintenanceScheduler( task=self._run_retention, check_interval_seconds=retention_cfg.get("check_interval_seconds", 3600), run_at_hour=retention_cfg.get("run_at_hour", 3), ) self._maintenance.start() def _stop_maintenance(self): """Detiene el mantenimiento diario.""" if self._maintenance: self._maintenance.stop() self._maintenance = None def _restart_maintenance(self): """Reinicia el mantenimiento diario con nueva configuración.""" if self._running: self._stop_maintenance() self._start_maintenance() def _run_retention(self): """Ejecuta una corrida de retención (tarea del mantenimiento diario).""" try: RetentionCleaner(self._config).run() except Exception as e: app_logger.error(f"Error en la retención: {e}", exc_info=True) EventRepository.create("ERROR", f"Error en la retención: {e}") def trigger_maintenance_now(self): """Corre la retención de inmediato en segundo plano (acción manual de la UI). Respeta el dry_run de la configuración y no consume el turno del día calendario. """ def _run(): if self._maintenance: self._maintenance.trigger_now() else: self._run_retention() Thread(target=_run, daemon=True).start() def _on_file_ready(self, file_path: str): """ Callback cuando un archivo está listo para procesar. Args: file_path: Ruta al archivo """ if self._paused: app_logger.info(f"Motor pausado, ignorando archivo: {file_path}") return try: # Calcular hash para evitar duplicados file_hash = calculate_file_hash(file_path) if JobRepository.has_blocking_job_by_hash(file_hash): app_logger.warning(f"Archivo ya procesado o en curso (hash duplicado): {file_path}") return # Limpia intentos FALLIDOS previos con este hash para permitir un reintento fresco # (antes un FAILED transitorio bloqueaba el reproceso de forma permanente). removed = JobRepository.delete_failed_by_hash(file_hash) if removed: app_logger.info( f"Reintento de {Path(file_path).name}: {removed} job(s) fallido(s) previo(s) eliminado(s)" ) # Crear job file_name = Path(file_path).name job_id = JobRepository.create(file_path, file_name, file_hash) app_logger.info(f"Job creado: {job_id} para archivo {file_name}") EventRepository.create("INFO", f"Nuevo job creado: {file_name}", job_id) self.signals.job_created.emit(job_id) # Crear worker y encolar worker = RestoreWorker( job_id=job_id, config=self._config, dry_run=self._config["features"]["dry_run_mode"] ) # Conectar señales worker.signals.job_completed.connect(self._on_job_completed) # Ejecutar en thread pool self._thread_pool.start(worker) app_logger.info(f"Worker encolado para job {job_id}") except Exception as e: app_logger.error(f"Error creando job para {file_path}: {e}", exc_info=True) EventRepository.create("ERROR", f"Error creando job: {e}") def _on_job_completed(self, job_id: str, success: bool): """ Callback cuando un job se completa. Args: job_id: ID del job success: Si fue exitoso """ app_logger.info(f"Job {job_id} completado - Éxito: {success}") # Actualizar estadísticas stats = self.get_stats() self.signals.stats_updated.emit(stats)