#!/usr/bin/env python3
"""
Script para sincronizar directorios en /data con Histories de la API.
Crea/actualiza Histories y sincroniza ficheros de entrada y salida.

Mejoras:
- Reintentos automáticos con backoff exponencial
- Modo dry-run para ver qué se haría sin hacerlo
- Estadísticas de ejecución detalladas
- Validación de integridad con checksums
- NUEVO: Detección de cambios (solo sincroniza lo modificado)
- NUEVO: Logging detallado en archivo
"""

import os
import sys
import requests
import hashlib
import time
import argparse
import json
import logging
import shutil
from pathlib import Path
from typing import Optional, Dict, List, Tuple
from datetime import datetime

# Configuración
API_BASE_URL = "https://mercedesbenz.fresbe.com"
API_DOCS_URL = f"{API_BASE_URL}/docs"

# Rutas relativas al directorio del script
PROJECT_DIR = Path(__file__).parent.resolve()
DATA_DIR = PROJECT_DIR / "data"
LOGS_DIR = PROJECT_DIR / "logs"
CACHE_FILE = PROJECT_DIR / ".sync_cache.json"

CREDENTIALS = {
    "email": "javiercabellos@fresbe.com",
    "password": "password",
    "device_name": "daimler-pipeline-sync"
}

# Configuración de reintentos
MAX_RETRIES = 3
RETRY_DELAY = 1  # segundos (con backoff exponencial)
TIMEOUT = 60  # segundos


def setup_logging(dry_run: bool = False) -> logging.Logger:
    """Configura logging a archivo y consola."""
    LOGS_DIR.mkdir(exist_ok=True)

    logger = logging.getLogger("sync_histories")
    logger.setLevel(logging.DEBUG)

    # Formato detallado
    formatter = logging.Formatter(
        "%(asctime)s - %(levelname)-8s - %(message)s",
        datefmt="%Y-%m-%d %H:%M:%S"
    )

    # Log a archivo
    log_file = LOGS_DIR / f"sync_{datetime.now().strftime('%Y%m%d')}.log"
    file_handler = logging.FileHandler(log_file)
    file_handler.setLevel(logging.DEBUG)
    file_handler.setFormatter(formatter)
    logger.addHandler(file_handler)

    # Log a consola
    console_handler = logging.StreamHandler()
    console_handler.setLevel(logging.INFO)
    console_formatter = logging.Formatter("%(message)s")
    console_handler.setFormatter(console_formatter)
    logger.addHandler(console_handler)

    if dry_run:
        logger.info("=" * 60)
        logger.info("[DRY-RUN] Modo simulación - no se realizarán cambios")
        logger.info("=" * 60)

    return logger


class FileCache:
    """Caché de hashes de ficheros para detectar cambios."""

    def __init__(self, cache_file: Path):
        self.cache_file = cache_file
        self.data = self._load()

    def _load(self) -> Dict:
        """Carga caché desde archivo."""
        if self.cache_file.exists():
            try:
                with open(self.cache_file, 'r') as f:
                    return json.load(f)
            except:
                return {}
        return {}

    def save(self):
        """Guarda caché a archivo."""
        with open(self.cache_file, 'w') as f:
            json.dump(self.data, f, indent=2)

    def get_hash(self, filepath: str) -> Optional[str]:
        """Obtiene hash almacenado de un fichero."""
        return self.data.get(filepath)

    def set_hash(self, filepath: str, file_hash: str):
        """Almacena hash de un fichero."""
        self.data[filepath] = file_hash

    def has_changed(self, filepath: str, current_hash: str) -> bool:
        """Comprueba si un fichero ha cambiado."""
        cached_hash = self.get_hash(filepath)
        if cached_hash is None:
            return True  # Nuevo fichero
        return cached_hash != current_hash


