From afb76ce4a651806648643246de095051ee244514 Mon Sep 17 00:00:00 2001 From: hreyes Date: Fri, 5 Jun 2026 10:49:05 -0600 Subject: [PATCH] feature/integracion-panel-restore-targets --- INTEGRACION_PANEL.md | 197 ++++++++++++++++ QUICKSTART.md | 4 + app/constants.py | 11 + app/db/job_repository.py | 9 + app/engine/engine.py | 37 ++- app/engine/restore_worker.py | 263 +++++++++++++++++---- app/panel/__init__.py | 1 + app/panel/panel_client.py | 402 ++++++++++++++++++++++++++++++++ app/sql/sql_manager.py | 140 ++++++----- app/transfer/__init__.py | 1 + app/transfer/sftp_copy.py | 160 +++++++++++++ app/ui/config_tab.py | 113 ++++++++- app/ui/nodes_tab.py | 2 +- requirements.txt | 9 + tests/conftest.py | 5 + tests/test_colocated_resolve.py | 41 ++++ tests/test_panel_client.py | 381 ++++++++++++++++++++++++++++++ tests/test_route_forward.py | 138 +++++++++++ tests/test_sftp_copy.py | 135 +++++++++++ tests/test_sql_manager_move.py | 77 ++++++ 20 files changed, 2022 insertions(+), 104 deletions(-) create mode 100644 INTEGRACION_PANEL.md create mode 100644 app/panel/__init__.py create mode 100644 app/panel/panel_client.py create mode 100644 app/transfer/__init__.py create mode 100644 app/transfer/sftp_copy.py create mode 100644 tests/conftest.py create mode 100644 tests/test_colocated_resolve.py create mode 100644 tests/test_panel_client.py create mode 100644 tests/test_route_forward.py create mode 100644 tests/test_sftp_copy.py create mode 100644 tests/test_sql_manager_move.py diff --git a/INTEGRACION_PANEL.md b/INTEGRACION_PANEL.md new file mode 100644 index 0000000..abbff17 --- /dev/null +++ b/INTEGRACION_PANEL.md @@ -0,0 +1,197 @@ +# Integración CloudRestoreAS ↔ PANEL_BASES_ANEXO24 + +## Variables de entorno y configuración + +| Panel (`.env`) | CloudRestoreAS (pestaña Config → PANEL) | Debe coincidir | +|----------------|-------------------------------------------|----------------| +| `CLOUDRESTORE_API_TOKEN` | `panel.api_token` | **Sí** — mismo valor en todas las instalaciones | +| `ENCRYPTION_KEY` | — | Solo panel (cifra credenciales de `restore_targets`) | +| URL del panel (ej. `https://panel:3000`) | `panel.api_url` | **Sí** — base URL sin barra final | +| — | `panel.instance_key` | Nombre de **este** servidor en el panel (`restore_targets.name`) | + +Generar token: + +```bash +node -e "console.log(require('crypto').randomBytes(32).toString('hex'))" +``` + +Generar clave de cifrado (panel): + +```bash +node -e "console.log(require('crypto').randomBytes(32).toString('base64'))" +``` + +`BACKUP_PATH` del panel es independiente: lista respaldos en el dashboard. Cada CloudRestoreAS reporta su carpeta local vía `instance-config` (clave = nombre del servidor). + +--- + +## Enrutamiento automático (todos los CRA) + +No hay modos que configurar. **Todo CRA con panel** hace lo mismo: + +1. Llega un ZIP a la carpeta de entrada. +2. `GET resolve-route?filename=X&instance=` — el panel identifica el nodo y el servidor asignado. +3. Si `action=restore_local` → extract + RESTORE aquí (credenciales SQL de este servidor en el panel). +4. Si `action=forward` → SFTP del ZIP al `input_folder` del destino (credenciales SSH del destino). +5. Mover ZIP a `processed`; reportar `completed` o `forwarded`. + +```mermaid +flowchart TB + ZIP[ZIP en carpeta local] + CRA[Any CRA con instance_key] + Panel["resolve-route"] + SFTP[SFTP al destino] + SQL[RESTORE local] + + ZIP --> CRA --> Panel + Panel -->|restore_local| SQL + Panel -->|forward| SFTP +``` + +Aplica igual en Alfa, Omega o un hub donde llegan todos los ZIP: + +| Situación | Qué hace el CRA | +|-----------|-----------------| +| Nodo asignado a **esta** instancia | RESTORE local | +| Nodo asignado a **otro** servidor | SFTP ZIP al destino | + +La decisión la toma el **panel** (nodo + asignación en Gestión BD), no el operador. + +### Qué configura cada instalación + +| Campo | Para qué | +|-------|----------| +| URL + token del panel | Conectar al panel | +| **Instancia** | Identidad de este servidor (`restore_targets.name`) | +| Carpeta entrada | Dónde vigila ZIPs esta máquina | + +El selector de instancia se llena desde `GET /api/restore/target-catalog`. No hay límite de servidores. + +### Hub donde llegan todos los ZIP + +Ejemplo: máquina **Alfa** recibe todos los archivos y también tiene bases propias: + +1. Panel: servidor Alfa registrado; bases de Alfa/Omega/Gamma asignadas en Gestión BD. +2. CRA en Alfa: instancia **Alfa**, misma URL/token que el resto. +3. ZIP de base Alfa → RESTORE local. +4. ZIP de base Omega → SFTP a carpeta de Omega (card **Reportada**). +5. CRA Omega detecta el ZIP y restaura localmente. + +No hay paso extra ni modo especial. + +### Sin panel (legacy local) + +Si `panel.api_url` está vacío, CloudRestoreAS usa nodos SQLite locales y SQL de la pestaña Config (sin reparto automático). + +### Config legacy `panel.mode` + +Valores antiguos (`orchestrator`, `hub_restore`, `colocated`) se **ignoran**. Se registra un aviso en log y se usa siempre `resolve-route`. El modo hub que subía `.bak` por SFTP ya no aplica. + +--- + +## Contrato API (`/api/restore/*`) + +Autenticación: `Authorization: Bearer `. + +### GET `/api/restore/target-catalog` + +**Cliente:** `panel_client.list_restore_target_names()` + +| Respuesta 200 | Descripción | +|---------------|-------------| +| `targets[]` | `{ "id", "name" }` — sin credenciales | + +### GET `/api/restore/resolve-route?filename=&instance=` + +**Cliente:** `panel_client.resolve_route()` + +| Query | Obligatorio | Descripción | +|-------|-------------|-------------| +| `filename` | Sí | Nombre del ZIP (panel resuelve nodo desde el stem) | +| `instance` | Sí | Instancia de este CRA (`restore_targets.name`) | + +| Respuesta 200 | Descripción | +|---------------|-------------| +| `action` | `restore_local` o `forward` | +| `db_name` | Base resuelta | +| `node_key` | Nodo en panel | +| `target` | Servidor destino (SQL/SSH + `input_folder`) | + +| Código | Significado | +|--------|-------------| +| 404 | Sin nodo/base o sin servidor asignado | +| 503 | Destino sin `input_folder` reportado | + +### GET `/api/restore/target-for?database=&instance=` + +Usado por utilidades legacy; el flujo principal de jobs usa `resolve-route`. + +### POST `/api/restore/job-result` + +`status`: `completed` | `failed` | `forwarded`. + +### POST `/api/restore/instance-config` + +Reporte de carpeta de entrada (`instance_key` = nombre del servidor). + +--- + +## Agregar servidores adicionales + +1. Panel → **Servidores de Restauración → + Nuevo servidor** +2. **Gestión de Bases de Datos** → dropdown servidor por base +3. Instalar CRA → Config → instancia → Guardar (card **Reportada**) + +--- + +## Despliegue inicial + +### Orden de arranque (migraciones automáticas) + +1. **a24c-postgres** — volumen Postgres. +2. **a24c-backend** — aplica `alembic upgrade head` al iniciar (incluye tablas CRA: + `restore_targets`, `restore_job_logs`, `cloudrestore_status`, `restore_target_id`). +3. **Panel** (`docker compose up`) — espera el esquema CRA en Postgres antes de abrir el puerto 3000. + +No hace falta ejecutar `database/migrations/*.sql` a mano: la fuente de verdad es Alembic en **a24c**. + +### Panel + +```bash +cd ~/dev/PANEL_BASES_ANEXO24 +cp .env.example .env +docker compose up -d +``` + +Servidores de restauración + asignación de bases en Gestión BD. + +### Cada Windows con CloudRestoreAS + +| Config | Valor | +|--------|--------| +| URL PANEL | misma en todos | +| API Token | mismo en todos | +| **Instancia** | nombre de **este** servidor en panel | +| Carpeta Entrada | local de esta máquina | + +### Verificaciones + +| Paso | Qué comprobar | +|------|----------------| +| Cards en panel | **Reportada** con carpeta | +| ZIP propio | RESTORE local | +| ZIP ajeno en esta carpeta | SFTP al destino; bitácora `forwarded` | +| Destino recibe ZIP | RESTORE local allí | + +### Troubleshooting + +- **Sin reportar**: CRA destino no guardó config o no alcanza el panel. +- **503 en forward**: destino sin `input_folder` reportado. +- **Job diferido**: panel caído, sin `instance_key`, o base sin asignación. +- **401**: token distinto entre panel y CRA. + +--- + +## Prueba local + +Un CRA con `instance_key` igual al nombre en panel. El panel muestra una card por servidor; las no reportadas aparecen como **Sin reportar**. diff --git a/QUICKSTART.md b/QUICKSTART.md index 317881f..40a3ab1 100644 --- a/QUICKSTART.md +++ b/QUICKSTART.md @@ -20,6 +20,10 @@ Get-OdbcDriver | Where-Object {$_.Name -like "*SQL Server*"} .\venv\Scripts\python.exe runner.py ``` +### Integración con PANEL_BASES_ANEXO24 + +Si usas el panel para asignar servidores de restauración, ver [INTEGRACION_PANEL.md](INTEGRACION_PANEL.md) para tokens, catálogo dinámico (`target-catalog`) y prueba end-to-end. + ### 3. Configuración Básica (en la aplicación) 1. **Tab Configuración**: diff --git a/app/constants.py b/app/constants.py index c69fdc4..7d9cc14 100644 --- a/app/constants.py +++ b/app/constants.py @@ -39,6 +39,8 @@ class StepType: FILELIST = "filelist" RESTORE = "restore" CLEANUP = "cleanup" + SFTP_COPY = "sftp_copy" # Transferencia del .bak al servidor SQL externo (SFTP/SSH) + FORWARD_ZIP = "forward_zip" # Reenvío del ZIP al input_folder del servidor destino # Configuración por defecto DEFAULT_CONFIG = { @@ -78,5 +80,14 @@ DEFAULT_CONFIG = { "dry_run_mode": False, "auto_scan_enabled": True, "scan_interval_seconds": 30 + }, + # Integración con PANEL_BASES_ANEXO24. Si api_url está vacío, CloudRestoreAS + # opera en modo local usando la sección "sql" (retrocompatibilidad). + "panel": { + "api_url": "", + "api_token": "", + # Nombre del restore_target en el panel (identidad de este CRA). + # Con panel configurado: enrutamiento automático restore_local | forward. + "instance_key": "", } } diff --git a/app/db/job_repository.py b/app/db/job_repository.py index ebea029..c69c7a3 100644 --- a/app/db/job_repository.py +++ b/app/db/job_repository.py @@ -176,6 +176,15 @@ class JobRepository: tuple(params) ) + @staticmethod + def delete(job_id: str): + """ + Elimina un job y su rastro de hash. Se usa para diferir un job cuando no + hay servidor de restauración activo: al borrarlo, el próximo escaneo del + FileWatcher vuelve a detectar el archivo y reintenta (G8 del plan). + """ + db.execute("DELETE FROM jobs WHERE job_id = ?", (job_id,)) + @staticmethod def exists_by_hash(source_hash: str) -> bool: """Verifica si existe un job con el hash dado.""" diff --git a/app/engine/engine.py b/app/engine/engine.py index 7a66dfd..d3c7479 100644 --- a/app/engine/engine.py +++ b/app/engine/engine.py @@ -1,15 +1,19 @@ """Motor principal de la aplicación.""" +import platform +import socket 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 .restore_worker import RestoreWorker +from .. import __version__ from ..db.job_repository import JobRepository from ..db.event_repository import EventRepository from ..db.config_repository import ConfigRepository from ..constants import JobStatus, DEFAULT_CONFIG +from ..panel import panel_client from ..utils.logger import app_logger @@ -41,7 +45,8 @@ class RestoreEngine(QObject): # Configuración self._config = self._load_config() - + self._report_instance_config_to_panel() + app_logger.info("RestoreEngine inicializado") def _load_config(self) -> dict: @@ -66,7 +71,8 @@ class RestoreEngine(QObject): ConfigRepository.set("app_config", config) self._config = config app_logger.info("Configuración guardada") - + self._report_instance_config_to_panel() + # Reconfigurar file watcher si está corriendo if self._running: self._restart_file_watcher() @@ -103,7 +109,8 @@ class RestoreEngine(QObject): self._running = True self._paused = False - + self._report_instance_config_to_panel() + EventRepository.create("INFO", "Motor iniciado") app_logger.info("Motor iniciado") @@ -190,6 +197,30 @@ class RestoreEngine(QObject): 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 + + input_folder = (self._config.get("paths") or {}).get("input_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 + panel_client.report_instance_config( + api_url=api_url, + api_token=api_token, + input_folder=input_folder, + host_name=host_name, + app_version=__version__, + instance_key=instance_key, + ) + def _validate_config(self) -> bool: """Valida que la configuración sea correcta.""" paths = self._config["paths"] diff --git a/app/engine/restore_worker.py b/app/engine/restore_worker.py index f687a3a..b705805 100644 --- a/app/engine/restore_worker.py +++ b/app/engine/restore_worker.py @@ -14,9 +14,19 @@ 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 @@ -47,9 +57,20 @@ class RestoreWorker(QRunnable): 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.""" @@ -58,46 +79,190 @@ class RestoreWorker(QRunnable): 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") - - # Pipeline de procesamiento - self._process_node_mapping(job) + + 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) @@ -192,95 +357,109 @@ class RestoreWorker(QRunnable): """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_config = self.config["sql"] sql_manager = SQLServerManager( - server=sql_config["server"], - use_windows_auth=sql_config["use_windows_auth"], - username=sql_config.get("username"), - password=sql_config.get("password") + 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( - self._bak_path, + 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: - data_folder = self.config["paths"]["data_sql_folder"] - success, stdout, error = sql_manager.restore_database( db_name=job.db_name, - backup_path=self._bak_path, + 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) diff --git a/app/panel/__init__.py b/app/panel/__init__.py new file mode 100644 index 0000000..f9c2ab2 --- /dev/null +++ b/app/panel/__init__.py @@ -0,0 +1 @@ +"""Integración con PANEL_BASES_ANEXO24 (servidor de restauración activo).""" diff --git a/app/panel/panel_client.py b/app/panel/panel_client.py new file mode 100644 index 0000000..cdf3db2 --- /dev/null +++ b/app/panel/panel_client.py @@ -0,0 +1,402 @@ +""" +Cliente HTTP para PANEL_BASES_ANEXO24. + +Cada base de datos está asignada en el PANEL a un servidor de restauración +(registrado en restore_targets). CloudRestoreAS, tras resolver el db_name del archivo, pregunta +al PANEL (polling) qué servidor corresponde a ESA base y restaura ahí; al terminar, +reporta el resultado del job. + +Política ante fallo (G8/G9): si el PANEL no responde o la base no tiene servidor +asignado, get_target_for_database devuelve None y el worker deja el job en cola para +reintentar; nunca se restaura con datos inciertos. +""" + +from typing import Optional +from urllib.parse import urlparse, quote + +import requests + +from ..utils.logger import app_logger + +# Timeout (connect, read) en segundos para las llamadas al PANEL. +DEFAULT_TIMEOUT = (5, 10) + +# Campos obligatorios del servidor devuelto por el endpoint (modo hub central + SFTP). +REQUIRED_TARGET_FIELDS = ( + "server", + "username", + "password", + "data_folder", + "ssh_host", + "ssh_username", + "ssh_password", + "remote_inbox_path", +) + +# Modo colocado (CRA en el mismo servidor): solo credenciales SQL locales. +REQUIRED_COLOCATED_TARGET_FIELDS = ( + "name", + "server", + "username", + "password", + "data_folder", +) + +# Reenvío de ZIP por SFTP al input_folder del destino. +REQUIRED_FORWARD_TARGET_FIELDS = ( + "name", + "ssh_host", + "ssh_username", + "ssh_password", + "input_folder", +) + + +def _normalize_base_url(api_url: str) -> Optional[str]: + """Valida y normaliza la URL base del PANEL. Devuelve None si es inválida.""" + if not api_url or not api_url.strip(): + return None + base = api_url.strip().rstrip("/") + parsed = urlparse(base) + # Solo http/https; evita esquemas peligrosos (file://, etc.) + if parsed.scheme not in ("http", "https") or not parsed.netloc: + app_logger.error(f"URL de PANEL inválida (esquema/host): {api_url}") + return None + return base + + +def _validate_target(data: dict, *, colocated: bool = False) -> bool: + """Verifica que el servidor devuelto tenga todos los campos requeridos no vacíos.""" + fields = REQUIRED_COLOCATED_TARGET_FIELDS if colocated else REQUIRED_TARGET_FIELDS + for field in fields: + value = data.get(field) + if value is None or (isinstance(value, str) and not value.strip()): + app_logger.error(f"Servidor del PANEL sin campo requerido: '{field}'") + return False + return True + + +def get_target_for_database( + api_url: str, + api_token: str, + db_name: str, + instance_key: Optional[str] = None, +) -> Optional[dict]: + """ + Consulta GET /api/restore/target-for?database=: el servidor de + restauración asignado a esa base de datos. + + Returns: + dict con id, name, server, username, password, data_folder, ssh_host, + ssh_port, ssh_username, ssh_password, remote_inbox_path si la base tiene + servidor asignado; None si no tiene asignación, el PANEL no responde o la + config es inválida. + """ + base = _normalize_base_url(api_url) + if base is None: + return None + if not api_token or not api_token.strip(): + app_logger.error("Token de PANEL no configurado; no se puede consultar el servidor de la base") + return None + if not db_name or not db_name.strip(): + app_logger.error("db_name vacío; no se puede consultar el servidor de la base") + return None + + colocated = bool(instance_key and instance_key.strip()) + url = f"{base}/api/restore/target-for?database={quote(db_name.strip())}" + if colocated: + url += f"&instance={quote(instance_key.strip())}" + headers = {"Authorization": f"Bearer {api_token.strip()}"} + + try: + resp = requests.get(url, headers=headers, timeout=DEFAULT_TIMEOUT) + except requests.exceptions.RequestException as e: + app_logger.error(f"PANEL no accesible al consultar servidor de la base '{db_name}': {e}") + return None + + if resp.status_code == 404: + app_logger.warning(f"La base '{db_name}' no tiene servidor de restauración asignado en el PANEL") + return None + if resp.status_code == 401: + app_logger.error("PANEL rechazó el token de CloudRestoreAS (401)") + return None + if resp.status_code != 200: + app_logger.error(f"PANEL respondió {resp.status_code} al consultar servidor de la base") + return None + + try: + data = resp.json() + except ValueError: + app_logger.error("Respuesta del PANEL no es JSON válido") + return None + + if not _validate_target(data, colocated=colocated): + return None + + name = data.get("name") or data.get("server") + app_logger.info(f"Servidor para la base '{db_name}': {name} ({data.get('server')})") + return data + + +def list_restore_target_names(api_url: str, api_token: str) -> list[str]: + """ + Consulta GET /api/restore/target-catalog: nombres de servidores de restauración + registrados en el panel (para llenar el selector de instancia). + + Returns: + Lista de nombres; lista vacía si falla la conexión, el token es inválido o la + respuesta no es válida. + """ + base = _normalize_base_url(api_url) + if base is None: + return [] + if not api_token or not api_token.strip(): + app_logger.error("Token de PANEL no configurado; no se puede consultar el catálogo") + return [] + + url = f"{base}/api/restore/target-catalog" + headers = {"Authorization": f"Bearer {api_token.strip()}"} + + try: + resp = requests.get(url, headers=headers, timeout=DEFAULT_TIMEOUT) + except requests.exceptions.RequestException as e: + app_logger.error(f"PANEL no accesible al consultar catálogo de servidores: {e}") + return [] + + if resp.status_code == 401: + app_logger.error("PANEL rechazó el token de CloudRestoreAS (401) al consultar catálogo") + return [] + if resp.status_code != 200: + app_logger.error(f"PANEL respondió {resp.status_code} al consultar catálogo de servidores") + return [] + + try: + data = resp.json() + except ValueError: + app_logger.error("Respuesta del PANEL (target-catalog) no es JSON válido") + return [] + + targets = data.get("targets") + if not isinstance(targets, list): + app_logger.error("Respuesta del PANEL (target-catalog) sin lista 'targets'") + return [] + + names: list[str] = [] + for item in targets: + if isinstance(item, dict): + name = item.get("name") + if isinstance(name, str) and name.strip(): + names.append(name.strip()) + return names + + +def resolve_route( + api_url: str, + api_token: str, + filename: str, + instance_key: Optional[str] = None, +) -> Optional[dict]: + """ + Consulta GET /api/restore/resolve-route: resuelve nodo/destino desde el nombre + del archivo y devuelve action (restore_local | forward) más target completo. + + Returns: + dict con action, db_name, node_key, target; None si el panel no responde, + token inválido o la respuesta no cumple el contrato. + """ + base = _normalize_base_url(api_url) + if base is None: + return None + if not api_token or not api_token.strip(): + app_logger.error("Token de PANEL no configurado; no se puede resolver ruta") + return None + if not filename or not filename.strip(): + app_logger.error("filename vacío; no se puede resolver ruta") + return None + + url = f"{base}/api/restore/resolve-route?filename={quote(filename.strip())}" + key = (instance_key or "").strip() + if key: + url += f"&instance={quote(key)}" + headers = {"Authorization": f"Bearer {api_token.strip()}"} + + try: + resp = requests.get(url, headers=headers, timeout=DEFAULT_TIMEOUT) + except requests.exceptions.RequestException as e: + app_logger.error(f"PANEL no accesible al resolver ruta para '{filename}': {e}") + return None + + if resp.status_code == 404: + app_logger.warning(f"El archivo '{filename}' no tiene nodo/base asignado en el PANEL") + return None + if resp.status_code == 503: + app_logger.warning( + f"Destino aún sin carpeta de entrada reportada para '{filename}'" + ) + return None + if resp.status_code == 401: + app_logger.error("PANEL rechazó el token de CloudRestoreAS (401) al resolver ruta") + return None + if resp.status_code != 200: + app_logger.error(f"PANEL respondió {resp.status_code} al resolver ruta") + return None + + try: + data = resp.json() + except ValueError: + app_logger.error("Respuesta del PANEL (resolve-route) no es JSON válido") + return None + + action = data.get("action") + target = data.get("target") + db_name = data.get("db_name") + if action not in ("restore_local", "forward") or not isinstance(target, dict): + app_logger.error("Respuesta del PANEL (resolve-route) con action o target inválidos") + return None + if not isinstance(db_name, str) or not db_name.strip(): + app_logger.error("Respuesta del PANEL (resolve-route) sin db_name válido") + return None + + if action == "restore_local": + if not _validate_target(target, colocated=True): + return None + elif not _validate_target_fields(target, REQUIRED_FORWARD_TARGET_FIELDS): + return None + + app_logger.info( + f"Ruta para '{filename}': action={action}, destino={target.get('name')}, db={db_name}" + ) + return data + + +def _validate_target_fields(data: dict, fields: tuple) -> bool: + """Verifica campos requeridos no vacíos (lista explícita).""" + for field in fields: + value = data.get(field) + if value is None or (isinstance(value, str) and not value.strip()): + app_logger.error(f"Servidor del PANEL sin campo requerido para forward: '{field}'") + return False + return True + + +def report_instance_config( + api_url: str, + api_token: str, + input_folder: str, + host_name: Optional[str] = None, + app_version: Optional[str] = None, + instance_key: Optional[str] = None, +) -> bool: + """ + Reporta la carpeta de entrada vigente a POST /api/restore/instance-config (best-effort). + El panel solo la muestra; CloudRestoreAS es la única fuente de escritura. + + Returns: + True si el PANEL aceptó el reporte (200), False en cualquier otro caso. + """ + base = _normalize_base_url(api_url) + if base is None or not api_token or not api_token.strip(): + return False + if not input_folder or not input_folder.strip(): + app_logger.warning("input_folder vacío; no se reporta al PANEL") + return False + + url = f"{base}/api/restore/instance-config" + headers = {"Authorization": f"Bearer {api_token.strip()}"} + key = (instance_key or "").strip() or None + payload = { + "input_folder": input_folder.strip(), + "host_name": (host_name or "").strip() or None, + "app_version": (app_version or "").strip() or None, + } + if key: + payload["instance_key"] = key + + try: + resp = requests.post(url, json=payload, headers=headers, timeout=DEFAULT_TIMEOUT) + except requests.exceptions.RequestException as e: + app_logger.error(f"No se pudo reportar la carpeta de entrada al PANEL: {e}") + return False + + if resp.status_code == 200: + app_logger.info(f"Carpeta de entrada reportada al PANEL: {input_folder.strip()}") + return True + app_logger.error( + f"PANEL respondió {resp.status_code} al reportar la carpeta de entrada" + ) + return False + + +def report_job_result( + api_url: str, + api_token: str, + filename: str, + status: str, + restore_target_id: Optional[int] = None, + db_name: Optional[str] = None, + duration_ms: Optional[int] = None, + error_message: Optional[str] = None, +) -> bool: + """ + Reporta el resultado de un job a POST /api/restore/job-result (best-effort). + No lanza excepciones: un fallo aquí no debe romper el flujo de restauración. + + Returns: + True si el PANEL aceptó el reporte (201), False en cualquier otro caso. + """ + base = _normalize_base_url(api_url) + if base is None or not api_token or not api_token.strip(): + return False + + url = f"{base}/api/restore/job-result" + headers = {"Authorization": f"Bearer {api_token.strip()}"} + payload = { + "filename": filename, + "restore_target_id": restore_target_id, + "db_name": db_name, + "status": status, + "duration_ms": duration_ms, + "error_message": error_message, + } + + try: + resp = requests.post(url, json=payload, headers=headers, timeout=DEFAULT_TIMEOUT) + except requests.exceptions.RequestException as e: + app_logger.error(f"No se pudo reportar el resultado del job al PANEL: {e}") + return False + + if resp.status_code == 201: + return True + app_logger.error(f"PANEL respondió {resp.status_code} al reportar el resultado del job") + return False + + +def test_connection(api_url: str, api_token: str) -> tuple[bool, Optional[str]]: + """ + Prueba conectividad y autenticación contra el PANEL para la UI de configuración. + Consulta /target-for con un nombre de base inexistente: un 404 confirma que el + PANEL responde y el token es válido (la base simplemente no existe). + + Returns: + (True, mensaje) si el PANEL responde y el token es aceptado; + (False, "mensaje de error") si no se pudo conectar o el token fue rechazado. + """ + base = _normalize_base_url(api_url) + if base is None: + return False, "URL de PANEL inválida (debe iniciar con http:// o https://)" + if not api_token or not api_token.strip(): + return False, "Falta el token de API del PANEL" + + url = f"{base}/api/restore/target-for?database={quote('__cloudrestore_ping__')}" + headers = {"Authorization": f"Bearer {api_token.strip()}"} + try: + resp = requests.get(url, headers=headers, timeout=DEFAULT_TIMEOUT) + except requests.exceptions.RequestException as e: + return False, f"No se pudo conectar al PANEL: {e}" + + # 200 (base de ping existiera) o 404 (no existe) → conexión y token OK. + if resp.status_code in (200, 404): + return True, "Conexión y token correctos" + if resp.status_code == 401: + return False, "Token rechazado por el PANEL (401)" + return False, f"El PANEL respondió con código {resp.status_code}" diff --git a/app/sql/sql_manager.py b/app/sql/sql_manager.py index e05fda4..0774087 100644 --- a/app/sql/sql_manager.py +++ b/app/sql/sql_manager.py @@ -34,7 +34,8 @@ class SQLServerManager: username: Usuario SQL (si no usa Windows Auth) password: Contraseña SQL (si no usa Windows Auth) """ - self.server = server + # ODBC usa coma para el puerto (ip,puerto), no dos puntos + self.server = server.replace(":", ",") if ":" in server else server self.use_windows_auth = use_windows_auth self.username = username self.password = password @@ -168,85 +169,112 @@ class SQLServerManager: Returns: Tupla (éxito, stdout, error) """ + conn = None try: - # Construir las cláusulas MOVE - move_clauses = [] - for lf in logical_files: - if lf.type == 'D': # Data file - new_path = f"{data_folder}\\{db_name}.mdf" - elif lf.type == 'L': # Log file - new_path = f"{data_folder}\\{db_name}_log.ldf" - else: - # Archivos adicionales (filestream, etc.) - continue - - move_clauses.append(f"MOVE N'{lf.logical_name}' TO N'{new_path}'") - + # Construir las cláusulas MOVE con un destino único por archivo lógico. + # Renombrar todos los data files a {db}.mdf colisiona si el backup tiene + # varios archivos; aquí cada archivo recibe un nombre distinto (G4). + move_clauses = self._build_move_clauses(db_name, data_folder, logical_files) if not move_clauses: - return False, None, "No se pudieron determinar los archivos de datos y log" - - # Construir el comando RESTORE - restore_cmd = f""" --- Poner la base de datos en modo single user -ALTER DATABASE [{db_name}] SET SINGLE_USER WITH ROLLBACK IMMEDIATE; + return False, None, "No se pudieron determinar los archivos del backup" --- Restaurar -RESTORE DATABASE [{db_name}] -FROM DISK = N'{backup_path}' -WITH {', '.join(move_clauses)}, REPLACE; + restore_query = ( + f"RESTORE DATABASE [{db_name}] FROM DISK = N'{backup_path}' " + f"WITH {', '.join(move_clauses)}, REPLACE" + ) + + app_logger.info(f"Comando RESTORE generado:\n{restore_query}") --- Volver a modo multi user -ALTER DATABASE [{db_name}] SET MULTI_USER; -""" - - app_logger.info(f"Comando RESTORE generado:\n{restore_cmd}") - if dry_run: app_logger.info("Modo DRY RUN: No se ejecutará el RESTORE") - return True, restore_cmd, None - + return True, restore_query, None + # Ejecutar RESTORE conn_str = self.get_connection_string() conn = pyodbc.connect(conn_str, timeout=timeout_minutes * 60) conn.autocommit = True # Necesario para ALTER DATABASE cursor = conn.cursor() - + start_time = time.time() app_logger.info(f"Ejecutando RESTORE DATABASE [{db_name}]...") - - # Ejecutar en múltiples pasos + output_lines = [] - - # 1. Single user try: - cursor.execute(f"ALTER DATABASE [{db_name}] SET SINGLE_USER WITH ROLLBACK IMMEDIATE") - output_lines.append("Base de datos configurada en modo SINGLE_USER") - except Exception as e: - app_logger.warning(f"Error configurando SINGLE_USER (puede no existir la DB): {e}") - - # 2. RESTORE - restore_query = f"RESTORE DATABASE [{db_name}] FROM DISK = N'{backup_path}' WITH {', '.join(move_clauses)}, REPLACE" - cursor.execute(restore_query) - output_lines.append("RESTORE DATABASE completado") - - # 3. Multi user - cursor.execute(f"ALTER DATABASE [{db_name}] SET MULTI_USER") - output_lines.append("Base de datos configurada en modo MULTI_USER") - + # 1. Single user (puede no existir la DB en una primera restauración) + try: + cursor.execute( + f"ALTER DATABASE [{db_name}] SET SINGLE_USER WITH ROLLBACK IMMEDIATE" + ) + output_lines.append("Base de datos configurada en modo SINGLE_USER") + except Exception as e: + app_logger.warning( + f"No se pudo configurar SINGLE_USER (puede no existir la DB): {e}" + ) + + # 2. RESTORE + cursor.execute(restore_query) + output_lines.append("RESTORE DATABASE completado") + finally: + # 3. Volver SIEMPRE a MULTI_USER, incluso si el RESTORE falló, para + # no dejar la BD inaccesible en un servidor remoto compartido (G7). + try: + cursor.execute(f"ALTER DATABASE [{db_name}] SET MULTI_USER") + output_lines.append("Base de datos configurada en modo MULTI_USER") + except Exception as e: + app_logger.error( + f"No se pudo volver a MULTI_USER la BD [{db_name}] " + f"(¿quedó en estado RESTORING tras un fallo?): {e}. " + "Requiere intervención manual del DBA." + ) + elapsed = time.time() - start_time output_lines.append(f"Restauración completada en {elapsed:.2f}s") - - conn.close() - + output = "\n".join(output_lines) app_logger.info(output) - + return True, output, None - + except Exception as e: error_msg = f"Error restaurando base de datos: {str(e)}" app_logger.error(error_msg) return False, None, error_msg + finally: + if conn is not None: + try: + conn.close() + except Exception: + app_logger.warning("No se pudo cerrar la conexión SQL tras el RESTORE") + + @staticmethod + def _build_move_clauses(db_name, data_folder, logical_files) -> list: + """ + Genera una cláusula MOVE por archivo lógico con destino único, evitando + colisiones cuando el backup tiene múltiples data files o logs (G4): + - 1er data → {db}.mdf, siguientes → {db}_N.ndf + - 1er log → {db}_log.ldf, siguientes → {db}_log_N.ldf + - otros tipos (FILESTREAM/full-text) → {db}_{nombre_lógico_saneado} + """ + clauses = [] + data_idx = 0 + log_idx = 0 + for lf in logical_files: + if lf.type == 'D': + suffix = "" if data_idx == 0 else f"_{data_idx}" + ext = "mdf" if data_idx == 0 else "ndf" + new_path = f"{data_folder}\\{db_name}{suffix}.{ext}" + data_idx += 1 + elif lf.type == 'L': + suffix = "" if log_idx == 0 else f"_{log_idx}" + new_path = f"{data_folder}\\{db_name}_log{suffix}.ldf" + log_idx += 1 + else: + # No descartar otros tipos: moverlos preservando el nombre lógico. + safe = "".join(c if c.isalnum() else "_" for c in lf.logical_name) + new_path = f"{data_folder}\\{db_name}_{safe}" + + clauses.append(f"MOVE N'{lf.logical_name}' TO N'{new_path}'") + return clauses def database_exists(self, db_name: str) -> bool: """ diff --git a/app/transfer/__init__.py b/app/transfer/__init__.py new file mode 100644 index 0000000..ac474c3 --- /dev/null +++ b/app/transfer/__init__.py @@ -0,0 +1 @@ +"""Transferencia de archivos hacia los servidores SQL remotos (SMB).""" diff --git a/app/transfer/sftp_copy.py b/app/transfer/sftp_copy.py new file mode 100644 index 0000000..b1bdb00 --- /dev/null +++ b/app/transfer/sftp_copy.py @@ -0,0 +1,160 @@ +""" +Transferencia del .bak al servidor SQL externo vía SFTP/SSH (paramiko). + +Los servidores SQL destino son máquinas externas independientes (no comparten red +local con CloudRestoreAS), por lo que el .bak se sube por SFTP al servidor y el SQL +Server restaura desde su disco local (remote_inbox_path). Tras el job, el .bak del +servidor se elimina siempre (G10). + +Nota de seguridad: se usa AutoAddPolicy para las host keys (TOFU). Para endurecer en +producción conviene fijar/known_hosts las claves de cada servidor. +""" + +import ntpath +import posixpath +from pathlib import Path +from typing import Optional + +import paramiko + +from ..utils.logger import app_logger + +# Timeout de conexión SSH en segundos. +SSH_TIMEOUT = 30 + + +class SFTPCopyError(Exception): + """Error al transferir o limpiar el .bak en el servidor remoto vía SFTP.""" + + +def _sftp_path(remote_inbox_path: str, filename: str) -> str: + """ + Ruta estilo POSIX para SFTP a partir de la carpeta destino (que puede venir en + formato Windows, p. ej. C:\\RestoreInbox). OpenSSH en Windows acepta 'C:/...'. + """ + posix_dir = remote_inbox_path.replace("\\", "/").rstrip("/") + return f"{posix_dir}/{filename}" + + +def windows_restore_path(remote_inbox_path: str, filename: str) -> str: + """Ruta Windows que usará RESTORE DATABASE en el servidor (C:\\RestoreInbox\\x.bak).""" + return ntpath.join(remote_inbox_path, filename) + + +def _connect(cfg: dict) -> paramiko.SSHClient: + """Abre una conexión SSH con las credenciales del servidor (cfg del PANEL).""" + client = paramiko.SSHClient() + client.set_missing_host_key_policy(paramiko.AutoAddPolicy()) + try: + client.connect( + hostname=cfg["ssh_host"], + port=int(cfg.get("ssh_port") or 22), + username=cfg["ssh_username"], + password=cfg["ssh_password"], + timeout=SSH_TIMEOUT, + allow_agent=False, + look_for_keys=False, + ) + except Exception as e: + raise SFTPCopyError( + f"No se pudo conectar por SSH a {cfg.get('ssh_host')}:{cfg.get('ssh_port')}: {e}" + ) from e + return client + + +def upload_to_remote(local_bak: str, cfg: dict) -> str: + """ + Sube `local_bak` al servidor remoto vía SFTP, a remote_inbox_path/filename. + + Args: + local_bak: ruta local del .bak ya extraído. + cfg: dict con ssh_host, ssh_port, ssh_username, ssh_password, remote_inbox_path. + + Returns: + Ruta SFTP (POSIX) del .bak en el servidor, para limpieza posterior. + + Raises: + SFTPCopyError: si el origen no existe o falla la conexión/transferencia. + """ + remote_folder = cfg.get("remote_inbox_path") or "" + return upload_file_to_folder(local_bak, cfg, remote_folder) + + +def upload_file_to_folder(local_file: str, cfg: dict, remote_folder: str) -> str: + """ + Sube un archivo local a remote_folder/filename vía SFTP. + + Args: + local_file: ruta local del archivo. + cfg: dict con ssh_host, ssh_port, ssh_username, ssh_password. + remote_folder: carpeta destino en el servidor remoto (Windows o POSIX). + + Returns: + Ruta SFTP (POSIX) del archivo en el servidor. + + Raises: + SFTPCopyError: si el origen no existe o falla la transferencia. + """ + src = Path(local_file) + if not src.is_file(): + raise SFTPCopyError(f"El archivo local no existe: {local_file}") + if not remote_folder or not str(remote_folder).strip(): + raise SFTPCopyError("La carpeta remota destino está vacía") + + remote_sftp = _sftp_path(str(remote_folder).strip(), src.name) + client = _connect(cfg) + try: + sftp = client.open_sftp() + try: + app_logger.info(f"Subiendo archivo por SFTP a {cfg['ssh_host']}: {remote_sftp}") + sftp.put(str(src), remote_sftp) + finally: + sftp.close() + except SFTPCopyError: + raise + except Exception as e: + raise SFTPCopyError(f"Fallo al subir archivo por SFTP ({remote_sftp}): {e}") from e + finally: + client.close() + + return remote_sftp + + +def upload_zip_parts(local_paths: list[str], cfg: dict, remote_folder: str) -> list[str]: + """ + Sube uno o más archivos ZIP (incl. multipart) al input_folder del destino. + + Returns: + Lista de rutas SFTP subidas. + """ + uploaded: list[str] = [] + for local_path in local_paths: + uploaded.append(upload_file_to_folder(local_path, cfg, remote_folder)) + return uploaded + + +def cleanup_remote(cfg: dict, remote_sftp_path: Optional[str]) -> None: + """ + Elimina el .bak del servidor remoto vía SFTP. Best-effort: registra pero no lanza, + para no enmascarar el resultado real del job. Se invoca siempre en `finally`. + """ + if not remote_sftp_path: + return + try: + client = _connect(cfg) + except SFTPCopyError as e: + app_logger.error(f"No se pudo conectar para limpiar el .bak remoto: {e}") + return + try: + sftp = client.open_sftp() + try: + sftp.remove(remote_sftp_path) + app_logger.info(f"Inbox remoto limpiado: {remote_sftp_path}") + except FileNotFoundError: + pass + finally: + sftp.close() + except Exception as e: + app_logger.error(f"No se pudo limpiar el .bak remoto ({remote_sftp_path}): {e}") + finally: + client.close() diff --git a/app/ui/config_tab.py b/app/ui/config_tab.py index 7f55cad..ae39644 100644 --- a/app/ui/config_tab.py +++ b/app/ui/config_tab.py @@ -3,13 +3,14 @@ from PySide6.QtWidgets import ( QWidget, QVBoxLayout, QHBoxLayout, QGroupBox, QFormLayout, QLineEdit, QPushButton, QSpinBox, - QCheckBox, QFileDialog, QMessageBox, QScrollArea + QCheckBox, QFileDialog, QMessageBox, QScrollArea, QComboBox ) from pathlib import Path from ..utils.crypto import encrypt_password, decrypt_password from ..extract.seven_zip import SevenZipExtractor from ..sql.sql_manager import SQLServerManager +from ..panel import panel_client class ConfigTab(QWidget): @@ -20,6 +21,7 @@ class ConfigTab(QWidget): super().__init__() self.engine = engine self.config = {} + self._catalog_names: list[str] = [] self._setup_ui() def _setup_ui(self): @@ -118,7 +120,38 @@ class ConfigTab(QWidget): sql_group.setLayout(sql_layout) layout.addWidget(sql_group) - + + # Grupo: PANEL de Control (servidor de restauración activo) + panel_group = QGroupBox("PANEL de Control") + panel_layout = QFormLayout() + + self.panel_url_input = QLineEdit() + self.panel_url_input.setPlaceholderText("http://ip-del-panel:3000") + panel_layout.addRow("URL del PANEL:", self.panel_url_input) + + self.panel_token_input = QLineEdit() + self.panel_token_input.setEchoMode(QLineEdit.EchoMode.Password) + panel_layout.addRow("API Token:", self.panel_token_input) + + self.panel_instance_combo = QComboBox() + self.panel_instance_combo.setToolTip( + "Nombre de este servidor en el panel. El CRA restaurará aquí lo que le " + "corresponda y reenviará el resto por SFTP automáticamente." + ) + self.refresh_instance_btn = QPushButton("Actualizar lista de servidores") + self.refresh_instance_btn.clicked.connect(self._refresh_instance_combo) + instance_row = QHBoxLayout() + instance_row.addWidget(self.panel_instance_combo, stretch=1) + instance_row.addWidget(self.refresh_instance_btn) + panel_layout.addRow("Instancia (servidor):", instance_row) + + test_panel_btn = QPushButton("🔌 Probar Conexión al PANEL") + test_panel_btn.clicked.connect(self._test_panel_connection) + panel_layout.addRow("", test_panel_btn) + + panel_group.setLayout(panel_layout) + layout.addWidget(panel_group) + # Grupo: Concurrencia concurrency_group = QGroupBox("Concurrencia") concurrency_layout = QFormLayout() @@ -269,7 +302,42 @@ class ConfigTab(QWidget): self.dry_run_checkbox.setChecked(features.get("dry_run_mode", False)) self.auto_scan_checkbox.setChecked(features.get("auto_scan_enabled", True)) self.scan_interval_spinbox.setValue(features.get("scan_interval_seconds", 30)) + + # PANEL de Control + panel = config.get("panel", {}) + self.panel_url_input.setText(panel.get("api_url", "")) + self.panel_token_input.setText(panel.get("api_token", "")) + self._refresh_instance_combo() + def _refresh_instance_combo(self): + """Recarga los nombres de servidores desde el catálogo del panel.""" + current = self.panel_instance_combo.currentText().strip() + saved_key = (self.config.get("panel", {}).get("instance_key") or "").strip() + preserve = current or saved_key + + api_url = self.panel_url_input.text().strip() + api_token = self.panel_token_input.text().strip() + + names: list[str] = [] + if api_url and api_token: + names = panel_client.list_restore_target_names(api_url, api_token) + + self._catalog_names = names + + items = [""] + for name in names: + if name not in items: + items.append(name) + if preserve and preserve not in items: + items.append(preserve) + + self.panel_instance_combo.blockSignals(True) + self.panel_instance_combo.clear() + self.panel_instance_combo.addItems(items) + idx = self.panel_instance_combo.findText(preserve) + self.panel_instance_combo.setCurrentIndex(idx if idx >= 0 else 0) + self.panel_instance_combo.blockSignals(False) + def _save_config(self): """Guarda la configuración.""" # Construir config @@ -307,9 +375,31 @@ class ConfigTab(QWidget): "dry_run_mode": self.dry_run_checkbox.isChecked(), "auto_scan_enabled": self.auto_scan_checkbox.isChecked(), "scan_interval_seconds": self.scan_interval_spinbox.value() + }, + "panel": { + "api_url": self.panel_url_input.text().strip(), + "api_token": self.panel_token_input.text().strip(), + "instance_key": self.panel_instance_combo.currentText().strip(), } } + panel_url = self.panel_url_input.text().strip() + panel_token = self.panel_token_input.text().strip() + instance_key = self.panel_instance_combo.currentText().strip() + if panel_url and panel_token and not instance_key: + QMessageBox.warning( + self, + "Advertencia", + "Con panel configurado debe seleccionar la instancia (servidor).", + ) + if panel_url and instance_key and self._catalog_names and instance_key not in self._catalog_names: + QMessageBox.warning( + self, + "Advertencia", + f'La instancia "{instance_key}" no está en el catálogo del panel. ' + "Verifique el nombre o cree el servidor en Servidores de Restauración.", + ) + # Cifrar password si no usa Windows Auth if not self.sql_windows_auth_checkbox.isChecked(): password = self.sql_password_input.text() @@ -356,6 +446,25 @@ class ConfigTab(QWidget): "No se pudo auto-detectar 7-Zip. Por favor, selecciona manualmente." ) + def _test_panel_connection(self): + """Prueba la conexión al PANEL y muestra el servidor de restauración activo.""" + api_url = self.panel_url_input.text().strip() + api_token = self.panel_token_input.text().strip() + + if not api_url: + QMessageBox.warning(self, "Error", "Debe especificar la URL del PANEL.") + return + + success, info = panel_client.test_connection(api_url, api_token) + + if success: + self._refresh_instance_combo() + QMessageBox.information( + self, "Éxito", f"Conexión al PANEL exitosa.\n{info or ''}" + ) + else: + QMessageBox.critical(self, "Error", f"No se pudo conectar al PANEL:\n{info}") + def _on_auth_changed(self, checked: bool): """Maneja el cambio en el tipo de autenticación.""" self.sql_username_input.setEnabled(not checked) diff --git a/app/ui/nodes_tab.py b/app/ui/nodes_tab.py index adc177b..2aad9e5 100644 --- a/app/ui/nodes_tab.py +++ b/app/ui/nodes_tab.py @@ -4,7 +4,7 @@ from PySide6.QtWidgets import ( QWidget, QVBoxLayout, QHBoxLayout, QTableWidget, QTableWidgetItem, QPushButton, QDialog, QFormLayout, QLineEdit, QTextEdit, QCheckBox, QDialogButtonBox, - QMessageBox + QMessageBox, QLabel ) from PySide6.QtCore import Qt diff --git a/requirements.txt b/requirements.txt index 30aba27..ab8c6a0 100644 --- a/requirements.txt +++ b/requirements.txt @@ -9,5 +9,14 @@ pyodbc>=5.0.0 # pywin32 para DPAPI (cifrado de passwords) pywin32>=306 +# requests para consultar al PANEL el servidor de restauración de cada base +requests>=2.31.0 + +# paramiko para transferir el .bak por SFTP/SSH a los servidores SQL externos +paramiko>=3.4.0 + # Para empaquetado (opcional) pyinstaller>=6.0.0 + +# Pruebas (opcional) +pytest>=8.0.0 diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..57495fd --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,5 @@ +"""Configuración de pytest: asegura que el paquete `app` sea importable.""" +import os +import sys + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) diff --git a/tests/test_colocated_resolve.py b/tests/test_colocated_resolve.py new file mode 100644 index 0000000..f5d16e3 --- /dev/null +++ b/tests/test_colocated_resolve.py @@ -0,0 +1,41 @@ +""" +Pruebas de validación modo colocado (target.name vs instance_key). +""" +import pytest + +from app.panel import panel_client + + +def test_validate_colocado_acepta_sin_ssh(): + data = { + "id": 1, + "name": "Omega", + "server": "192.168.1.10", + "username": "sa", + "password": "x", + "data_folder": "C:\\SQLData", + } + assert panel_client._validate_target(data, colocated=True) is True + + +def test_validate_colocado_rechaza_sin_name(): + data = { + "id": 1, + "server": "192.168.1.10", + "username": "sa", + "password": "x", + "data_folder": "C:\\SQLData", + } + assert panel_client._validate_target(data, colocated=True) is False + + +def test_validate_hub_exige_ssh(): + data = { + "id": 1, + "name": "Omega", + "server": "192.168.1.10", + "username": "sa", + "password": "x", + "data_folder": "C:\\SQLData", + } + assert panel_client._validate_target(data, colocated=False) is False diff --git a/tests/test_panel_client.py b/tests/test_panel_client.py new file mode 100644 index 0000000..dad29ab --- /dev/null +++ b/tests/test_panel_client.py @@ -0,0 +1,381 @@ +""" +Pruebas del cliente del PANEL (servidor de restauración asignado por base de datos). +Se mockea `requests` para no depender de un PANEL real. +""" +import pytest + +from app.panel import panel_client + + +class FakeResponse: + def __init__(self, status_code, json_data=None, raise_json=False): + self.status_code = status_code + self._json = json_data + self._raise_json = raise_json + + def json(self): + if self._raise_json: + raise ValueError("no es JSON") + return self._json + + +VALID_TARGET = { + "id": 2, + "name": "Omega", + "server": "192.168.1.100,1433", + "username": "sa", + "password": "secreto", + "data_folder": "C:\\SQLData", + "ssh_host": "192.168.1.100", + "ssh_port": 22, + "ssh_username": "Administrator", + "ssh_password": "ssh-secreto", + "remote_inbox_path": "C:\\RestoreInbox", +} + +URL = "http://panel:3000" +TOKEN = "tok" +DB = "EMPRESA_DB" + + +def test_target_url_vacia_devuelve_none(): + assert panel_client.get_target_for_database("", TOKEN, DB) is None + + +def test_target_url_invalida_devuelve_none(): + assert panel_client.get_target_for_database("ftp://x", TOKEN, DB) is None + + +def test_target_token_vacio_devuelve_none(): + assert panel_client.get_target_for_database(URL, "", DB) is None + + +def test_target_db_vacio_devuelve_none(): + assert panel_client.get_target_for_database(URL, TOKEN, " ") is None + + +def test_target_ok(monkeypatch): + captured = {} + + def fake_get(url, headers=None, timeout=None): + captured["url"] = url + return FakeResponse(200, VALID_TARGET) + + monkeypatch.setattr(panel_client.requests, "get", fake_get) + target = panel_client.get_target_for_database(URL, TOKEN, DB) + assert target is not None + assert target["name"] == "Omega" + # el db_name viaja como query param url-encoded + assert "database=EMPRESA_DB" in captured["url"] + + +def test_target_db_con_espacios_se_encodea(monkeypatch): + captured = {} + monkeypatch.setattr( + panel_client.requests, "get", + lambda url, **k: (captured.__setitem__("url", url), FakeResponse(200, VALID_TARGET))[1] + ) + panel_client.get_target_for_database(URL, TOKEN, "MI BASE") + assert "MI%20BASE" in captured["url"] + + +def test_target_falta_campo_requerido(monkeypatch): + incompleto = dict(VALID_TARGET) + del incompleto["password"] + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(200, incompleto)) + assert panel_client.get_target_for_database(URL, TOKEN, DB) is None + + +def test_target_campo_vacio_es_invalido(monkeypatch): + vacio = dict(VALID_TARGET, password=" ") + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(200, vacio)) + assert panel_client.get_target_for_database(URL, TOKEN, DB) is None + + +@pytest.mark.parametrize("code", [404, 401, 500, 503]) +def test_target_codigos_no_200_devuelven_none(monkeypatch, code): + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(code)) + assert panel_client.get_target_for_database(URL, TOKEN, DB) is None + + +def test_target_error_red_devuelve_none(monkeypatch): + def boom(*a, **k): + raise panel_client.requests.exceptions.ConnectionError("caído") + + monkeypatch.setattr(panel_client.requests, "get", boom) + assert panel_client.get_target_for_database(URL, TOKEN, DB) is None + + +def test_target_json_invalido_devuelve_none(monkeypatch): + monkeypatch.setattr( + panel_client.requests, "get", lambda *a, **k: FakeResponse(200, raise_json=True) + ) + assert panel_client.get_target_for_database(URL, TOKEN, DB) is None + + +def test_report_job_result_201_true(monkeypatch): + captured = {} + + def fake_post(url, json=None, headers=None, timeout=None): + captured["json"] = json + return FakeResponse(201) + + monkeypatch.setattr(panel_client.requests, "post", fake_post) + ok = panel_client.report_job_result( + URL, TOKEN, filename="empresa.bak", status="completed", + restore_target_id=2, db_name="EMP", duration_ms=1000 + ) + assert ok is True + assert captured["json"]["filename"] == "empresa.bak" + assert captured["json"]["status"] == "completed" + + +def test_report_job_result_otro_codigo_false(monkeypatch): + monkeypatch.setattr(panel_client.requests, "post", lambda *a, **k: FakeResponse(500)) + assert panel_client.report_job_result(URL, TOKEN, "x.bak", "failed") is False + + +def test_report_job_result_error_red_false(monkeypatch): + def boom(*a, **k): + raise panel_client.requests.exceptions.Timeout("timeout") + + monkeypatch.setattr(panel_client.requests, "post", boom) + assert panel_client.report_job_result(URL, TOKEN, "x.bak", "failed") is False + + +@pytest.mark.parametrize("code", [200, 404]) +def test_test_connection_ok(monkeypatch, code): + # 200 o 404 → conexión y token correctos + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(code)) + ok, info = panel_client.test_connection(URL, TOKEN) + assert ok is True + + +def test_test_connection_token_rechazado(monkeypatch): + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(401)) + ok, info = panel_client.test_connection(URL, TOKEN) + assert ok is False + assert "401" in info + + +def test_test_connection_url_invalida(): + ok, info = panel_client.test_connection("noesurl", TOKEN) + assert ok is False + + +def test_report_instance_config_200_true(monkeypatch): + captured = {} + + def fake_post(url, json=None, headers=None, timeout=None): + captured["url"] = url + captured["json"] = json + return FakeResponse(200) + + monkeypatch.setattr(panel_client.requests, "post", fake_post) + ok = panel_client.report_instance_config( + URL, TOKEN, r"D:\Backups\Entrada", host_name="WIN-01", app_version="1.0.0" + ) + assert ok is True + assert captured["url"].endswith("/api/restore/instance-config") + assert captured["json"]["input_folder"] == r"D:\Backups\Entrada" + assert captured["json"]["host_name"] == "WIN-01" + + +def test_report_instance_config_input_vacio_false(): + assert panel_client.report_instance_config(URL, TOKEN, " ") is False + + +def test_report_instance_config_sin_token_false(): + assert panel_client.report_instance_config(URL, "", r"D:\In") is False + + +def test_report_instance_config_otro_codigo_false(monkeypatch): + monkeypatch.setattr(panel_client.requests, "post", lambda *a, **k: FakeResponse(500)) + assert panel_client.report_instance_config(URL, TOKEN, r"D:\In") is False + + +def test_report_instance_config_error_red_false(monkeypatch): + def boom(*a, **k): + raise panel_client.requests.exceptions.ConnectionError("caído") + + monkeypatch.setattr(panel_client.requests, "post", boom) + assert panel_client.report_instance_config(URL, TOKEN, r"D:\In") is False + + +def test_target_con_instance_envia_param(monkeypatch): + captured = {} + monkeypatch.setattr( + panel_client.requests, "get", + lambda url, **k: (captured.__setitem__("url", url), FakeResponse(200, VALID_TARGET))[1] + ) + panel_client.get_target_for_database(URL, TOKEN, DB, instance_key="Alfa") + assert "instance=Alfa" in captured["url"] + + +COLOCATED_TARGET = { + "id": 1, + "name": "Alfa", + "server": "localhost", + "username": "sa", + "password": "secreto", + "data_folder": "C:\\SQLData", +} + + +def test_target_colocado_no_exige_ssh(monkeypatch): + monkeypatch.setattr( + panel_client.requests, "get", lambda *a, **k: FakeResponse(200, COLOCATED_TARGET) + ) + target = panel_client.get_target_for_database(URL, TOKEN, DB, instance_key="Alfa") + assert target is not None + assert target["name"] == "Alfa" + + +def test_report_instance_config_con_instance_key(monkeypatch): + captured = {} + monkeypatch.setattr( + panel_client.requests, "post", + lambda url, json=None, **k: (captured.update({"json": json}), FakeResponse(200))[1] + ) + ok = panel_client.report_instance_config( + URL, TOKEN, r"D:\Alfa\In", instance_key="Alfa" + ) + assert ok is True + assert captured["json"]["instance_key"] == "Alfa" + + +CATALOG_RESPONSE = { + "targets": [ + {"id": 1, "name": "Alfa"}, + {"id": 2, "name": "Omega"}, + {"id": 4, "name": "Delta"}, + ] +} + + +def test_list_restore_target_names_ok(monkeypatch): + captured = {} + + def fake_get(url, headers=None, timeout=None): + captured["url"] = url + return FakeResponse(200, CATALOG_RESPONSE) + + monkeypatch.setattr(panel_client.requests, "get", fake_get) + names = panel_client.list_restore_target_names(URL, TOKEN) + assert names == ["Alfa", "Omega", "Delta"] + assert captured["url"].endswith("/api/restore/target-catalog") + + +def test_list_restore_target_names_url_vacia(): + assert panel_client.list_restore_target_names("", TOKEN) == [] + + +def test_list_restore_target_names_token_vacio(): + assert panel_client.list_restore_target_names(URL, "") == [] + + +@pytest.mark.parametrize("code", [401, 500, 503]) +def test_list_restore_target_names_codigos_no_200(monkeypatch, code): + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(code)) + assert panel_client.list_restore_target_names(URL, TOKEN) == [] + + +def test_list_restore_target_names_error_red(monkeypatch): + def boom(*a, **k): + raise panel_client.requests.exceptions.ConnectionError("caído") + + monkeypatch.setattr(panel_client.requests, "get", boom) + assert panel_client.list_restore_target_names(URL, TOKEN) == [] + + +def test_list_restore_target_names_json_invalido(monkeypatch): + monkeypatch.setattr( + panel_client.requests, "get", lambda *a, **k: FakeResponse(200, raise_json=True) + ) + assert panel_client.list_restore_target_names(URL, TOKEN) == [] + + +def test_list_restore_target_names_sin_lista_targets(monkeypatch): + monkeypatch.setattr( + panel_client.requests, "get", lambda *a, **k: FakeResponse(200, {"foo": []}) + ) + assert panel_client.list_restore_target_names(URL, TOKEN) == [] + + +ROUTE_FORWARD = { + "action": "forward", + "db_name": "GENERICA-TEST", + "node_key": "GENERICA-TEST", + "target": { + "id": 2, + "name": "Alfa", + "server": "192.168.1.10,1433", + "username": "sa", + "password": "secreto", + "data_folder": "C:\\SQLData", + "ssh_host": "192.168.1.10", + "ssh_port": 22, + "ssh_username": "Administrator", + "ssh_password": "ssh-secreto", + "remote_inbox_path": "C:\\RestoreInbox", + "input_folder": "D:\\Restore\\Alfa\\Entrada", + }, +} + +ROUTE_RESTORE_LOCAL = { + "action": "restore_local", + "db_name": "GENERICA-TEST", + "node_key": "GENERICA-TEST", + "target": { + "id": 2, + "name": "Alfa", + "server": "192.168.1.10,1433", + "username": "sa", + "password": "secreto", + "data_folder": "C:\\SQLData", + "ssh_host": "192.168.1.10", + "ssh_port": 22, + "ssh_username": "Administrator", + "ssh_password": "ssh-secreto", + "remote_inbox_path": "C:\\RestoreInbox", + "input_folder": "D:\\Restore\\Alfa\\Entrada", + }, +} + + +def test_resolve_route_forward_ok(monkeypatch): + captured = {} + + def fake_get(url, headers=None, timeout=None): + captured["url"] = url + return FakeResponse(200, ROUTE_FORWARD) + + monkeypatch.setattr(panel_client.requests, "get", fake_get) + route = panel_client.resolve_route(URL, TOKEN, "GENERICA-TEST.ZIP", instance_key="Omega") + assert route is not None + assert route["action"] == "forward" + assert "resolve-route" in captured["url"] + assert "instance=Omega" in captured["url"] + + +def test_resolve_route_restore_local_ok(monkeypatch): + monkeypatch.setattr( + panel_client.requests, "get", lambda *a, **k: FakeResponse(200, ROUTE_RESTORE_LOCAL) + ) + route = panel_client.resolve_route(URL, TOKEN, "GENERICA-TEST.ZIP", instance_key="Alfa") + assert route is not None + assert route["action"] == "restore_local" + + +def test_resolve_route_503_devuelve_none(monkeypatch): + monkeypatch.setattr(panel_client.requests, "get", lambda *a, **k: FakeResponse(503)) + assert panel_client.resolve_route(URL, TOKEN, "X.ZIP") is None + + +def test_resolve_route_forward_sin_input_folder_invalido(monkeypatch): + bad = dict(ROUTE_FORWARD) + bad["target"] = dict(bad["target"], input_folder=" ") + monkeypatch.setattr( + panel_client.requests, "get", lambda *a, **k: FakeResponse(200, bad) + ) + assert panel_client.resolve_route(URL, TOKEN, "X.ZIP") is None diff --git a/tests/test_route_forward.py b/tests/test_route_forward.py new file mode 100644 index 0000000..cb10df5 --- /dev/null +++ b/tests/test_route_forward.py @@ -0,0 +1,138 @@ +""" +Pruebas de enrutamiento automático (restore_local | forward) en RestoreWorker. +""" +from pathlib import Path + +import pytest + +from app.engine.restore_worker import RestoreDeferred, RestoreWorker + + +@pytest.fixture +def base_config(tmp_path): + return { + "paths": { + "input_folder": str(tmp_path / "in"), + "processed_folder": str(tmp_path / "processed"), + "failed_folder": str(tmp_path / "failed"), + "extract_folder": str(tmp_path / "extract"), + "data_sql_folder": str(tmp_path / "data"), + "seven_zip_exe": "C:\\Program Files\\7-Zip\\7z.exe", + }, + "sql": {"server": "localhost", "use_windows_auth": True}, + "timeouts": {"extract_minutes": 30, "restore_minutes": 60}, + "panel": { + "api_url": "http://panel:3000", + "api_token": "tok", + "instance_key": "Alfa", + }, + } + + +def test_collect_zip_paths_simple(tmp_path, base_config): + z = tmp_path / "in" / "NODO.ZIP" + z.parent.mkdir(parents=True) + z.write_bytes(b"z") + worker = RestoreWorker("job-1", base_config) + paths = worker._collect_zip_paths(str(z)) + assert paths == [str(z)] + + +def test_collect_zip_paths_multipart(tmp_path, base_config): + folder = tmp_path / "in" + folder.mkdir(parents=True) + p1 = folder / "NODO.ZIP.001" + p2 = folder / "NODO.ZIP.002" + p1.write_bytes(b"1") + p2.write_bytes(b"2") + worker = RestoreWorker("job-1", base_config) + paths = worker._collect_zip_paths(str(p1)) + assert len(paths) == 2 + assert all(Path(p).exists() for p in paths) + + +def test_resolve_route_envia_instance_key(monkeypatch, base_config): + worker = RestoreWorker("job-1", base_config) + captured = {} + + def fake_resolve(api_url, api_token, filename, instance_key=None): + captured["instance_key"] = instance_key + return { + "action": "restore_local", + "db_name": "DB1", + "node_key": "NODO", + "target": { + "id": 1, + "name": "Alfa", + "server": "localhost", + "username": "sa", + "password": "x", + "data_folder": "C:\\Data", + }, + } + + monkeypatch.setattr( + "app.engine.restore_worker.panel_client.resolve_route", fake_resolve + ) + monkeypatch.setattr( + "app.engine.restore_worker.JobRepository.update_node_and_db", lambda *a, **k: None + ) + + class FakeJob: + source_name = "NODO.ZIP" + source_path = "C:\\in\\NODO.ZIP" + + route = worker._resolve_route(FakeJob()) + assert route["action"] == "restore_local" + assert captured["instance_key"] == "Alfa" + + +def test_resolve_route_sin_instance_key_diferido(base_config): + base_config["panel"]["instance_key"] = "" + worker = RestoreWorker("job-1", base_config) + + class FakeJob: + source_name = "NODO.ZIP" + source_path = "C:\\in\\NODO.ZIP" + + with pytest.raises(RestoreDeferred, match="instance_key"): + worker._resolve_route(FakeJob()) + + +def test_legacy_mode_no_cambia_comportamiento(monkeypatch, base_config): + base_config["panel"]["mode"] = "orchestrator" + worker = RestoreWorker("job-1", base_config) + assert worker._legacy_panel_mode() == "orchestrator" + + captured = {} + + def fake_resolve(api_url, api_token, filename, instance_key=None): + captured["instance_key"] = instance_key + return { + "action": "forward", + "db_name": "DB1", + "node_key": "NODO", + "target": { + "id": 2, + "name": "Omega", + "ssh_host": "h", + "ssh_username": "u", + "ssh_password": "p", + "input_folder": "D:\\In", + }, + } + + monkeypatch.setattr( + "app.engine.restore_worker.panel_client.resolve_route", fake_resolve + ) + monkeypatch.setattr( + "app.engine.restore_worker.JobRepository.update_node_and_db", lambda *a, **k: None + ) + + class FakeJob: + source_name = "NODO.ZIP" + source_path = "C:\\in\\NODO.ZIP" + + route = worker._resolve_route(FakeJob()) + assert route["action"] == "forward" + assert captured["instance_key"] == "Alfa" diff --git a/tests/test_sftp_copy.py b/tests/test_sftp_copy.py new file mode 100644 index 0000000..7037545 --- /dev/null +++ b/tests/test_sftp_copy.py @@ -0,0 +1,135 @@ +""" +Pruebas de la transferencia SFTP al servidor remoto. Se mockea paramiko para no +requerir un servidor SSH real; se valida la conversión de rutas y el flujo de subida. +""" +import pytest + +from app.transfer import sftp_copy +from app.transfer.sftp_copy import SFTPCopyError + +CFG = { + "ssh_host": "192.168.1.100", + "ssh_port": 22, + "ssh_username": "Administrator", + "ssh_password": "ssh-secreto", + "remote_inbox_path": "C:\\RestoreInbox", +} + + +def test_windows_restore_path(): + assert sftp_copy.windows_restore_path("C:\\RestoreInbox", "empresa.bak") == "C:\\RestoreInbox\\empresa.bak" + + +def test_sftp_path_convierte_windows_a_posix(): + # ruta destino en formato Windows → SFTP usa forward slashes + assert sftp_copy._sftp_path("C:\\RestoreInbox", "empresa.bak") == "C:/RestoreInbox/empresa.bak" + # ya en posix se respeta (sin doble slash) + assert sftp_copy._sftp_path("/srv/inbox/", "x.bak") == "/srv/inbox/x.bak" + + +def test_upload_origen_inexistente(tmp_path): + with pytest.raises(SFTPCopyError, match="no existe"): + sftp_copy.upload_to_remote(str(tmp_path / "noexiste.bak"), CFG) + + +class FakeSFTP: + def __init__(self, store): + self.store = store + + def put(self, local, remote): + self.store["put"] = (local, remote) + + def remove(self, remote): + self.store["removed"] = remote + + def close(self): + self.store["sftp_closed"] = True + + +class FakeClient: + def __init__(self, store): + self.store = store + + def set_missing_host_key_policy(self, policy): + self.store["policy_set"] = True + + def connect(self, **kwargs): + self.store["connect"] = kwargs + + def open_sftp(self): + return FakeSFTP(self.store) + + def close(self): + self.store["client_closed"] = True + + +def _fake_paramiko(store): + class FakeParamiko: + AutoAddPolicy = object + SSHClient = lambda self=None: FakeClient(store) + fp = FakeParamiko() + # SSHClient() debe devolver FakeClient + fp.SSHClient = lambda: FakeClient(store) + return fp + + +def test_upload_ok(tmp_path, monkeypatch): + bak = tmp_path / "empresa.bak" + bak.write_bytes(b"data") + store: dict = {} + monkeypatch.setattr(sftp_copy, "paramiko", _fake_paramiko(store)) + + remote = sftp_copy.upload_to_remote(str(bak), CFG) + assert remote == "C:/RestoreInbox/empresa.bak" + assert store["put"][1] == "C:/RestoreInbox/empresa.bak" + assert store["connect"]["hostname"] == "192.168.1.100" + assert store["connect"]["port"] == 22 + assert store["client_closed"] is True + + +def test_cleanup_remote_borra(monkeypatch): + store: dict = {} + monkeypatch.setattr(sftp_copy, "paramiko", _fake_paramiko(store)) + sftp_copy.cleanup_remote(CFG, "C:/RestoreInbox/empresa.bak") + assert store["removed"] == "C:/RestoreInbox/empresa.bak" + assert store["client_closed"] is True + + +def test_cleanup_remote_none_no_falla(): + sftp_copy.cleanup_remote(CFG, None) # no debe conectar ni lanzar + + +def test_upload_file_to_folder_ok(tmp_path, monkeypatch): + zf = tmp_path / "backup.zip" + zf.write_bytes(b"zipdata") + store: dict = {} + monkeypatch.setattr(sftp_copy, "paramiko", _fake_paramiko(store)) + + remote = sftp_copy.upload_file_to_folder( + str(zf), CFG, "D:\\Restore\\Alfa\\Entrada" + ) + assert remote == "D:/Restore/Alfa/Entrada/backup.zip" + assert store["put"][1] == "D:/Restore/Alfa/Entrada/backup.zip" + + +def test_upload_zip_parts_multipart(tmp_path, monkeypatch): + p1 = tmp_path / "big.zip.001" + p2 = tmp_path / "big.zip.002" + p1.write_bytes(b"a") + p2.write_bytes(b"b") + store: dict = {} + monkeypatch.setattr(sftp_copy, "paramiko", _fake_paramiko(store)) + + uploaded = sftp_copy.upload_zip_parts( + [str(p1), str(p2)], CFG, "D:\\In" + ) + assert len(uploaded) == 2 + assert uploaded[0].endswith("/big.zip.001") + assert uploaded[1].endswith("/big.zip.002") + + +def test_upload_file_to_folder_vacio_falla(tmp_path): + f = tmp_path / "x.zip" + f.write_bytes(b"x") + with pytest.raises(SFTPCopyError, match="carpeta remota"): + sftp_copy.upload_file_to_folder(str(f), CFG, " ") diff --git a/tests/test_sql_manager_move.py b/tests/test_sql_manager_move.py new file mode 100644 index 0000000..fb422fc --- /dev/null +++ b/tests/test_sql_manager_move.py @@ -0,0 +1,77 @@ +""" +Pruebas de _build_move_clauses: nombres únicos por archivo lógico (G4), +evitando colisiones con múltiples data files / logs y sin descartar otros tipos. +""" +from app.sql.sql_manager import SQLServerManager, LogicalFile + + +def _paths(clauses): + """Extrae las rutas destino (lo que va tras 'TO N') de cada cláusula MOVE.""" + out = [] + for c in clauses: + # MOVE N'logico' TO N'ruta' + ruta = c.split("TO N'")[1].rstrip("'") + out.append(ruta) + return out + + +def test_un_data_un_log(): + files = [ + LogicalFile("EMP_dat", "X.mdf", "D"), + LogicalFile("EMP_log", "X.ldf", "L"), + ] + clauses = SQLServerManager._build_move_clauses("EMP", "C:\\D", files) + rutas = _paths(clauses) + assert rutas == ["C:\\D\\EMP.mdf", "C:\\D\\EMP_log.ldf"] + + +def test_multiples_data_files_sin_colision(): + files = [ + LogicalFile("d1", "a.mdf", "D"), + LogicalFile("d2", "b.ndf", "D"), + LogicalFile("d3", "c.ndf", "D"), + LogicalFile("l1", "a.ldf", "L"), + ] + clauses = SQLServerManager._build_move_clauses("EMP", "C:\\D", files) + rutas = _paths(clauses) + assert rutas == [ + "C:\\D\\EMP.mdf", + "C:\\D\\EMP_1.ndf", + "C:\\D\\EMP_2.ndf", + "C:\\D\\EMP_log.ldf", + ] + # No debe haber rutas duplicadas (la causa del bug original). + assert len(set(rutas)) == len(rutas) + + +def test_multiples_logs_sin_colision(): + files = [ + LogicalFile("d1", "a.mdf", "D"), + LogicalFile("l1", "a.ldf", "L"), + LogicalFile("l2", "b.ldf", "L"), + ] + rutas = _paths(SQLServerManager._build_move_clauses("EMP", "C:\\D", files)) + assert rutas == ["C:\\D\\EMP.mdf", "C:\\D\\EMP_log.ldf", "C:\\D\\EMP_log_1.ldf"] + assert len(set(rutas)) == len(rutas) + + +def test_tipo_no_data_ni_log_no_se_descarta(): + # Tipos como FILESTREAM/full-text (p. ej. 'S') deben moverse, no ignorarse. + files = [ + LogicalFile("d1", "a.mdf", "D"), + LogicalFile("l1", "a.ldf", "L"), + LogicalFile("fs stream!", "fs", "S"), + ] + clauses = SQLServerManager._build_move_clauses("EMP", "C:\\D", files) + assert len(clauses) == 3 # ninguno descartado + rutas = _paths(clauses) + # El nombre lógico se sanea (no alfanumérico → '_'). + assert "C:\\D\\EMP_fs_stream_" in rutas[2] + assert len(set(rutas)) == len(rutas) + + +def test_cada_archivo_logico_genera_una_clausula(): + files = [LogicalFile(f"f{i}", f"f{i}", "D") for i in range(5)] + clauses = SQLServerManager._build_move_clauses("DB", "C:\\D", files) + assert len(clauses) == 5 + assert len(set(_paths(clauses))) == 5 # todas únicas