from pymongo import MongoClient
from pymongo.errors import PyMongoError, CollectionInvalid, BulkWriteError
from datetime import datetime, timedelta
from typing import Dict, Any, List, Union, Optional


FLUX_DATABASE_NAME = "FLUX__TS"
#Tiempo para considerar desconexión y generar valor fantasma para FE
UMBRAL_SEGUNDOS = 300
# hora_default = " --"
DEFAULT_LAST_SYNC_TIME =  " --" # str(datetime.now().strftime("%Y-%m-%d %H:%M:%S"))

def get_flux_collection(db, centro_id):
    """
    Helper para obtener/crear la colección TimeSeries correcta.
    """
    collection_name = f"{centro_id.lower()}_flux_ts"
    
    # Verificación Lazy de existencia
    if collection_name not in db.list_collection_names():
        try:
            db.create_collection(
                collection_name,
                timeseries={
                    'timeField': 'timestamp',
                    'granularity': 'minutes'
                }
            )
        except CollectionInvalid:
            pass # Ya existía (condición de carrera)
            
    return db[collection_name]


def parse_flux_document(item):
    """
    Helper para validar y transformar un raw json en un documento mongo.
    """
    raw_timestamp = item.get("timestamp")
    if not raw_timestamp:
        raise ValueError("Campo 'timestamp' faltante")

    try:
        # Ajustar formato según lo que envíe el Host
        dt_object = datetime.strptime(raw_timestamp, '%Y-%m-%d %H:%M:%S')
    except ValueError:
        raise ValueError(f"Formato fecha inválido: {raw_timestamp}")

    external_id = item.get("id")

    return {
        "timestamp": dt_object,
        "flow_rate": item.get("Flow rate"),
        "cumulant": item.get("Cumulant float point nummber"),
        "connected": item.get("connected"),
        "received_at": datetime.utcnow(),
        "original_id": external_id # Guardamos el ID original por trazabilidad
    }


def save_flux_single(client, data, centro_id):
    """
    Procesa UN solo registro (Modo Síncrono) con manejo de errores.
    - Si tiene éxito: Retorna None.
    - Si falla: Retorna False.
    """
    try:
        db = client[FLUX_DATABASE_NAME]
        collection = get_flux_collection(db, centro_id)

        # 1. Parsear documento
        doc = parse_flux_document(data)
        
        # 2. Guardar
        collection.insert_one(doc)
        
        # 3. Retorno Exitoso (None solicitado)
        return None

    except Exception as e:
        # Logueamos el error localmente
        print(f"[ERROR] save_flux_single: {str(e)}")
        # Retornamos False para que el controlador sepa que falló
        return False


def save_flux_batch(client, data_list, centro_id):
    """
    Procesa un ARRAY de registros (Modo Asíncrono).
    Siempre guarda. Retorna lista PLANA UNIDIMENSIONAL de IDs encontrados.
    """
    db = client[FLUX_DATABASE_NAME]
    collection = get_flux_collection(db, centro_id)

    docs_to_insert = []
    ids_to_ack = []

    # 1. Procesar lista en memoria
    for item in data_list:
        try:
            doc = parse_flux_document(item)
            docs_to_insert.append(doc)
            
            # Lógica de Aplanado (Flatten)
            oid = doc["original_id"]
            
            if oid is not None:
                if isinstance(oid, list):
                    ids_to_ack.extend(oid) # Si es lista [1,2], agrega 1 y 2 al nivel actual
                else:
                    ids_to_ack.append(oid) # Si es escalar 1, agrega 1
                
        except ValueError as e:
            print(f"Advertencia: Saltando item corrupto en batch Flux: {e}")
            continue

    # 2. Inserción Masiva
    if docs_to_insert:
        try:
            collection.insert_many(docs_to_insert, ordered=False)
        except BulkWriteError as bwe:
            print(f"Advertencia Bulk: {bwe.details}")


    return ids_to_ack

