212 lines
6.6 KiB
Python
212 lines
6.6 KiB
Python
"""File watcher para detectar nuevos archivos ZIP."""
|
|
|
|
import os
|
|
import time
|
|
import hashlib
|
|
from pathlib import Path
|
|
from typing import Optional, Dict, Callable
|
|
from threading import Thread, Event
|
|
from ..utils.logger import app_logger
|
|
|
|
|
|
class FileStabilityChecker:
|
|
"""Verifica que un archivo esté estable (no siendo copiado)."""
|
|
|
|
def __init__(
|
|
self,
|
|
check_interval: int = 5,
|
|
stable_duration: int = 10,
|
|
use_ready_marker: bool = False
|
|
):
|
|
"""
|
|
Inicializa el verificador de estabilidad.
|
|
|
|
Args:
|
|
check_interval: Intervalo entre verificaciones (segundos)
|
|
stable_duration: Tiempo que debe estar estable (segundos)
|
|
use_ready_marker: Si True, requiere archivo .ready
|
|
"""
|
|
self.check_interval = check_interval
|
|
self.stable_duration = stable_duration
|
|
self.use_ready_marker = use_ready_marker
|
|
|
|
def is_file_ready(self, file_path: str) -> bool:
|
|
"""
|
|
Verifica si un archivo está listo para procesarse.
|
|
|
|
Args:
|
|
file_path: Ruta al archivo
|
|
|
|
Returns:
|
|
True si está listo
|
|
"""
|
|
path = Path(file_path)
|
|
|
|
if not path.exists():
|
|
return False
|
|
|
|
# Verificar marcador .ready si está habilitado
|
|
if self.use_ready_marker:
|
|
ready_marker = path.parent / f"{path.name}.ready"
|
|
if not ready_marker.exists():
|
|
app_logger.debug(f"Esperando marcador .ready para {path.name}")
|
|
return False
|
|
|
|
# Verificar estabilidad de tamaño
|
|
try:
|
|
initial_size = path.stat().st_size
|
|
initial_mtime = path.stat().st_mtime
|
|
|
|
time.sleep(self.stable_duration)
|
|
|
|
if not path.exists():
|
|
return False
|
|
|
|
final_size = path.stat().st_size
|
|
final_mtime = path.stat().st_mtime
|
|
|
|
is_stable = (initial_size == final_size and initial_mtime == final_mtime)
|
|
|
|
if not is_stable:
|
|
app_logger.debug(
|
|
f"Archivo {path.name} aún está cambiando "
|
|
f"(size: {initial_size} -> {final_size})"
|
|
)
|
|
|
|
return is_stable
|
|
|
|
except Exception as e:
|
|
app_logger.error(f"Error verificando estabilidad de {file_path}: {e}")
|
|
return False
|
|
|
|
|
|
class FileWatcher:
|
|
"""Vigila una carpeta por nuevos archivos ZIP."""
|
|
|
|
def __init__(
|
|
self,
|
|
watch_folder: str,
|
|
on_file_ready: Callable[[str], None],
|
|
scan_interval: int = 30,
|
|
stability_checker: Optional[FileStabilityChecker] = None
|
|
):
|
|
"""
|
|
Inicializa el file watcher.
|
|
|
|
Args:
|
|
watch_folder: Carpeta a vigilar
|
|
on_file_ready: Callback cuando un archivo está listo
|
|
scan_interval: Intervalo de escaneo (segundos)
|
|
stability_checker: Verificador de estabilidad (opcional)
|
|
"""
|
|
self.watch_folder = watch_folder
|
|
self.on_file_ready = on_file_ready
|
|
self.scan_interval = scan_interval
|
|
self.stability_checker = stability_checker or FileStabilityChecker()
|
|
|
|
self._stop_event = Event()
|
|
self._thread: Optional[Thread] = None
|
|
self._known_files: Dict[str, float] = {} # path -> mtime
|
|
|
|
def start(self):
|
|
"""Inicia el watcher en un thread separado."""
|
|
if self._thread and self._thread.is_alive():
|
|
app_logger.warning("FileWatcher ya está corriendo")
|
|
return
|
|
|
|
self._stop_event.clear()
|
|
self._thread = Thread(target=self._run, daemon=True)
|
|
self._thread.start()
|
|
app_logger.info(f"FileWatcher iniciado en: {self.watch_folder}")
|
|
|
|
def stop(self):
|
|
"""Detiene el watcher."""
|
|
if self._thread:
|
|
self._stop_event.set()
|
|
self._thread.join(timeout=5)
|
|
app_logger.info("FileWatcher detenido")
|
|
|
|
def scan_now(self):
|
|
"""Fuerza un escaneo inmediato."""
|
|
app_logger.info("Escaneo manual solicitado")
|
|
self._scan_folder()
|
|
|
|
def _run(self):
|
|
"""Loop principal del watcher."""
|
|
while not self._stop_event.is_set():
|
|
try:
|
|
self._scan_folder()
|
|
except Exception as e:
|
|
app_logger.error(f"Error en FileWatcher: {e}", exc_info=True)
|
|
|
|
# Esperar con posibilidad de interrumpir
|
|
self._stop_event.wait(self.scan_interval)
|
|
|
|
def _scan_folder(self):
|
|
"""Escanea la carpeta por archivos nuevos."""
|
|
folder = Path(self.watch_folder)
|
|
|
|
if not folder.exists():
|
|
app_logger.warning(f"Carpeta de entrada no existe: {self.watch_folder}")
|
|
return
|
|
|
|
# Buscar archivos .zip y .zip.001
|
|
patterns = ["*.zip", "*.zip.001"]
|
|
found_files = []
|
|
|
|
for pattern in patterns:
|
|
found_files.extend(folder.glob(pattern))
|
|
|
|
for file_path in found_files:
|
|
file_str = str(file_path)
|
|
|
|
# Ignorar archivos ya procesados recientemente
|
|
if file_str in self._known_files:
|
|
continue
|
|
|
|
# Verificar si está listo
|
|
if not self.stability_checker.is_file_ready(file_str):
|
|
continue
|
|
|
|
# Marcar como conocido
|
|
self._known_files[file_str] = file_path.stat().st_mtime
|
|
|
|
# Notificar
|
|
app_logger.info(f"Archivo nuevo detectado y listo: {file_path.name}")
|
|
try:
|
|
self.on_file_ready(file_str)
|
|
except Exception as e:
|
|
app_logger.error(f"Error procesando archivo {file_path.name}: {e}")
|
|
|
|
# Limpiar archivos conocidos que ya no existen
|
|
self._clean_known_files()
|
|
|
|
def _clean_known_files(self):
|
|
"""Limpia archivos conocidos que ya no existen."""
|
|
to_remove = []
|
|
for file_path in list(self._known_files.keys()):
|
|
if not Path(file_path).exists():
|
|
to_remove.append(file_path)
|
|
|
|
for file_path in to_remove:
|
|
del self._known_files[file_path]
|
|
|
|
|
|
def calculate_file_hash(file_path: str) -> str:
|
|
"""
|
|
Calcula el SHA256 hash de un archivo.
|
|
|
|
Args:
|
|
file_path: Ruta al archivo
|
|
|
|
Returns:
|
|
Hash hexadecimal
|
|
"""
|
|
sha256 = hashlib.sha256()
|
|
|
|
with open(file_path, 'rb') as f:
|
|
while chunk := f.read(8192):
|
|
sha256.update(chunk)
|
|
|
|
return sha256.hexdigest()
|