class Stats:
    """Estadísticas de ejecución."""

    def __init__(self, logger):
        self.logger = logger
        self.start_time = datetime.now()
        self.histories_created = 0
        self.histories_updated = 0
        self.files_uploaded = 0
        self.files_skipped = 0
        self.total_size = 0
        self.end_time = None

    def add_file(self, size: int):
        """Registra un fichero subido."""
        self.files_uploaded += 1
        self.total_size += size

    def skip_file(self):
        """Registra un fichero omitido (no cambió)."""
        self.files_skipped += 1

    def finish(self):
        """Marca el fin de la ejecución."""
        self.end_time = datetime.now()

    def get_duration(self) -> float:
        """Retorna duración en segundos."""
        if self.end_time:
            return (self.end_time - self.start_time).total_seconds()
        return (datetime.now() - self.start_time).total_seconds()

    def get_size_mb(self) -> float:
        """Retorna tamaño total en MB."""
        return self.total_size / (1024 * 1024)

    def get_speed_mbps(self) -> float:
        """Retorna velocidad promedio en MB/s."""
        duration = self.get_duration()
        if duration > 0:
            return self.get_size_mb() / duration
        return 0

    def print_summary(self):
        """Imprime resumen de estadísticas."""
        duration = self.get_duration()
        size_mb = self.get_size_mb()
        speed = self.get_speed_mbps()

        self.logger.info("=" * 60)
        self.logger.info("ESTADÍSTICAS DE EJECUCIÓN")
        self.logger.info("=" * 60)
        self.logger.info(f"  Histories creadas: {self.histories_created}")
        self.logger.info(f"  Histories actualizadas: {self.histories_updated}")
        self.logger.info(f"  Ficheros sincronizados: {self.files_uploaded}")
        self.logger.info(f"  Ficheros omitidos (sin cambios): {self.files_skipped}")
        self.logger.info(f"  Tamaño total sincronizado: {size_mb:.2f} MB")
        self.logger.info(f"  Tiempo total: {duration:.2f} segundos")
        if speed > 0:
            self.logger.info(f"  Velocidad promedio: {speed:.2f} MB/s")
        self.logger.info("=" * 60)