def get_flux_historical_data(db, collection_name, sensor_key, N=360, step=1):
    """
    Obtiene datos históricos de Flux.
    Si no hay datos, retorna la estructura vacía estandarizada.
    Si el último dato es > 5 min, agrega estado desconectado al final.
    """
    
    # Estructura vacía por defecto
    empty_result = {
        sensor_key: {
            "flow_rate": [],
            "cumulant": [],
            "connection_status": [],
            "TS": [DEFAULT_LAST_SYNC_TIME],
        }
    }

    try:
        # Verificación de existencia de colección
        if collection_name not in db.list_collection_names():
            return empty_result

        collection = db[collection_name]

        # Definir Rango de Tiempo (Día actual)
        ahora = datetime.now()
        inicio_hoy = ahora.replace(hour=0, minute=0, second=0, microsecond=0)
        inicio_manana = inicio_hoy + timedelta(days=1)

        # Consulta
        query = {
            "timestamp": {"$gte": inicio_hoy, "$lt": inicio_manana}
        }
        
        projection = {
            "flow_rate": 1, 
            "cumulant": 1, 
            "connected": 1, 
            "timestamp": 1, 
            "_id": 0
        }

        # Traemos los N más recientes (Descendente para obtener el último rápido)
        cursor = collection.find(query, projection).sort("timestamp", -1).limit(N)
        
        docs = list(cursor)

        if not docs:
            return empty_result

        # Procesamiento (Invertir para tener orden cronológico)
        # docs[0] es el más reciente de la DB, docs[-1] el más antiguo
        # Al invertirlo: docs_asc[-1] será el más reciente.
        docs_asc = list(reversed(docs))
        
        # Aplicamos el step
        docs_sampled = docs_asc[::step]

        # Extracción de Vectores
        flow_rate_vals = []
        cumulant_vals = []
        connected_vals = []
        ts_vals = []

        for doc in docs_sampled:
            raw_ts = doc.get("timestamp")

            # Formateo string para el JSON
            ts_str = raw_ts.isoformat() if isinstance(raw_ts, datetime) else str(raw_ts)
            
            flow_rate_vals.append(doc.get("flow_rate"))
            cumulant_vals.append(doc.get("cumulant"))
            connected_vals.append(doc.get("connected")) 
            ts_vals.append(ts_str)

        # Empaquetamos en un diccionario temporal para pasarlo a la auxiliar
        vectors_package = {
            "flow_rate": flow_rate_vals,
            "cumulant": cumulant_vals,
            "connection_status": connected_vals,
            "TS": ts_vals
        }

        # Construcción de Estructura de Respuesta Final
        resultado = {
            sensor_key: vectors_package
        }
        
        return resultado

    except Exception as e:
        print(f'Error obteniendo historicos Flux: {e}')
        return empty_result


def get_flux_heartbeat(db, collection_name):
    """
    Busca el último registro recibido HOY para determinar el estado de conexión.
    Retorna la fecha formateada o ' --' si no hay datos.
    """

    hora_default = DEFAULT_LAST_SYNC_TIME
    # print(f"{hora_default=}")
    try:
        # Verificación rápida de existencia
        if collection_name not in db.list_collection_names():
            return hora_default

        collection = db[collection_name]

        # Rango de hoy (para filtrar datos del día operativo)
        ahora = datetime.now()
        inicio_hoy = ahora.replace(hour=0, minute=0, second=0, microsecond=0)
        inicio_manana = inicio_hoy + timedelta(days=1)

        # Consulta optimizada: Filtra por día, ordena por recepción
        last_heartbeat_doc = collection.find_one(
            {
                "timestamp": {"$gte": inicio_hoy, "$lt": inicio_manana}
            },
            sort=[("timestamp", -1)], # El más reciente recibido
            projection={"timestamp": 1, "_id": 0} 
        )

        if last_heartbeat_doc and "timestamp" in last_heartbeat_doc:
            rec_at = last_heartbeat_doc["timestamp"]
            
            # Formateo seguro
            if isinstance(rec_at, datetime):
                return rec_at.strftime("%Y-%m-%d %H:%M:%S")
            return str(rec_at)
            
        return hora_default

    except Exception as e:
        print(f"Advertencia en get_flux_heartbeat: {e}")
        return hora_default

