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
```
### 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)
1. **Tab Configuración**:

View File

@@ -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,16 @@ 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": "",
# False en dev: el panel local usa HTTPS con certificado propio.
"verify_ssl": False,
}
}

View File

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

View File

@@ -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,6 +45,7 @@ class RestoreEngine(QObject):
# Configuración
self._config = self._load_config()
self._report_instance_config_to_panel()
app_logger.info("RestoreEngine inicializado")
@@ -66,6 +71,7 @@ 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:
@@ -103,6 +109,7 @@ 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"]

View File

@@ -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
@@ -50,6 +60,17 @@ class RestoreWorker(QRunnable):
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."""
@@ -64,8 +85,26 @@ class RestoreWorker(QRunnable):
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)
@@ -80,8 +119,19 @@ class RestoreWorker(QRunnable):
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)
@@ -95,9 +145,124 @@ class RestoreWorker(QRunnable):
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)
@@ -199,16 +364,30 @@ class RestoreWorker(QRunnable):
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
@@ -222,13 +401,15 @@ class RestoreWorker(QRunnable):
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"]
)
@@ -254,11 +435,9 @@ class RestoreWorker(QRunnable):
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"],

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)
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
@@ -65,6 +66,93 @@ class SQLServerManager:
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]]:
"""
Prueba la conexión a SQL Server.
@@ -168,42 +256,25 @@ 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"
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;
restore_query = (
f"RESTORE DATABASE [{db_name}] FROM DISK = N'{backup_path}' "
f"WITH {', '.join(move_clauses)}, REPLACE"
)
-- Restaurar
RESTORE DATABASE [{db_name}]
FROM DISK = N'{backup_path}'
WITH {', '.join(move_clauses)}, REPLACE;
-- Volver a modo multi user
ALTER DATABASE [{db_name}] SET MULTI_USER;
"""
app_logger.info(f"Comando RESTORE generado:\n{restore_cmd}")
app_logger.info(f"Comando RESTORE generado:\n{restore_query}")
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()
@@ -214,39 +285,105 @@ ALTER DATABASE [{db_name}] SET MULTI_USER;
start_time = time.time()
app_logger.info(f"Ejecutando RESTORE DATABASE [{db_name}]...")
# Ejecutar en múltiples pasos
output_lines = []
restore_error: Optional[str] = None
# 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")
# 1. Limpiar BD atascada o poner SINGLE_USER si ya existe ONLINE
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:
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")
restore_error = str(e)
app_logger.error(f"RESTORE DATABASE falló para [{db_name}]: {e}")
finally:
# 3. MULTI_USER solo si la BD quedó ONLINE (G7)
final_state = self.get_database_state(db_name)
if final_state == "ONLINE":
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:
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
output_lines.append(f"Restauración completada en {elapsed:.2f}s")
conn.close()
output_lines.append(f"Restauración finalizada en {elapsed:.2f}s")
output = "\n".join(output_lines)
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
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:
"""

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 (
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):
@@ -119,6 +121,37 @@ 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()
@@ -270,6 +303,41 @@ class ConfigTab(QWidget):
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)

View File

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

View File

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

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