class APIClient:
    """Cliente para interactuar con la API con reintentos automáticos."""

    def __init__(self, base_url: str, logger: logging.Logger, dry_run: bool = False):
        self.base_url = base_url
        self.session = requests.Session()
        self.token = None
        self.dry_run = dry_run
        self.logger = logger
        # Headers por defecto
        self.session.headers.update({
            "Accept": "application/json"
        })

    def _retry_request(self, method: str, url: str, **kwargs) -> Optional[requests.Response]:
        """Realiza solicitud HTTP con reintentos automáticos."""
        for attempt in range(MAX_RETRIES):
            try:
                if method.upper() == "GET":
                    response = self.session.get(url, timeout=TIMEOUT, **kwargs)
                elif method.upper() == "POST":
                    response = self.session.post(url, timeout=TIMEOUT, **kwargs)
                elif method.upper() == "PUT":
                    response = self.session.put(url, timeout=TIMEOUT, **kwargs)
                else:
                    return None

                return response

            except requests.exceptions.RequestException as e:
                if attempt < MAX_RETRIES - 1:
                    delay = RETRY_DELAY * (2 ** attempt)
                    self.logger.warning(
                        f"Error en solicitud (intento {attempt + 1}/{MAX_RETRIES}): {str(e)[:50]}"
                    )
                    self.logger.warning(f"Reintentando en {delay} segundos...")
                    time.sleep(delay)
                else:
                    self.logger.error(f"Error después de {MAX_RETRIES} intentos: {e}")
                    return None

        return None

    def login(self, email: str, password: str, device_name: str) -> bool:
        """Autentica con la API y obtiene token."""
        try:
            url = f"{self.base_url}/api/login"
            payload = {
                "email": email,
                "password": password,
                "device_name": device_name
            }

            self.logger.debug(f"Intentando login: {email}")
            response = self._retry_request("POST", url, json=payload)

            if response and response.status_code == 200:
                data = response.json()
                self.token = data.get("token")

                if not self.token:
                    self.logger.error("No se encontró token en la respuesta")
                    return False

                self.session.headers.update({"Authorization": f"Bearer {self.token}"})
                user_name = data.get("user", {}).get("name", "usuario")
                self.logger.info(f"✓ Autenticación exitosa: {user_name} ({email})")
                return True
            else:
                status = response.status_code if response else "Sin respuesta"
                self.logger.error(f"Error de autenticación: {status}")
                return False
        except Exception as e:
            self.logger.error(f"Error al conectar con la API: {e}")
            return False

    def get_histories(self) -> List[Dict]:
        """Obtiene lista de todas las Histories."""
        try:
            url = f"{self.base_url}/api/histories"
            response = self._retry_request("GET", url)

            if response and response.status_code == 200:
                data = response.json()
                if isinstance(data, dict) and "data" in data:
                    return data["data"]
                elif isinstance(data, list):
                    return data
                else:
                    return [data] if data else []
            return []
        except Exception as e:
            self.logger.error(f"Error al obtener Histories: {e}")
            return []

    def get_history_by_name(self, history_name: str) -> Optional[Dict]:
        """Obtiene una History por nombre."""
        try:
            histories = self.get_histories()
            for history in histories:
                if history.get("name") == history_name:
                    return history
            return None
        except Exception as e:
            self.logger.error(f"Error al buscar History '{history_name}': {e}")
            return None

    def create_history(self, history_name: str,
                      input_files: List[Tuple[str, Path]] = None,
                      output_files: List[Tuple[str, Path]] = None) -> Optional[Dict]:
        """Crea una nueva History con ficheros."""
        file_objects = []
        try:
            url = f"{self.base_url}/api/histories"

            data = {
                "name": history_name,
                "description": f"Procesamiento {history_name}"
            }

            files = []

            # Agregar ficheros de entrada
            if input_files:
                for filename, filepath in input_files:
                    if filepath.exists():
                        if self.dry_run:
                            files.append((filename, filepath.stat().st_size))
                        else:
                            file_obj = open(filepath, 'rb')
                            file_objects.append(file_obj)
                            files.append(('input_files[]', (filename, file_obj)))

            # Agregar ficheros de salida
            if output_files:
                for filename, filepath in output_files:
                    if filepath.exists():
                        if self.dry_run:
                            files.append((filename, filepath.stat().st_size))
                        else:
                            file_obj = open(filepath, 'rb')
                            file_objects.append(file_obj)
                            files.append(('output_files[]', (filename, file_obj)))

            if self.dry_run:
                self.logger.info(f"[DRY-RUN] History sería creada: {history_name}")
                return {"id": -1, "name": history_name}

            self.logger.debug(f"Creando History: {history_name}")
            response = self._retry_request("POST", url, data=data, files=files)

            if response and response.status_code in [200, 201]:
                try:
                    response_data = response.json()
                    history = response_data.get("data") if isinstance(response_data, dict) and "data" in response_data else response_data
                    self.logger.info(f"✓ History creada: {history_name} (ID: {history.get('id')})")
                    return history
                except:
                    self.logger.error(f"Error al parsear respuesta JSON")
                    return None
            else:
                status = response.status_code if response else "Sin respuesta"
                self.logger.error(f"Error al crear History '{history_name}': {status}")
                return None
        except Exception as e:
            self.logger.error(f"Error al crear History '{history_name}': {e}")
            return None
        finally:
            for file_obj in file_objects:
                try:
                    file_obj.close()
                except:
                    pass

    def update_history(self, history_id: str,
                      input_files: List[Tuple[str, Path]] = None,
                      output_files: List[Tuple[str, Path]] = None) -> bool:
        """Actualiza los ficheros de una History."""
        file_objects = []
        try:
            url = f"{self.base_url}/api/histories/{history_id}"

            data = {}
            files = []

            if input_files:
                for filename, filepath in input_files:
                    if filepath.exists():
                        if self.dry_run:
                            files.append((filename, filepath.stat().st_size))
                        else:
                            file_obj = open(filepath, 'rb')
                            file_objects.append(file_obj)
                            files.append(('input_files[]', (filename, file_obj)))

            if output_files:
                for filename, filepath in output_files:
                    if filepath.exists():
                        if self.dry_run:
                            files.append((filename, filepath.stat().st_size))
                        else:
                            file_obj = open(filepath, 'rb')
                            file_objects.append(file_obj)
                            files.append(('output_files[]', (filename, file_obj)))

            if not files and not data:
                return True

            if self.dry_run:
                self.logger.info(f"[DRY-RUN] History {history_id} sería actualizada con {len(files)} ficheros")
                return True

            self.logger.debug(f"Actualizando History {history_id}")
            response = self._retry_request("PUT", url, data=data if data else None,
                                          files=files if files else None)

            if response and response.status_code in [200, 204]:
                return True
            else:
                status = response.status_code if response else "Sin respuesta"
                self.logger.error(f"Error al actualizar History '{history_id}': {status}")
                return False
        except Exception as e:
            self.logger.error(f"Error al actualizar History '{history_id}': {e}")
            return False
        finally:
            for file_obj in file_objects:
                try:
                    file_obj.close()
                except:
                    pass