def serialize_timestamp(timestamp):
    """
    @brief Verifica si un timestamp es un objeto datetime y lo convierte a ISO 8601.

    @param timestamp Puede ser un objeto datetime o una cadena.
    @return string con el timestamp en formato ISO 8601 si es datetime.
            Si ya es string, se devuelve sin cambios.
    @throws ValueError si el timestamp no es un datetime o string válido.
    """
    if isinstance(timestamp, datetime):  # Si es un objeto datetime (ISODate en MongoDB)
        return timestamp.isoformat()
    if isinstance(timestamp, str):  # Si ya es un string, lo retorna directamente
        return timestamp

    raise ValueError(f"Formato de timestamp inválido: {timestamp} ({type(timestamp)})")

#----------LEGACY----------------#
def legacy_get_flux_historical_data(db, collection_name, sensor_key, N=360, step=1):
    """
    Obtiene datos históricos de Flux.
    Si no hay datos, retorna la estructura vacía estandarizada.
    """
    # 1. Definir la estructura de respuesta vacía por defecto
    empty_result = {
        sensor_key: {
            "flow_rate": [None],
            "cumulant": [None],
            "connection_status": [False],
            "TS": [DEFAULT_LAST_SYNC_TIME],
        }
    }

    try:
        # 2. Verificación de existencia de colección
        if collection_name not in db.list_collection_names():
            # Retornamos la estructura vacía en lugar de {}
            return empty_result

        collection = db[collection_name]

        # 3. Definir Rango de Tiempo (Día actual)
        ahora = datetime.now()
        inicio_hoy = ahora.replace(hour=0, minute=0, second=0, microsecond=0)
        inicio_manana = inicio_hoy + timedelta(days=1)

        # 4. Consulta
        query = {
            "timestamp": {"$gte": inicio_hoy, "$lt": inicio_manana}
        }
        
        projection = {
            "flow_rate": 1, 
            "cumulant": 1, 
            "connected": 1, 
            "timestamp": 1, 
            "_id": 0
        }

        # Traemos los N más recientes
        cursor = collection.find(query, projection).sort("timestamp", -1).limit(N)
        
        docs = list(cursor)

        if not docs:
            # Retornamos la estructura vacía en lugar de {}
            return empty_result

        # 5. Procesamiento (Invertir y Muestrear)
        # Invertimos para tener orden cronológico (Ascendente: Mañana -> Tarde)
        docs_asc = list(reversed(docs))
        
        # Aplicamos el step
        docs_sampled = docs_asc[::step]

        # 6. Extracción de Vectores
        flow_rate_vals = []
        cumulant_vals = []
        connected_vals = []
        ts_vals = []

        for doc in docs_sampled:
            ts_str = doc.get("timestamp").isoformat() if isinstance(doc.get("timestamp"), datetime) else str(doc.get("timestamp"))
            
            flow_rate_vals.append(doc.get("flow_rate"))
            cumulant_vals.append(doc.get("cumulant"))
            connected_vals.append(doc.get("connected")) 
            ts_vals.append(ts_str)

        # 7. Construcción de Estructura de Respuesta con datos
        resultado = {
            sensor_key: {
                "flow_rate": flow_rate_vals,
                "cumulant": cumulant_vals,
                "connection_status": connected_vals,
                "TS": ts_vals,
            }
        }
        
        return resultado

    except Exception as e:
        print(f'Error obteniendo historicos Flux: {e}')
        # En caso de error, también es seguro retornar la estructura vacía para no romper el front
        return empty_result
