feature/generador-instaladores-linux-windows

This commit is contained in:
2026-07-30 07:34:17 -06:00
parent cafe3f1b87
commit c3f1d70e23
34 changed files with 3498 additions and 252 deletions

View File

@@ -2,18 +2,21 @@
import platform
import socket
from threading import Thread
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 .maintenance_scheduler import DailyMaintenanceScheduler
from .restore_worker import RestoreWorker
from .retention import RetentionCleaner
from .. import __version__
from ..db.job_repository import JobRepository
from ..db.event_repository import EventRepository
from ..db.config_repository import ConfigRepository
from ..config.env_loader import apply_env_overrides, load_env_file
from ..constants import JobStatus, DEFAULT_CONFIG
from ..constants import APP_ARCH, APP_DIR, APP_PLATFORM, JobStatus, DEFAULT_CONFIG
from ..panel import panel_client
from ..utils.logger import app_logger
@@ -39,7 +42,10 @@ class RestoreEngine(QObject):
# File watcher
self._file_watcher: Optional[FileWatcher] = None
# Mantenimiento diario (retención de respaldos obsoletos)
self._maintenance: Optional[DailyMaintenanceScheduler] = None
# Estado
self._running = False
self._paused = False
@@ -75,9 +81,10 @@ class RestoreEngine(QObject):
app_logger.info("Configuración guardada")
self._report_instance_config_to_panel()
# Reconfigurar file watcher si está corriendo
# Reconfigurar file watcher y mantenimiento si está corriendo
if self._running:
self._restart_file_watcher()
self._restart_maintenance()
def get_config(self) -> dict:
"""Obtiene la configuración actual."""
@@ -93,7 +100,12 @@ class RestoreEngine(QObject):
if not self._validate_config():
app_logger.error("Configuración inválida, no se puede iniciar el motor")
return
# Barrido de Temp: al arrancar no hay jobs corriendo, así que cualquier subcarpeta en
# extract_folder es un remanente huérfano de una corrida previa (los caminos de
# fallo/diferido no siempre limpiaban). Se elimina para que Temp no se acumule.
self._purge_stale_temp()
# Configurar thread pool
extract_workers = self._config["concurrency"]["extract_workers"]
restore_workers = self._config["concurrency"]["restore_workers"]
@@ -108,7 +120,10 @@ class RestoreEngine(QObject):
# Iniciar file watcher
if self._config["features"]["auto_scan_enabled"]:
self._start_file_watcher()
# Iniciar el mantenimiento diario (limpieza de respaldos obsoletos)
self._start_maintenance()
self._running = True
self._paused = False
self._report_instance_config_to_panel()
@@ -116,6 +131,28 @@ class RestoreEngine(QObject):
EventRepository.create("INFO", "Motor iniciado")
app_logger.info("Motor iniciado")
def _purge_stale_temp(self):
"""Elimina subcarpetas huérfanas en la carpeta temporal de extracción (al arrancar)."""
import shutil
extract_folder = (self._config.get("paths") or {}).get("extract_folder")
if not extract_folder:
return
base = Path(extract_folder)
if not base.is_dir():
return
removed = 0
for child in base.iterdir():
if not child.is_dir():
continue
try:
shutil.rmtree(child, ignore_errors=True)
removed += 1
except Exception as e:
app_logger.warning(f"No se pudo limpiar Temp huérfano {child}: {e}")
if removed:
app_logger.info(f"Temp: {removed} carpeta(s) huérfana(s) eliminada(s) al iniciar")
def stop(self):
"""Detiene el motor."""
if not self._running:
@@ -125,7 +162,10 @@ class RestoreEngine(QObject):
if self._file_watcher:
self._file_watcher.stop()
self._file_watcher = None
# Detener el mantenimiento diario
self._stop_maintenance()
# Esperar a que terminen los workers
self._thread_pool.waitForDone(msecs=30000) # 30s timeout
@@ -207,28 +247,45 @@ class RestoreEngine(QObject):
if not api_url or not api_token:
return
input_folder = (self._config.get("paths") or {}).get("input_folder") or ""
paths = self._config.get("paths") or {}
input_folder = paths.get("input_folder") or ""
processed_folder = paths.get("processed_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
# platform/arch le dicen al PANEL qué artefacto le toca a este servidor cuando
# instala o actualiza (a24c.cras_releases se llavea por version+platform+arch).
panel_client.report_instance_config(
api_url=api_url,
api_token=api_token,
input_folder=input_folder,
processed_folder=processed_folder,
host_name=host_name,
app_version=__version__,
instance_key=instance_key,
platform_name=APP_PLATFORM,
arch=APP_ARCH,
# Dónde vive el ejecutable. El PANEL actualiza en esta ruta; si instalara en la
# default crearía una segunda instalación y dejaría huérfano este config/.env.
install_path=str(APP_DIR),
)
def _validate_config(self) -> bool:
"""Valida que la configuración sea correcta."""
paths = self._config["paths"]
# Validar carpetas requeridas
required_paths = ["input_folder", "extract_folder", "data_sql_folder"]
# Validar carpetas requeridas. processed_folder/failed_folder se incluyen para que la
# reubicación de ZIP (éxito/fallo) nunca falle por carpeta inexistente.
required_paths = [
"input_folder",
"extract_folder",
"data_sql_folder",
"processed_folder",
"failed_folder",
]
for key in required_paths:
if not paths.get(key):
app_logger.error(f"Falta configurar: {key}")
@@ -280,6 +337,52 @@ class RestoreEngine(QObject):
"""Reinicia el file watcher con nueva configuración."""
if self._running and not self._paused:
self._start_file_watcher()
def _start_maintenance(self):
"""Inicia el mantenimiento diario si la retención está habilitada."""
retention_cfg = self._config.get("retention", {})
if not retention_cfg.get("enabled", True):
app_logger.info("Retención deshabilitada; no se inicia el mantenimiento diario")
return
self._maintenance = DailyMaintenanceScheduler(
task=self._run_retention,
check_interval_seconds=retention_cfg.get("check_interval_seconds", 3600),
run_at_hour=retention_cfg.get("run_at_hour", 3),
)
self._maintenance.start()
def _stop_maintenance(self):
"""Detiene el mantenimiento diario."""
if self._maintenance:
self._maintenance.stop()
self._maintenance = None
def _restart_maintenance(self):
"""Reinicia el mantenimiento diario con nueva configuración."""
if self._running:
self._stop_maintenance()
self._start_maintenance()
def _run_retention(self):
"""Ejecuta una corrida de retención (tarea del mantenimiento diario)."""
try:
RetentionCleaner(self._config).run()
except Exception as e:
app_logger.error(f"Error en la retención: {e}", exc_info=True)
EventRepository.create("ERROR", f"Error en la retención: {e}")
def trigger_maintenance_now(self):
"""Corre la retención de inmediato en segundo plano (acción manual de la UI).
Respeta el dry_run de la configuración y no consume el turno del día calendario.
"""
def _run():
if self._maintenance:
self._maintenance.trigger_now()
else:
self._run_retention()
Thread(target=_run, daemon=True).start()
def _on_file_ready(self, file_path: str):
"""
@@ -296,10 +399,18 @@ class RestoreEngine(QObject):
# Calcular hash para evitar duplicados
file_hash = calculate_file_hash(file_path)
if JobRepository.exists_by_hash(file_hash):
app_logger.warning(f"Archivo ya procesado (hash duplicado): {file_path}")
if JobRepository.has_blocking_job_by_hash(file_hash):
app_logger.warning(f"Archivo ya procesado o en curso (hash duplicado): {file_path}")
return
# Limpia intentos FALLIDOS previos con este hash para permitir un reintento fresco
# (antes un FAILED transitorio bloqueaba el reproceso de forma permanente).
removed = JobRepository.delete_failed_by_hash(file_hash)
if removed:
app_logger.info(
f"Reintento de {Path(file_path).name}: {removed} job(s) fallido(s) previo(s) eliminado(s)"
)
# Crear job
file_name = Path(file_path).name
job_id = JobRepository.create(file_path, file_name, file_hash)

View File

@@ -0,0 +1,125 @@
"""Programador de mantenimiento diario in-process.
Corre una tarea (p.ej. la retención) UNA VEZ por día calendario, dentro del propio proceso de
CloudRestoreAS. No depende de cron/systemd externos (el despliegue es embedded-only) ni del
event loop de Qt: usa un hilo daemon con el mismo patrón que ``FileWatcher``.
La marca de la última corrida se persiste en la tabla ``config`` (vía ``ConfigRepository``), de
modo que:
- corre a lo más una vez por día calendario ("claim" de la fecha ANTES de ejecutar), y
- si el servicio estuvo caído se "pone al día" en el primer arranque de un día nuevo.
"""
from datetime import datetime
from threading import Event, Thread
from typing import Callable, Optional
from ..db.config_repository import ConfigRepository
from ..db.event_repository import EventRepository
from ..utils.logger import app_logger
class DailyMaintenanceScheduler:
"""Ejecuta ``task`` una vez al día en un hilo daemon."""
def __init__(
self,
task: Callable[[], None],
*,
check_interval_seconds: int = 3600,
run_at_hour: Optional[int] = 3,
state_key: str = "retention_last_run",
clock: Callable[[], datetime] = datetime.now,
) -> None:
self._task = task
self._check_interval = max(60, int(check_interval_seconds))
self._run_at_hour = run_at_hour
self._state_key = state_key
self._clock = clock
self._stop_event = Event()
self._thread: Optional[Thread] = None
def start(self) -> None:
"""Inicia el hilo de mantenimiento (hace un chequeo inmediato de 'catch-up')."""
if self._thread and self._thread.is_alive():
app_logger.warning("DailyMaintenanceScheduler ya está corriendo")
return
self._stop_event.clear()
self._thread = Thread(target=self._run, daemon=True)
self._thread.start()
app_logger.info(
f"Mantenimiento diario iniciado (hora={self._run_at_hour}, "
f"cada {self._check_interval}s)"
)
def stop(self, timeout: float = 30.0) -> None:
"""Detiene el hilo (timeout amplio: la limpieza puede tardar)."""
if self._thread:
self._stop_event.set()
self._thread.join(timeout=timeout)
app_logger.info("Mantenimiento diario detenido")
def trigger_now(self) -> None:
"""Corre la tarea de inmediato sin importar la fecha (uso manual/validación)."""
self._execute(force=True)
# -- Interno ---------------------------------------------------------------------
def _run(self) -> None:
while not self._stop_event.is_set():
try:
self._run_if_due()
except Exception as e:
app_logger.error(
f"Error en DailyMaintenanceScheduler: {e}", exc_info=True
)
self._stop_event.wait(self._check_interval)
def _run_if_due(self) -> None:
now = self._clock()
if self._is_due(now, self._load_state()):
self._execute(now=now)
def _is_due(self, now: datetime, state: dict) -> bool:
if state.get("last_run_date") == now.date().isoformat():
return False # ya corrió hoy
if self._run_at_hour is None:
return True # primera oportunidad de un día nuevo
return now.hour >= self._run_at_hour
def _execute(self, now: Optional[datetime] = None, force: bool = False) -> None:
now = now or self._clock()
today = now.date().isoformat()
# Claim al inicio: marca la fecha ANTES de correr para garantizar "máximo 1/día"
# aunque la corrida falle o el proceso muera a mitad (la tarea es idempotente).
if not force:
self._save_state(today, now, "running")
try:
self._task()
status = "ok"
except Exception as e:
app_logger.error(
f"Fallo en la tarea de mantenimiento diaria: {e}", exc_info=True
)
EventRepository.create("ERROR", f"Fallo en la limpieza diaria: {e}")
status = f"error: {e}"
if not force:
self._save_state(today, now, status)
def _load_state(self) -> dict:
state = ConfigRepository.get(self._state_key, {})
return state if isinstance(state, dict) else {}
def _save_state(self, run_date: str, now: datetime, status: str) -> None:
ConfigRepository.set(
self._state_key,
{
"last_run_date": run_date,
"last_run_at": now.isoformat(),
"last_status": status,
},
)

View File

@@ -146,11 +146,23 @@ class RestoreWorker(QRunnable):
EventRepository.create("ERROR", f"Job falló: {error_msg}", self.job_id)
job = JobRepository.get(self.job_id)
# Mover el ZIP fallido a Fallados/<fecha>/ para diagnóstico y descarga desde el panel.
try:
if job:
self._move_zip_to_failed(job)
except Exception as move_err:
app_logger.warning(f"No se pudo mover ZIP a Fallados: {move_err}")
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)
finally:
# Garantiza la limpieza del Temp en TODOS los caminos (éxito, fallo, diferido,
# forward). En éxito _cleanup ya lo eliminó; aquí es red de seguridad.
self._purge_extract_dir()
def _resolve_route(self, job) -> dict:
"""
Enrutamiento automático vía panel: restore_local o forward según nodo/asignación.
@@ -185,12 +197,41 @@ class RestoreWorker(QRunnable):
return route
def _collect_zip_paths(self, source_path: str) -> list[str]:
"""Rutas locales del ZIP (incluye todas las partes multipart)."""
"""
Rutas locales del ZIP (incluye todas las partes multipart).
La comparación es insensible a mayúsculas en TODOS los pasos, porque en este dominio
los respaldos llegan como .ZIP con frecuencia. Antes se hacía
`path.stem.split(".zip")[0]`, que con "EMPRESA.ZIP.001" dejaba base_name="EMPRESA.ZIP"
y armaba el glob "EMPRESA.ZIP.zip.*": no encontraba nada y devolvía lista vacía. Como
_move_zip_to_processed y _move_zip_to_failed iteran sobre este resultado, las partes
nunca salían de Entrada y se acumulaban mezcladas con los pendientes.
Tampoco se usa glob(): en Linux distingue mayúsculas, así que un patrón en minúsculas
seguiría sin encontrar las partes en MAYÚSCULAS. Se filtra iterdir() comparando en
minúsculas, que funciona igual en Windows y en Linux.
"""
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]
if not SevenZipExtractor.is_multipart(str(path)):
return [str(path)]
# "EMPRESA.ZIP.001" -> stem "EMPRESA.ZIP" -> base "EMPRESA" (sin importar la caja).
stem = path.stem
base_name = stem[:-4] if stem.lower().endswith(".zip") else stem
prefix = f"{base_name}.zip.".lower()
parts = sorted(
(item for item in path.parent.iterdir() if item.name.lower().startswith(prefix)),
key=lambda item: item.name.lower(),
)
if parts:
return [str(item) for item in parts]
# Sin partes localizadas se devuelve el archivo original: es preferible mover solo esa
# parte a no mover nada y dejarla atorada en Entrada para siempre.
app_logger.warning(
f"No se localizaron las partes multipart de {path.name}; se usa solo ese archivo"
)
return [str(path)]
def _move_zip_to_processed(self, job):
@@ -205,15 +246,53 @@ class RestoreWorker(QRunnable):
shutil.move(str(src), str(dest))
app_logger.info(f"Movido: {src.name} -> {dest}")
def _move_zip_to_failed(self, job):
"""Mueve el ZIP (y partes multipart) a Fallados/<fecha>/ para diagnóstico y descarga.
Antes los fallidos quedaban en Entrada (mezclados con pendientes y bloqueando el
pickup por dedup). Al moverlos a Fallados, el panel puede listarlos/descargarlos por
restaurador con la misma lógica relativa que Procesados.
"""
failed_folder = Path(self.config["paths"]["failed_folder"])
date_folder = failed_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)
if not src.exists():
continue
dest = date_folder / src.name
shutil.move(str(src), str(dest))
app_logger.info(f"Movido a Fallados: {src.name} -> {dest}")
def _purge_extract_dir(self):
"""Elimina la carpeta temporal de extracción si quedó (best-effort, cualquier salida).
_cleanup() solo corre en éxito; sin esto, los caminos de fallo/diferido dejan
`Temp/<job_id>/` huérfano acumulándose. Se invoca en el `finally` de run().
"""
try:
if self._extract_dir and Path(self._extract_dir).exists():
shutil.rmtree(self._extract_dir, ignore_errors=True)
app_logger.info(f"Temp de extracción purgado: {self._extract_dir}")
except Exception as e:
app_logger.warning(f"No se pudo purgar Temp {self._extract_dir}: {e}")
def _forward_zip(self, job, route: dict, start_time: float):
"""Reenvía el ZIP al input_folder del servidor destino vía SFTP."""
"""Reenvía el ZIP al input_folder del servidor destino vía SFTP.
La entrega exitosa por SFTP es el punto de no retorno: en cuanto la subida se confirma,
el job se marca COMPLETED/forwarded ANTES de cualquier tarea de limpieza local. Así, un
error POSTERIOR a la entrega (p.ej. mover el ZIP a Procesados) ya no degrada el job a
fallido ni reporta 'failed' al panel — el respaldo sí llegó al destino.
"""
target = route["target"]
self._target = target
remote_folder = target["input_folder"]
zip_paths = self._collect_zip_paths(job.source_path)
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,
@@ -221,15 +300,18 @@ class RestoreWorker(QRunnable):
stdout=f"Destino: {target.get('name')} ({len(uploaded)} archivo(s))",
)
except Exception as e:
# Envío parcial: limpia best-effort las partes ya subidas para no dejar una
# restauración a medias en el destino, marca el step fallido y re-lanza.
partial = getattr(e, "uploaded", None)
if partial:
self._cleanup_partial_forward(target, partial)
JobStepRepository.complete(step_id, exit_code=1, error=str(e))
raise
self._move_zip_to_processed(job)
# --- Entrega confirmada: commit del éxito ANTES de cualquier limpieza local ---
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"
)
@@ -239,8 +321,33 @@ class RestoreWorker(QRunnable):
self.job_id,
)
self._report_to_panel(job, "forwarded", duration_ms=total_ms)
# Housekeeping best-effort: si el move falla, el job SIGUE siendo forwarded y el ZIP
# queda en Entrada (bloqueado por dedup de hash 'completed', no se reenvía en bucle).
try:
self._move_zip_to_processed(job)
except Exception as e:
app_logger.warning(
f"Reenvío OK pero no se pudo mover el ZIP a Procesados: {e}"
)
EventRepository.create(
"WARNING",
f"Reenvío exitoso; el ZIP quedó en Entrada (no se pudo mover a Procesados): {e}",
self.job_id,
)
self.signals.job_completed.emit(self.job_id, True)
def _cleanup_partial_forward(self, target: dict, uploaded: list) -> None:
"""Borra best-effort del destino las partes ya subidas tras un fallo de reenvío."""
for remote_path in uploaded:
try:
sftp_copy.cleanup_remote(target, remote_path)
except Exception as e:
app_logger.warning(
f"No se pudo limpiar la parte remota {remote_path}: {e}"
)
def _report_to_panel(
self,
job,
@@ -479,28 +586,9 @@ class RestoreWorker(QRunnable):
shutil.rmtree(self._extract_dir)
app_logger.info(f"Carpeta de extracción eliminada: {self._extract_dir}")
# Mover ZIP a Processed
processed_folder = Path(self.config["paths"]["processed_folder"])
date_folder = processed_folder / datetime.now().strftime("%Y-%m-%d")
date_folder.mkdir(parents=True, exist_ok=True)
source_path = Path(job.source_path)
dest_path = date_folder / source_path.name
# Si es multipart, mover todas las partes
if SevenZipExtractor.is_multipart(str(source_path)):
# Buscar todas las partes
base_name = source_path.stem.split('.zip')[0]
parts = list(source_path.parent.glob(f"{base_name}.zip.*"))
for part in parts:
part_dest = date_folder / part.name
shutil.move(str(part), str(part_dest))
app_logger.info(f"Movido: {part.name} -> {part_dest}")
else:
shutil.move(str(source_path), str(dest_path))
app_logger.info(f"Movido: {source_path.name} -> {dest_path}")
# Mover ZIP (y partes multipart) a Processed y registrar rel_path/tamaño.
self._move_zip_to_processed(job)
JobStepRepository.complete(step_id, exit_code=0)
except Exception as e:

253
app/engine/retention.py Normal file
View File

@@ -0,0 +1,253 @@
"""Retención diaria de respaldos aplicados para no saturar el disco del servidor.
Dos limpiezas independientes sobre las carpetas locales de CloudRestoreAS:
- ``Procesados/``: por NODO. Para cada nodo se toma su restauración más reciente como
referencia y se conservan las de los últimos ``days``; se borran las anteriores (respaldos
ya aplicados y obsoletos). Nunca se borra la más reciente ni nodos con una sola restauración.
- ``Fallados/``: por ANTIGÜEDAD absoluta. Se borran los ZIP cuya carpeta-fecha sea anterior a
``hoy - failed_days`` (los fallos no tienen semántica de "última restauración exitosa por nodo").
Salvaguardas: solo ``unlink`` de archivos dentro de la carpeta configurada; nunca ``rmtree``;
no sigue symlinks; y jamás toca la carpeta-fecha de la restauración más reciente de un nodo.
El borrado físico se correlaciona con la tabla ``jobs`` (la BD no guarda la ruta destino), por
lo que la carpeta-fecha se deriva de ``finished_at`` con tolerancia de ±1 día por el desfase
UTC/local del momento del movimiento.
"""
from dataclasses import dataclass
from datetime import date, datetime, timedelta, timezone
from pathlib import Path
from typing import Callable, Optional
from ..db.event_repository import EventRepository
from ..db.job_repository import Job, JobRepository
from ..extract.seven_zip import SevenZipExtractor
from ..utils.logger import app_logger
@dataclass
class RetentionResult:
"""Resumen de una corrida de retención."""
dry_run: bool
deleted_files: int = 0
freed_bytes: int = 0
purged_jobs: int = 0
missing_files: int = 0
errors: int = 0
class RetentionCleaner:
"""Aplica la política de retención sobre las carpetas Procesados/ y Fallados/."""
def __init__(
self,
config: dict,
*,
dry_run: Optional[bool] = None,
clock: Callable[[], datetime] = datetime.now,
) -> None:
retention = config.get("retention", {})
paths = config.get("paths", {})
processed_raw = (paths.get("processed_folder") or "").strip()
failed_raw = (paths.get("failed_folder") or "").strip()
self._processed_folder = Path(processed_raw) if processed_raw else None
self._failed_folder = Path(failed_raw) if failed_raw else None
self._days = int(retention.get("days", 2))
self._failed_days = int(retention.get("failed_days", 7))
# El dry_run explícito (p.ej. botón manual) gana sobre la config.
self._dry_run = bool(retention.get("dry_run", True)) if dry_run is None else dry_run
self._clock = clock
# -- Orquestación ----------------------------------------------------------------
def run(self) -> RetentionResult:
"""Ejecuta ambas limpiezas y devuelve el resumen."""
result = RetentionResult(dry_run=self._dry_run)
self._clean_processed(result)
self._clean_failed(result)
mode = "SIMULACRO" if self._dry_run else "real"
freed_mb = result.freed_bytes / (1024 * 1024)
verb = "se borrarían" if self._dry_run else "borrados"
message = (
f"Retención ({mode}): {result.deleted_files} archivo(s) {verb} "
f"({freed_mb:.1f} MB), {result.purged_jobs} job(s) marcados, "
f"{result.missing_files} no hallado(s), {result.errors} error(es)"
)
app_logger.info(message)
EventRepository.create("INFO", message)
return result
# -- Procesados (por nodo) -------------------------------------------------------
def _clean_processed(self, result: RetentionResult) -> None:
if not self._processed_folder or not self._processed_folder.is_dir():
return
obsolete = JobRepository.get_obsolete_completed_by_node(self._days)
if not obsolete:
return
# Fecha-carpeta local de la restauración más reciente de cada nodo: intocable.
ref_dates = {
node: self._local_date_from_iso(finished_at)
for node, finished_at in JobRepository.get_latest_completed_per_node().items()
}
for job in obsolete:
local_date = self._local_date_from_iso(job.finished_at)
if local_date is None:
continue
ref_date = ref_dates.get(job.node_name)
parts, date_folder = self._resolve_processed_paths(job, local_date, ref_date)
if not parts:
# El archivo ya no está en disco (o el nombre no coincide): lo damos por
# purgado para no re-escanearlo indefinidamente.
result.missing_files += 1
self._mark_purged(job, result)
continue
self._delete_paths(parts, self._processed_folder, result)
self._mark_purged(job, result)
if not self._dry_run and date_folder is not None:
self._cleanup_empty_dir(date_folder)
def _resolve_processed_paths(
self, job: Job, local_date: date, ref_date: Optional[date]
) -> tuple[list[Path], Optional[Path]]:
"""Localiza el/los archivo(s) del job en Procesados/ derivando la carpeta-fecha.
Devuelve (partes_existentes, carpeta_fecha) o ([], None) si no se localizó. Nunca
considera la carpeta-fecha de la restauración más reciente del nodo (``ref_date``).
"""
for candidate in self._candidate_dates(local_date):
if ref_date is not None and candidate == ref_date:
continue
date_folder = self._processed_folder / candidate.isoformat()
if not date_folder.is_dir():
continue
parts = self._collect_parts_in_folder(date_folder, job.source_name)
if parts:
return parts, date_folder
return [], None
# -- Fallados (por antigüedad absoluta) ------------------------------------------
def _clean_failed(self, result: RetentionResult) -> None:
if not self._failed_folder or not self._failed_folder.is_dir():
return
cutoff = self._clock().date() - timedelta(days=self._failed_days)
for date_folder in sorted(self._failed_folder.iterdir()):
if not date_folder.is_dir():
continue
folder_date = self._parse_date_folder(date_folder.name)
if folder_date is None:
# Carpeta con nombre que no es una fecha: no la tocamos.
continue
if folder_date >= cutoff:
continue
files = [p for p in date_folder.iterdir() if p.is_file()]
self._delete_paths(files, self._failed_folder, result)
if not self._dry_run:
self._cleanup_empty_dir(date_folder)
# -- Helpers de borrado ----------------------------------------------------------
def _delete_paths(
self, paths: list[Path], root: Path, result: RetentionResult
) -> None:
for path in paths:
try:
if path.is_symlink():
app_logger.warning(f"Retención: se omite symlink {path}")
continue
if not self._is_inside(path, root):
app_logger.warning(
f"Retención: se omite ruta fuera de {root}: {path}"
)
continue
if not path.is_file():
continue
size = path.stat().st_size
if self._dry_run:
app_logger.info(f"[SIMULACRO] Se borraría {path} ({size} bytes)")
else:
path.unlink()
app_logger.info(f"Retención: borrado {path} ({size} bytes)")
result.deleted_files += 1
result.freed_bytes += size
except FileNotFoundError:
result.missing_files += 1
except OSError as e:
app_logger.warning(f"Retención: no se pudo borrar {path}: {e}")
result.errors += 1
def _cleanup_empty_dir(self, folder: Path) -> None:
"""Borra la carpeta-fecha solo si quedó vacía (best-effort, nunca rmtree)."""
try:
if folder.is_dir() and not any(folder.iterdir()):
folder.rmdir()
app_logger.info(f"Retención: carpeta vacía eliminada {folder}")
except OSError as e:
app_logger.warning(f"Retención: no se pudo eliminar carpeta {folder}: {e}")
def _mark_purged(self, job: Job, result: RetentionResult) -> None:
if self._dry_run:
return
JobRepository.mark_purged(job.job_id)
result.purged_jobs += 1
# -- Helpers puros ---------------------------------------------------------------
@staticmethod
def _collect_parts_in_folder(date_folder: Path, source_name: str) -> list[Path]:
"""Rutas del archivo (y sus partes multipart) dentro de una carpeta-fecha."""
if SevenZipExtractor.is_multipart(source_name):
base_name = Path(source_name).stem.split(".zip")[0]
return sorted(date_folder.glob(f"{base_name}.zip.*"))
candidate = date_folder / source_name
return [candidate] if candidate.exists() else []
@staticmethod
def _candidate_dates(local_date: date) -> list[date]:
"""Fecha exacta y ±1 día, para absorber el desfase de medianoche/zona horaria."""
return [local_date, local_date - timedelta(days=1), local_date + timedelta(days=1)]
@staticmethod
def _local_date_from_iso(finished_at: Optional[str]) -> Optional[date]:
"""Convierte un finished_at (ISO UTC naive) a la fecha local del movimiento."""
if not finished_at:
return None
try:
dt = datetime.fromisoformat(finished_at)
except ValueError:
return None
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt.astimezone().date()
@staticmethod
def _parse_date_folder(name: str) -> Optional[date]:
try:
return datetime.strptime(name, "%Y-%m-%d").date()
except ValueError:
return None
@staticmethod
def _is_inside(path: Path, root: Path) -> bool:
"""True si `path` resuelve dentro de `root` (anti path-traversal)."""
try:
resolved = path.resolve()
root_resolved = root.resolve()
except OSError:
return False
return resolved == root_resolved or root_resolved in resolved.parents