def calculate_file_hash(filepath: Path, algorithm: str = 'md5') -> str:
    """Calcula el hash de un fichero."""
    hash_func = hashlib.new(algorithm)
    with open(filepath, 'rb') as f:
        for chunk in iter(lambda: f.read(4096), b''):
            hash_func.update(chunk)
    return hash_func.hexdigest()


def get_directory_files(dir_path: Path, allowed_extensions: Optional[List[str]] = None) -> List[Tuple[str, Path]]:
    """Obtiene lista de ficheros en un directorio con sus rutas."""
    if not dir_path.exists():
        return []

    files = []
    for file in dir_path.iterdir():
        if file.is_file():
            if allowed_extensions:
                if not any(file.suffix.lower() == ext.lower() for ext in allowed_extensions):
                    continue
            files.append((file.name, file))

    return sorted(files, key=lambda x: x[0])


def delete_directory(dir_path: Path, logger: logging.Logger) -> bool:
    """Borra un directorio completamente."""
    try:
        if dir_path.exists():
            size_mb = sum(f.stat().st_size for f in dir_path.rglob('*') if f.is_file()) / (1024 * 1024)
            logger.info(f"  Borrando directorio: {dir_path.name} ({size_mb:.2f} MB)")
            shutil.rmtree(dir_path)
            logger.info(f"  ✓ Directorio borrado: {dir_path.name}")
            return True
        return False
    except Exception as e:
        logger.error(f"  ✗ Error al borrar directorio '{dir_path.name}': {e}")
        return False


