504 lines
19 KiB
Python
504 lines
19 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
|
|
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))
|