2 Commits

Author SHA1 Message Date
072be5b5db feature/integracion-panel-restore-targets 2026-06-05 12:28:04 -06:00
afb76ce4a6 feature/integracion-panel-restore-targets 2026-06-05 10:49:05 -06:00
20 changed files with 2167 additions and 104 deletions

197
INTEGRACION_PANEL.md Normal file
View File

@@ -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=<esta instancia>` — 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 <CLOUDRESTORE_API_TOKEN>`.
### 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=<zip>&instance=<nombre>`
**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=<db_name>&instance=<nombre>`
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**.

View File

@@ -20,6 +20,19 @@ Get-OdbcDriver | Where-Object {$_.Name -like "*SQL Server*"}
.\venv\Scripts\python.exe runner.py .\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.
**Prueba local en tu PC** (carpetas + config automática):
```powershell
.\scripts\prepare_local_test.ps1
.\venv\Scripts\python.exe runner.py
```
Copia el ZIP a `C:\CloudRestore\Entrada` (ver `LEEME.txt` en esa carpeta).
### 3. Configuración Básica (en la aplicación) ### 3. Configuración Básica (en la aplicación)
1. **Tab Configuración**: 1. **Tab Configuración**:

View File

@@ -39,6 +39,8 @@ class StepType:
FILELIST = "filelist" FILELIST = "filelist"
RESTORE = "restore" RESTORE = "restore"
CLEANUP = "cleanup" 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 # Configuración por defecto
DEFAULT_CONFIG = { DEFAULT_CONFIG = {
@@ -78,5 +80,16 @@ DEFAULT_CONFIG = {
"dry_run_mode": False, "dry_run_mode": False,
"auto_scan_enabled": True, "auto_scan_enabled": True,
"scan_interval_seconds": 30 "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": "",
# False en dev: el panel local usa HTTPS con certificado propio.
"verify_ssl": False,
} }
} }

View File

@@ -176,6 +176,15 @@ class JobRepository:
tuple(params) 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 @staticmethod
def exists_by_hash(source_hash: str) -> bool: def exists_by_hash(source_hash: str) -> bool:
"""Verifica si existe un job con el hash dado.""" """Verifica si existe un job con el hash dado."""

View File

@@ -1,15 +1,19 @@
"""Motor principal de la aplicación.""" """Motor principal de la aplicación."""
import platform
import socket
from typing import Optional from typing import Optional
from pathlib import Path from pathlib import Path
from PySide6.QtCore import QObject, Signal, QThreadPool from PySide6.QtCore import QObject, Signal, QThreadPool
from .file_watcher import FileWatcher, FileStabilityChecker, calculate_file_hash from .file_watcher import FileWatcher, FileStabilityChecker, calculate_file_hash
from .restore_worker import RestoreWorker from .restore_worker import RestoreWorker
from .. import __version__
from ..db.job_repository import JobRepository from ..db.job_repository import JobRepository
from ..db.event_repository import EventRepository from ..db.event_repository import EventRepository
from ..db.config_repository import ConfigRepository from ..db.config_repository import ConfigRepository
from ..constants import JobStatus, DEFAULT_CONFIG from ..constants import JobStatus, DEFAULT_CONFIG
from ..panel import panel_client
from ..utils.logger import app_logger from ..utils.logger import app_logger
@@ -41,7 +45,8 @@ class RestoreEngine(QObject):
# Configuración # Configuración
self._config = self._load_config() self._config = self._load_config()
self._report_instance_config_to_panel()
app_logger.info("RestoreEngine inicializado") app_logger.info("RestoreEngine inicializado")
def _load_config(self) -> dict: def _load_config(self) -> dict:
@@ -66,7 +71,8 @@ class RestoreEngine(QObject):
ConfigRepository.set("app_config", config) ConfigRepository.set("app_config", config)
self._config = config self._config = config
app_logger.info("Configuración guardada") app_logger.info("Configuración guardada")
self._report_instance_config_to_panel()
# Reconfigurar file watcher si está corriendo # Reconfigurar file watcher si está corriendo
if self._running: if self._running:
self._restart_file_watcher() self._restart_file_watcher()
@@ -103,7 +109,8 @@ class RestoreEngine(QObject):
self._running = True self._running = True
self._paused = False self._paused = False
self._report_instance_config_to_panel()
EventRepository.create("INFO", "Motor iniciado") EventRepository.create("INFO", "Motor iniciado")
app_logger.info("Motor iniciado") app_logger.info("Motor iniciado")
@@ -190,6 +197,30 @@ class RestoreEngine(QObject):
return stats 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: def _validate_config(self) -> bool:
"""Valida que la configuración sea correcta.""" """Valida que la configuración sea correcta."""
paths = self._config["paths"] paths = self._config["paths"]

View File

@@ -14,9 +14,19 @@ from ..db.event_repository import EventRepository
from ..db.node_repository import NodeRepository from ..db.node_repository import NodeRepository
from ..extract.seven_zip import SevenZipExtractor from ..extract.seven_zip import SevenZipExtractor
from ..sql.sql_manager import SQLServerManager from ..sql.sql_manager import SQLServerManager
from ..panel import panel_client
from ..transfer import sftp_copy
from ..utils.logger import app_logger 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): class RestoreWorkerSignals(QObject):
"""Señales para comunicación con la UI.""" """Señales para comunicación con la UI."""
job_started = Signal(str) # job_id job_started = Signal(str) # job_id
@@ -47,9 +57,20 @@ class RestoreWorker(QRunnable):
self.config = config self.config = config
self.dry_run = dry_run self.dry_run = dry_run
self.signals = RestoreWorkerSignals() self.signals = RestoreWorkerSignals()
self._extract_dir: Optional[str] = None self._extract_dir: Optional[str] = None
self._bak_path: 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): def run(self):
"""Ejecuta el procesamiento del job.""" """Ejecuta el procesamiento del job."""
@@ -58,46 +79,190 @@ class RestoreWorker(QRunnable):
try: try:
app_logger.info(f"Iniciando procesamiento de job {self.job_id}") app_logger.info(f"Iniciando procesamiento de job {self.job_id}")
self.signals.job_started.emit(self.job_id) self.signals.job_started.emit(self.job_id)
# Obtener job # Obtener job
job = JobRepository.get(self.job_id) job = JobRepository.get(self.job_id)
if not job: if not job:
raise ValueError(f"Job {self.job_id} no encontrado") raise ValueError(f"Job {self.job_id} no encontrado")
# Pipeline de procesamiento if not self._panel_configured():
self._process_node_mapping(job) 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._extract_backup(job)
self._restore_database(job) self._restore_database(job)
self._cleanup(job) self._cleanup(job)
# Actualizar tiempos # Actualizar tiempos
total_ms = int((time.time() - start_time) * 1000) total_ms = int((time.time() - start_time) * 1000)
JobRepository.update_timing(self.job_id, total_ms=total_ms) JobRepository.update_timing(self.job_id, total_ms=total_ms)
# Marcar como completado # Marcar como completado
JobRepository.update_status(self.job_id, JobStatus.COMPLETED) JobRepository.update_status(self.job_id, JobStatus.COMPLETED)
app_logger.info( app_logger.info(
f"Job {self.job_id} completado exitosamente en {total_ms}ms" 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) 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: except Exception as e:
error_msg = str(e) error_msg = str(e)
app_logger.error(f"Error procesando job {self.job_id}: {error_msg}", exc_info=True) app_logger.error(f"Error procesando job {self.job_id}: {error_msg}", exc_info=True)
JobRepository.update_status( JobRepository.update_status(
self.job_id, self.job_id,
JobStatus.FAILED, JobStatus.FAILED,
error=error_msg, error=error_msg,
increment_attempts=True increment_attempts=True
) )
EventRepository.create("ERROR", f"Job falló: {error_msg}", self.job_id) 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.error_occurred.emit(self.job_id, error_msg)
self.signals.job_completed.emit(self.job_id, False) 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): def _process_node_mapping(self, job):
"""Procesa el mapeo de nodo a base de datos.""" """Procesa el mapeo de nodo a base de datos."""
step_id = JobStepRepository.create(self.job_id, StepType.NODE_MAPPING) 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.""" """Restaura la base de datos desde el backup."""
JobRepository.update_status(self.job_id, JobStatus.RESTORING) JobRepository.update_status(self.job_id, JobStatus.RESTORING)
self.signals.job_progress.emit(self.job_id, JobStatus.RESTORING) self.signals.job_progress.emit(self.job_id, JobStatus.RESTORING)
# Recargar job para obtener db_name actualizado # Recargar job para obtener db_name actualizado
job = JobRepository.get(self.job_id) job = JobRepository.get(self.job_id)
if not job.db_name: if not job.db_name:
raise ValueError("DB name no está configurado en el job") 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 # Conectar a SQL
connect_step_id = JobStepRepository.create(self.job_id, StepType.SQL_CONNECT) connect_step_id = JobStepRepository.create(self.job_id, StepType.SQL_CONNECT)
try: try:
sql_config = self.config["sql"]
sql_manager = SQLServerManager( sql_manager = SQLServerManager(
server=sql_config["server"], server=server,
use_windows_auth=sql_config["use_windows_auth"], use_windows_auth=use_windows_auth,
username=sql_config.get("username"), username=username,
password=sql_config.get("password") password=password
) )
# Test connection # Test connection
success, error = sql_manager.test_connection() success, error = sql_manager.test_connection()
if not success: if not success:
raise RuntimeError(f"Conexión SQL falló: {error}") raise RuntimeError(f"Conexión SQL falló: {error}")
JobStepRepository.complete(connect_step_id, exit_code=0) JobStepRepository.complete(connect_step_id, exit_code=0)
except Exception as e: except Exception as e:
JobStepRepository.complete(connect_step_id, exit_code=1, error=str(e)) JobStepRepository.complete(connect_step_id, exit_code=1, error=str(e))
raise raise
sql_backup_path = self._bak_path
# Obtener FILELISTONLY # Obtener FILELISTONLY
filelist_step_id = JobStepRepository.create(self.job_id, StepType.FILELIST) filelist_step_id = JobStepRepository.create(self.job_id, StepType.FILELIST)
filelist_start = time.time() filelist_start = time.time()
try: try:
logical_files, error = sql_manager.get_filelist_from_backup( logical_files, error = sql_manager.get_filelist_from_backup(
self._bak_path, sql_backup_path,
timeout_minutes=self.config["timeouts"]["restore_minutes"] timeout_minutes=self.config["timeouts"]["restore_minutes"]
) )
filelist_ms = int((time.time() - filelist_start) * 1000) filelist_ms = int((time.time() - filelist_start) * 1000)
JobRepository.update_timing(self.job_id, filelist_ms=filelist_ms) JobRepository.update_timing(self.job_id, filelist_ms=filelist_ms)
if error: if error:
raise RuntimeError(error) raise RuntimeError(error)
files_info = ", ".join([f"{lf.logical_name}({lf.type})" for lf in logical_files]) files_info = ", ".join([f"{lf.logical_name}({lf.type})" for lf in logical_files])
JobStepRepository.complete( JobStepRepository.complete(
filelist_step_id, filelist_step_id,
exit_code=0, exit_code=0,
stdout=files_info[:1000] stdout=files_info[:1000]
) )
except Exception as e: except Exception as e:
JobStepRepository.complete(filelist_step_id, exit_code=1, error=str(e)) JobStepRepository.complete(filelist_step_id, exit_code=1, error=str(e))
raise raise
# RESTORE DATABASE # RESTORE DATABASE
restore_step_id = JobStepRepository.create(self.job_id, StepType.RESTORE) restore_step_id = JobStepRepository.create(self.job_id, StepType.RESTORE)
restore_start = time.time() restore_start = time.time()
try: try:
data_folder = self.config["paths"]["data_sql_folder"]
success, stdout, error = sql_manager.restore_database( success, stdout, error = sql_manager.restore_database(
db_name=job.db_name, db_name=job.db_name,
backup_path=self._bak_path, backup_path=sql_backup_path,
data_folder=data_folder, data_folder=data_folder,
logical_files=logical_files, logical_files=logical_files,
timeout_minutes=self.config["timeouts"]["restore_minutes"], timeout_minutes=self.config["timeouts"]["restore_minutes"],
dry_run=self.dry_run dry_run=self.dry_run
) )
restore_ms = int((time.time() - restore_start) * 1000) restore_ms = int((time.time() - restore_start) * 1000)
JobRepository.update_timing(self.job_id, restore_ms=restore_ms) JobRepository.update_timing(self.job_id, restore_ms=restore_ms)
if not success: if not success:
raise RuntimeError(error or "RESTORE falló sin mensaje de error") raise RuntimeError(error or "RESTORE falló sin mensaje de error")
JobStepRepository.complete( JobStepRepository.complete(
restore_step_id, restore_step_id,
exit_code=0, exit_code=0,
stdout=stdout[:1000] if stdout else None stdout=stdout[:1000] if stdout else None
) )
except Exception as e: except Exception as e:
JobStepRepository.complete(restore_step_id, exit_code=1, error=str(e)) JobStepRepository.complete(restore_step_id, exit_code=1, error=str(e))
raise raise
def _cleanup(self, job): def _cleanup(self, job):
"""Limpia archivos temporales y mueve el ZIP.""" """Limpia archivos temporales y mueve el ZIP."""
JobRepository.update_status(self.job_id, JobStatus.CLEANING) JobRepository.update_status(self.job_id, JobStatus.CLEANING)

1
app/panel/__init__.py Normal file
View File

@@ -0,0 +1 @@
"""Integración con PANEL_BASES_ANEXO24 (servidor de restauración activo)."""

427
app/panel/panel_client.py Normal file
View File

@@ -0,0 +1,427 @@
"""
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.
"""
import os
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)
def _http_verify(verify_ssl: Optional[bool] = None) -> bool:
"""Verificación TLS para requests. Por defecto False (panel dev con cert propio)."""
if verify_ssl is not None:
return verify_ssl
env = os.getenv("CLOUDRESTORE_PANEL_VERIFY_SSL", "").strip().lower()
if env in ("1", "true", "yes"):
return True
if env in ("0", "false", "no"):
return False
return False
# 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=<db_name>: 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, verify=_http_verify()
)
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, verify=_http_verify()
)
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, verify=_http_verify()
)
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, verify=_http_verify()
)
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, verify=_http_verify()
)
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, verify=_http_verify()
)
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}"

View File

@@ -34,7 +34,8 @@ class SQLServerManager:
username: Usuario SQL (si no usa Windows Auth) username: Usuario SQL (si no usa Windows Auth)
password: Contraseña 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.use_windows_auth = use_windows_auth
self.username = username self.username = username
self.password = password self.password = password
@@ -65,6 +66,93 @@ class SQLServerManager:
return ";".join(parts) return ";".join(parts)
@staticmethod
def _drain_cursor(cursor) -> None:
"""Consume todos los result sets de comandos largos (RESTORE, etc.)."""
while True:
try:
if cursor.description:
cursor.fetchall()
except Exception:
pass
if not cursor.nextset():
break
def wait_for_database_state(
self,
db_name: str,
target_state: str = "ONLINE",
timeout_seconds: int = 300,
poll_seconds: float = 2.0,
) -> Optional[str]:
"""Espera hasta que la BD alcance target_state o agote el timeout."""
deadline = time.time() + timeout_seconds
last_state: Optional[str] = None
while time.time() < deadline:
last_state = self.get_database_state(db_name)
if last_state == target_state:
return last_state
if last_state is None and target_state == "ONLINE":
# Puede tardar en aparecer en sys.databases al inicio del RESTORE.
pass
time.sleep(poll_seconds)
return last_state
def get_database_state(self, db_name: str) -> Optional[str]:
"""Devuelve state_desc de sys.databases o None si no existe."""
try:
conn_str = self.get_connection_string()
conn = pyodbc.connect(conn_str, timeout=10)
cursor = conn.cursor()
cursor.execute(
"SELECT state_desc FROM sys.databases WHERE name = ?",
(db_name,),
)
row = cursor.fetchone()
conn.close()
return str(row[0]) if row else None
except Exception as e:
app_logger.error(f"Error consultando estado de DB [{db_name}]: {e}")
return None
def _prepare_database_for_restore(self, cursor, db_name: str) -> None:
"""
Limpia una BD atascada en RESTORING u offline antes de un RESTORE nuevo.
"""
state = self.get_database_state(db_name)
if not state:
return
if state == "RESTORING":
app_logger.warning(
f"BD [{db_name}] en RESTORING; intentando WITH RECOVERY..."
)
try:
cursor.execute(f"RESTORE DATABASE [{db_name}] WITH RECOVERY")
self._drain_cursor(cursor)
state = self.get_database_state(db_name)
except Exception as e:
app_logger.warning(f"RECOVERY falló para [{db_name}]: {e}")
if state == "RESTORING":
app_logger.warning(
f"BD [{db_name}] sigue en RESTORING; eliminando con DROP DATABASE..."
)
cursor.execute(f"DROP DATABASE [{db_name}]")
self._drain_cursor(cursor)
return
if state == "ONLINE":
try:
cursor.execute(
f"ALTER DATABASE [{db_name}] SET SINGLE_USER WITH ROLLBACK IMMEDIATE"
)
app_logger.info(f"BD [{db_name}] configurada en SINGLE_USER")
except Exception as e:
app_logger.warning(
f"No se pudo configurar SINGLE_USER en [{db_name}]: {e}"
)
def test_connection(self) -> Tuple[bool, Optional[str]]: def test_connection(self) -> Tuple[bool, Optional[str]]:
""" """
Prueba la conexión a SQL Server. Prueba la conexión a SQL Server.
@@ -168,85 +256,134 @@ class SQLServerManager:
Returns: Returns:
Tupla (éxito, stdout, error) Tupla (éxito, stdout, error)
""" """
conn = None
try: try:
# Construir las cláusulas MOVE # Construir las cláusulas MOVE con un destino único por archivo lógico.
move_clauses = [] # Renombrar todos los data files a {db}.mdf colisiona si el backup tiene
for lf in logical_files: # varios archivos; aquí cada archivo recibe un nombre distinto (G4).
if lf.type == 'D': # Data file move_clauses = self._build_move_clauses(db_name, data_folder, logical_files)
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}'")
if not move_clauses: if not move_clauses:
return False, None, "No se pudieron determinar los archivos de datos y log" return False, None, "No se pudieron determinar los archivos del backup"
# 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;
-- Restaurar restore_query = (
RESTORE DATABASE [{db_name}] f"RESTORE DATABASE [{db_name}] FROM DISK = N'{backup_path}' "
FROM DISK = N'{backup_path}' f"WITH {', '.join(move_clauses)}, REPLACE"
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: if dry_run:
app_logger.info("Modo DRY RUN: No se ejecutará el RESTORE") app_logger.info("Modo DRY RUN: No se ejecutará el RESTORE")
return True, restore_cmd, None return True, restore_query, None
# Ejecutar RESTORE # Ejecutar RESTORE
conn_str = self.get_connection_string() conn_str = self.get_connection_string()
conn = pyodbc.connect(conn_str, timeout=timeout_minutes * 60) conn = pyodbc.connect(conn_str, timeout=timeout_minutes * 60)
conn.autocommit = True # Necesario para ALTER DATABASE conn.autocommit = True # Necesario para ALTER DATABASE
cursor = conn.cursor() cursor = conn.cursor()
start_time = time.time() start_time = time.time()
app_logger.info(f"Ejecutando RESTORE DATABASE [{db_name}]...") app_logger.info(f"Ejecutando RESTORE DATABASE [{db_name}]...")
# Ejecutar en múltiples pasos
output_lines = [] output_lines = []
restore_error: Optional[str] = None
# 1. Single user
try: try:
cursor.execute(f"ALTER DATABASE [{db_name}] SET SINGLE_USER WITH ROLLBACK IMMEDIATE") # 1. Limpiar BD atascada o poner SINGLE_USER si ya existe ONLINE
output_lines.append("Base de datos configurada en modo SINGLE_USER") self._prepare_database_for_restore(cursor, db_name)
# 2. RESTORE (RECOVERY es el default; drenar result sets hasta completar)
cursor.execute(restore_query)
self._drain_cursor(cursor)
output_lines.append("RESTORE DATABASE ejecutado")
# Esperar a que SQL Server termine (evita borrar el .bak demasiado pronto)
waited_state = self.wait_for_database_state(
db_name,
target_state="ONLINE",
timeout_seconds=timeout_minutes * 60,
)
if waited_state != "ONLINE":
restore_error = (
f"BD [{db_name}] no alcanzó ONLINE tras RESTORE "
f"(estado: {waited_state})"
)
app_logger.error(restore_error)
except Exception as e: except Exception as e:
app_logger.warning(f"Error configurando SINGLE_USER (puede no existir la DB): {e}") restore_error = str(e)
app_logger.error(f"RESTORE DATABASE falló para [{db_name}]: {e}")
# 2. RESTORE finally:
restore_query = f"RESTORE DATABASE [{db_name}] FROM DISK = N'{backup_path}' WITH {', '.join(move_clauses)}, REPLACE" # 3. MULTI_USER solo si la BD quedó ONLINE (G7)
cursor.execute(restore_query) final_state = self.get_database_state(db_name)
output_lines.append("RESTORE DATABASE completado") if final_state == "ONLINE":
try:
# 3. Multi user cursor.execute(f"ALTER DATABASE [{db_name}] SET MULTI_USER")
cursor.execute(f"ALTER DATABASE [{db_name}] SET MULTI_USER") output_lines.append("Base de datos configurada en modo MULTI_USER")
output_lines.append("Base de datos configurada en modo MULTI_USER") except Exception as e:
restore_error = restore_error or (
f"No se pudo volver a MULTI_USER en [{db_name}]: {e}"
)
app_logger.error(restore_error)
elif final_state:
msg = (
f"BD [{db_name}] quedó en estado {final_state} tras RESTORE"
)
restore_error = restore_error or msg
app_logger.error(msg)
elapsed = time.time() - start_time elapsed = time.time() - start_time
output_lines.append(f"Restauración completada en {elapsed:.2f}s") output_lines.append(f"Restauración finalizada en {elapsed:.2f}s")
conn.close()
output = "\n".join(output_lines) output = "\n".join(output_lines)
app_logger.info(output) app_logger.info(output)
if restore_error or self.get_database_state(db_name) != "ONLINE":
state = self.get_database_state(db_name)
err = restore_error or f"BD [{db_name}] no quedó ONLINE (estado: {state})"
return False, output, err
return True, output, None return True, output, None
except Exception as e: except Exception as e:
error_msg = f"Error restaurando base de datos: {str(e)}" error_msg = f"Error restaurando base de datos: {str(e)}"
app_logger.error(error_msg) app_logger.error(error_msg)
return False, None, 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: def database_exists(self, db_name: str) -> bool:
""" """

1
app/transfer/__init__.py Normal file
View File

@@ -0,0 +1 @@
"""Transferencia de archivos hacia los servidores SQL remotos (SMB)."""

160
app/transfer/sftp_copy.py Normal file
View File

@@ -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()

View File

@@ -3,13 +3,14 @@
from PySide6.QtWidgets import ( from PySide6.QtWidgets import (
QWidget, QVBoxLayout, QHBoxLayout, QGroupBox, QWidget, QVBoxLayout, QHBoxLayout, QGroupBox,
QFormLayout, QLineEdit, QPushButton, QSpinBox, QFormLayout, QLineEdit, QPushButton, QSpinBox,
QCheckBox, QFileDialog, QMessageBox, QScrollArea QCheckBox, QFileDialog, QMessageBox, QScrollArea, QComboBox
) )
from pathlib import Path from pathlib import Path
from ..utils.crypto import encrypt_password, decrypt_password from ..utils.crypto import encrypt_password, decrypt_password
from ..extract.seven_zip import SevenZipExtractor from ..extract.seven_zip import SevenZipExtractor
from ..sql.sql_manager import SQLServerManager from ..sql.sql_manager import SQLServerManager
from ..panel import panel_client
class ConfigTab(QWidget): class ConfigTab(QWidget):
@@ -20,6 +21,7 @@ class ConfigTab(QWidget):
super().__init__() super().__init__()
self.engine = engine self.engine = engine
self.config = {} self.config = {}
self._catalog_names: list[str] = []
self._setup_ui() self._setup_ui()
def _setup_ui(self): def _setup_ui(self):
@@ -118,7 +120,38 @@ class ConfigTab(QWidget):
sql_group.setLayout(sql_layout) sql_group.setLayout(sql_layout)
layout.addWidget(sql_group) 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 # Grupo: Concurrencia
concurrency_group = QGroupBox("Concurrencia") concurrency_group = QGroupBox("Concurrencia")
concurrency_layout = QFormLayout() concurrency_layout = QFormLayout()
@@ -269,7 +302,42 @@ class ConfigTab(QWidget):
self.dry_run_checkbox.setChecked(features.get("dry_run_mode", False)) self.dry_run_checkbox.setChecked(features.get("dry_run_mode", False))
self.auto_scan_checkbox.setChecked(features.get("auto_scan_enabled", True)) self.auto_scan_checkbox.setChecked(features.get("auto_scan_enabled", True))
self.scan_interval_spinbox.setValue(features.get("scan_interval_seconds", 30)) 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): def _save_config(self):
"""Guarda la configuración.""" """Guarda la configuración."""
# Construir config # Construir config
@@ -307,9 +375,31 @@ class ConfigTab(QWidget):
"dry_run_mode": self.dry_run_checkbox.isChecked(), "dry_run_mode": self.dry_run_checkbox.isChecked(),
"auto_scan_enabled": self.auto_scan_checkbox.isChecked(), "auto_scan_enabled": self.auto_scan_checkbox.isChecked(),
"scan_interval_seconds": self.scan_interval_spinbox.value() "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 # Cifrar password si no usa Windows Auth
if not self.sql_windows_auth_checkbox.isChecked(): if not self.sql_windows_auth_checkbox.isChecked():
password = self.sql_password_input.text() password = self.sql_password_input.text()
@@ -356,6 +446,25 @@ class ConfigTab(QWidget):
"No se pudo auto-detectar 7-Zip. Por favor, selecciona manualmente." "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): def _on_auth_changed(self, checked: bool):
"""Maneja el cambio en el tipo de autenticación.""" """Maneja el cambio en el tipo de autenticación."""
self.sql_username_input.setEnabled(not checked) self.sql_username_input.setEnabled(not checked)

View File

@@ -4,7 +4,7 @@ from PySide6.QtWidgets import (
QWidget, QVBoxLayout, QHBoxLayout, QTableWidget, QWidget, QVBoxLayout, QHBoxLayout, QTableWidget,
QTableWidgetItem, QPushButton, QDialog, QFormLayout, QTableWidgetItem, QPushButton, QDialog, QFormLayout,
QLineEdit, QTextEdit, QCheckBox, QDialogButtonBox, QLineEdit, QTextEdit, QCheckBox, QDialogButtonBox,
QMessageBox QMessageBox, QLabel
) )
from PySide6.QtCore import Qt from PySide6.QtCore import Qt

View File

@@ -9,5 +9,14 @@ pyodbc>=5.0.0
# pywin32 para DPAPI (cifrado de passwords) # pywin32 para DPAPI (cifrado de passwords)
pywin32>=306 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) # Para empaquetado (opcional)
pyinstaller>=6.0.0 pyinstaller>=6.0.0
# Pruebas (opcional)
pytest>=8.0.0

5
tests/conftest.py Normal file
View File

@@ -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__))))

View File

@@ -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

381
tests/test_panel_client.py Normal file
View File

@@ -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

138
tests/test_route_forward.py Normal file
View File

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

135
tests/test_sftp_copy.py Normal file
View File

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

View File

@@ -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