def sync_directory(client: APIClient, dir_path: Path, stats: Stats, cache: FileCache, logger: logging.Logger) -> bool:
    """Sincroniza un directorio con una History en la API."""
    dir_name = dir_path.name
    origin_dir = dir_path / "origin"
    output_dir = dir_path / "output"

    logger.info("")
    logger.info("=" * 60)
    logger.info(f"Sincronizando: {dir_name}")
    logger.info("=" * 60)

    # Obtener ficheros locales
    input_files = get_directory_files(origin_dir)
    output_files = get_directory_files(output_dir)

    # Detectar cambios
    input_files_changed = []
    output_files_changed = []

    logger.info(f"Ficheros en origin/: {len(input_files)}")
    for filename, filepath in input_files:
        file_hash = calculate_file_hash(filepath)
        cache_key = f"{dir_name}/origin/{filename}"

        if cache.has_changed(cache_key, file_hash):
            size_mb = filepath.stat().st_size / (1024 * 1024)
            logger.info(f"  [NUEVO/CAMBIO] {filename} ({size_mb:.2f} MB)")
            input_files_changed.append((filename, filepath))
            cache.set_hash(cache_key, file_hash)
        else:
            logger.debug(f"  [SIN CAMBIOS] {filename}")
            stats.skip_file()

    logger.info(f"Ficheros en output/: {len(output_files)}")
    for filename, filepath in output_files:
        file_hash = calculate_file_hash(filepath)
        cache_key = f"{dir_name}/output/{filename}"

        if cache.has_changed(cache_key, file_hash):
            size_mb = filepath.stat().st_size / (1024 * 1024)
            logger.info(f"  [NUEVO/CAMBIO] {filename} ({size_mb:.2f} MB)")
            output_files_changed.append((filename, filepath))
            cache.set_hash(cache_key, file_hash)
        else:
            logger.debug(f"  [SIN CAMBIOS] {filename}")
            stats.skip_file()

    # Si no hay cambios, omitir sincronización
    if not input_files_changed and not output_files_changed:
        logger.info(f"  → Sin cambios detectados")
        return True

    # Buscar o crear History
    history = client.get_history_by_name(dir_name)

    if history:
        logger.info(f"✓ History encontrada: {history.get('id', 'unknown')}")
        history_id = history.get('id')
        is_new = False
    else:
        logger.info(f"→ Creando nueva History...")
        history = client.create_history(
            dir_name,
            input_files=input_files_changed if input_files_changed else None,
            output_files=output_files_changed if output_files_changed else None
        )

        if not history:
            logger.error(f"No se pudo crear History")
            return False

        history_id = history.get('id')
        is_new = True
        if not client.dry_run:
            stats.histories_created += 1
            for _, filepath in (input_files_changed or []) + (output_files_changed or []):
                stats.add_file(filepath.stat().st_size)

    # Actualizar ficheros si la History ya existía
    success = True

    if not is_new:
        if input_files_changed:
            logger.info(f"→ Actualizando {len(input_files_changed)} fichero(s) de entrada...")
            if client.update_history(history_id, input_files=input_files_changed):
                logger.info(f"  ✓ Ficheros de entrada sincronizados")
                if not client.dry_run:
                    for _, filepath in input_files_changed:
                        stats.add_file(filepath.stat().st_size)
                    stats.histories_updated += 1
            else:
                success = False

        if output_files_changed:
            logger.info(f"→ Actualizando {len(output_files_changed)} fichero(s) de salida...")
            if client.update_history(history_id, output_files=output_files_changed):
                logger.info(f"  ✓ Ficheros de salida sincronizados")
                if not client.dry_run:
                    for _, filepath in output_files_changed:
                        stats.add_file(filepath.stat().st_size)
                    stats.histories_updated += 1
            else:
                success = False
    else:
        if not client.dry_run and (input_files_changed or output_files_changed):
            logger.info(f"✓ Ficheros incluidos en creación")

    # Validación de integridad
    if not client.dry_run and success and (input_files_changed or output_files_changed):
        logger.info(f"→ Validando integridad...")
        for filename, filepath in (input_files_changed or []) + (output_files_changed or []):
            file_hash = calculate_file_hash(filepath)
            logger.debug(f"  {filename}: {file_hash[:8]}... ✓")
        logger.info(f"✓ Validación completada")

    return success


