Compare commits
2 Commits
854ffa116f
...
072be5b5db
| Author | SHA1 | Date | |
|---|---|---|---|
| 072be5b5db | |||
| afb76ce4a6 |
197
INTEGRACION_PANEL.md
Normal file
197
INTEGRACION_PANEL.md
Normal 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**.
|
||||
@@ -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**:
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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."""
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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
|
||||
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
1
app/panel/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
"""Integración con PANEL_BASES_ANEXO24 (servidor de restauración activo)."""
|
||||
427
app/panel/panel_client.py
Normal file
427
app/panel/panel_client.py
Normal 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}"
|
||||
@@ -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")
|
||||
except Exception as e:
|
||||
app_logger.warning(f"Error configurando SINGLE_USER (puede no existir la DB): {e}")
|
||||
# 1. Limpiar BD atascada o poner SINGLE_USER si ya existe ONLINE
|
||||
self._prepare_database_for_restore(cursor, db_name)
|
||||
|
||||
# 2. RESTORE
|
||||
restore_query = f"RESTORE DATABASE [{db_name}] FROM DISK = N'{backup_path}' WITH {', '.join(move_clauses)}, REPLACE"
|
||||
# 2. RESTORE (RECOVERY es el default; drenar result sets hasta completar)
|
||||
cursor.execute(restore_query)
|
||||
output_lines.append("RESTORE DATABASE completado")
|
||||
self._drain_cursor(cursor)
|
||||
output_lines.append("RESTORE DATABASE ejecutado")
|
||||
|
||||
# 3. Multi user
|
||||
# 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:
|
||||
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
1
app/transfer/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
"""Transferencia de archivos hacia los servidores SQL remotos (SMB)."""
|
||||
160
app/transfer/sftp_copy.py
Normal file
160
app/transfer/sftp_copy.py
Normal 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()
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
5
tests/conftest.py
Normal 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__))))
|
||||
41
tests/test_colocated_resolve.py
Normal file
41
tests/test_colocated_resolve.py
Normal 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
381
tests/test_panel_client.py
Normal 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
138
tests/test_route_forward.py
Normal 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
135
tests/test_sftp_copy.py
Normal 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, " ")
|
||||
77
tests/test_sql_manager_move.py
Normal file
77
tests/test_sql_manager_move.py
Normal 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
|
||||
Reference in New Issue
Block a user