Files
CloudRecoveryAS/app/engine/restore_worker.py

510 lines
20 KiB
Python

"""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. El backup puede venir con terminaciones extra
# (p.ej. ".KNOWNWORLD") tanto en el ZIP como en el propio .bak; se
# normaliza a "<nodo>.bak". El nodo es siempre el primer segmento del
# nombre del archivo de origen.
locate_step_id = JobStepRepository.create(self.job_id, StepType.LOCATE_BAK)
node_base = Path(job.source_name).name.split(".")[0].strip() or None
self._bak_path = SevenZipExtractor.find_bak_file(
self._extract_dir, node_name=node_base
)
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))