def main():
    """Función principal."""
    parser = argparse.ArgumentParser(
        description="Sincronizador de Histories - daimler-pipeline"
    )
    parser.add_argument(
        "--dry-run",
        action="store_true",
        help="Mostrar qué se haría sin hacerlo realmente"
    )
    parser.add_argument(
        "--cleanup",
        action="store_true",
        help="Borrar directorios de /data después de sincronizar exitosamente (sin confirmación)"
    )

    args = parser.parse_args()

    # Configurar logging
    logger = setup_logging(dry_run=args.dry_run)

    logger.info("╔════════════════════════════════════════════════════════════╗")
    logger.info("║  Sincronizador de Histories - daimler-pipeline            ║")
    if args.dry_run:
        logger.info("║  MODO: DRY-RUN (Sin hacer cambios)                       ║")
    if args.cleanup:
        logger.info("║  CLEANUP: Se borrarán directorios después de sincronizar  ║")
    logger.info("╚════════════════════════════════════════════════════════════╝")
    logger.info(f"API: {API_BASE_URL}")
    logger.info(f"Data dir: {DATA_DIR}")
    logger.info(f"Logs: {LOGS_DIR}")

    # Verificar directorio data
    if not DATA_DIR.exists():
        logger.error(f"El directorio {DATA_DIR} no existe")
        return False

    # Crear cliente API
    client = APIClient(API_BASE_URL, logger, dry_run=args.dry_run)

    # Autenticarse
    if not client.login(CREDENTIALS["email"], CREDENTIALS["password"], CREDENTIALS["device_name"]):
        logger.error("No se pudo autenticar con la API")
        return False

    # Cargar caché
    cache = FileCache(CACHE_FILE)
    logger.debug(f"Caché cargada: {len(cache.data)} entradas previas")

    # Listar directorios
    directories = [d for d in DATA_DIR.iterdir() if d.is_dir() and not d.name.startswith('.')]
    directories.sort()

    if not directories:
        logger.warning(f"No hay directorios en {DATA_DIR}")
        return True

    logger.info(f"\nEncontrados {len(directories)} directorio(s) para sincronizar:")
    for d in directories:
        logger.info(f"  - {d.name}")

    # Sincronizar
    results = []
    stats = Stats(logger)

    for dir_path in directories:
        success = sync_directory(client, dir_path, stats, cache, logger)
        results.append((dir_path.name, success))

    # Guardar caché actualizada
    cache.save()
    logger.debug(f"Caché guardada: {len(cache.data)} entradas")

    # Finalizar estadísticas
    stats.finish()

    # Resumen final
    logger.info("")
    logger.info("=" * 60)
    logger.info("RESUMEN DE SINCRONIZACIÓN")
    logger.info("=" * 60)

    successful = sum(1 for _, success in results if success)
    failed = len(results) - successful

    for dir_name, success in results:
        status = "✓ OK" if success else "✗ FALLO"
        logger.info(f"  {status}: {dir_name}")

    logger.info(f"\nTotal: {successful} exitoso(s), {failed} fallido(s)")

    # Mostrar estadísticas
    if not args.dry_run:
        stats.print_summary()
    else:
        logger.info(f"\n[DRY-RUN] Ejecución simulada completada sin cambios")

    # Cleanup: borrar directorios después de sincronizar exitosamente
    if args.cleanup and not args.dry_run and failed == 0:
        logger.info("")
        logger.info("=" * 60)
        logger.info("LIMPIEZA DE DIRECTORIOS")
        logger.info("=" * 60)

        # Calcular tamaño total a borrar
        total_size = sum(
            sum(f.stat().st_size for f in d.rglob('*') if f.is_file())
            for d in directories
        ) / (1024 * 1024)

        logger.warning(f"Borrando {len(directories)} directorio(s) ({total_size:.2f} MB)")

        # Ejecutar cleanup
        deleted_count = 0
        for dir_path in directories:
            if delete_directory(dir_path, logger):
                deleted_count += 1

        logger.info("")
        logger.info(f"✓ Limpieza completada: {deleted_count} directorio(s) borrado(s)")

    elif args.cleanup and args.dry_run:
        logger.info("")
        logger.info("[DRY-RUN] Los directorios NO serán borrados en modo simulación")

    logger.info(f"\nLog guardado en: {LOGS_DIR}/sync_{datetime.now().strftime('%Y%m%d')}.log")

    return failed == 0


if __name__ == "__main__":
    try:
        success = main()
        sys.exit(0 if success else 1)
    except KeyboardInterrupt:
        print("\n⚠ Sincronización cancelada por el usuario")
        sys.exit(1)
    except Exception as e:
        print(f"✗ Error inesperado: {e}")
        import traceback
        traceback.print_exc()
        sys.exit(1)
