Files
CloudRecoveryAS/app/engine/engine.py

455 lines
16 KiB
Python

"""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)