# src/processors/realtime.py
"""
Procesador de datos en tiempo real - MVP v1.0

Gestiona el procesamiento continuo de nuevas mediciones,
manteniendo estadísticas actualizadas por períodos de 15 minutos.

Este módulo implementa:
- Monitoreo continuo de nuevas mediciones
- Procesamiento en períodos de 15 minutos
- Gestión de períodos incompletos
- Control de timeouts
"""

from datetime import datetime, timedelta
import logging
from mysql.connector import Error 
from ..base import BaseProcessor
from ...db.database import (
    get_unprocessed_measurements,
    get_active_sensors,
    mark_as_processed_batch,  
    insert_registros_batch    
)
from ...models.entrada import Entrada
from ...config.settings import PROCESS_INTERVAL, BATCH_SIZE

logger = logging.getLogger(__name__)

class RealtimeProcessor(BaseProcessor):
    def __init__(self, db_name):
        super().__init__(db_name)
        self.entradas = {}  # {dmac: Entrada}
        self.last_process_time = None
        self.period_measurements = {}  # {dmac: count}
        self.closed_periods = {}  # {dmac: último_período_cerrado}
        self.initialize_entradas()
        
    def initialize_entradas(self):
        """
        Inicializa las entradas activas para el período actual.
        """
        if not self.check_connection():
            return

        active_sensors = get_active_sensors(self.conn)
        now = datetime.now()
        
        # Normalizar el período actual
        periodo_actual = now.replace(
            minute=(now.minute // PROCESS_INTERVAL) * PROCESS_INTERVAL,
            second=0,
            microsecond=0
        )
        
        for sensor_id, dmac, id_entrada in active_sensors:
            self.entradas[dmac] = Entrada(id_entrada, dmac, periodo_actual)
            self.period_measurements[dmac] = 0
            
            if self.debug:
                logger.debug(f"Entrada {dmac} inicializada para período {periodo_actual:%Y-%m-%d %H:%M}")

        logger.info(f"Inicializadas {len(self.entradas)} entradas activas")

    def process_new_data(self):
        """
        Procesa las nuevas mediciones en tiempo real.
        """
        try:
            if not self.check_connection():
                return

            now = datetime.now()
            if self.debug:
                logger.debug(f"Iniciando process_new_data en {now}")

            measurements = get_unprocessed_measurements(self.conn, self.debug)
            
            if self.debug:
                logger.debug(f"Obtenidas {len(measurements) if measurements else 0} mediciones sin procesar")
            
            # Procesar mediciones
            for m_id, dmac, temp, fecha, id_entrada in measurements:
                if dmac not in self.entradas:
                    if self.debug:
                        logger.debug(f"Medición para entrada no inicializada: {dmac}")
                    continue

                entrada = self.entradas[dmac]
                
                # Si es un nuevo período, cerrar el anterior
                if fecha >= entrada.periodo_fin:
                    if self.debug:
                        logger.debug(
                            f"Nueva medición {fecha} fuera del período actual "
                            f"({entrada.periodo_inicio:%Y-%m-%d %H:%M} - {entrada.periodo_fin:%Y-%m-%d %H:%M}) "
                            f"para entrada {dmac}"
                        )
                    self._close_period(entrada)
                    entrada.reiniciar_periodo(fecha)
                    self.period_measurements[dmac] = 0

                # Procesar medición
                if entrada.agregar_medicion(temp, fecha):
                    self.period_measurements[dmac] = self.period_measurements.get(dmac, 0) + 1

            self.last_process_time = now
            
            # Verificar períodos vencidos
            self._check_periods(now)

        except Exception as e:
            logger.error(f"Error procesando nuevos datos: {e}", exc_info=True)

    def show_measurement_status(self):
        """
        Muestra el estado actual de las mediciones.
        """
        print("\n=== Mediciones Actuales ===")
        total = sum(self.period_measurements.values())
        print(f"Total mediciones período: {total}")
        
        for dmac, entrada in self.entradas.items():
            count = self.period_measurements.get(dmac, 0)
            if count > 0:
                stats = entrada.get_estadisticas()
                if stats:
                    print(f"\nEntrada {dmac}:")
                    print(f"  Mediciones: {count}")
                    print(f"  Última temperatura: {stats['promedio']}°C")
                    print(f"  Máx/Min: {stats['maximo']}°C / {stats['minimo']}°C")

    def _check_periods(self, current_time):
        """
        Verifica y cierra períodos que hayan expirado.
        """
        if self.debug:
            logger.debug(f"Verificando períodos en: {current_time}")
            
        for dmac, entrada in self.entradas.items():
            if current_time >= entrada.periodo_fin:
                if self.debug:
                    logger.debug(
                        f"Cerrando período para entrada {dmac} - "
                        f"período actual: {entrada.periodo_inicio:%Y-%m-%d %H:%M} - {entrada.periodo_fin:%Y-%m-%d %H:%M}"
                    )
                self._close_period(entrada)
                entrada.reiniciar_periodo(entrada.periodo_fin)  # Importante: Usamos periodo_fin como inicio del siguiente
                self.period_measurements[dmac] = 0
    
    def _close_period(self, entrada):
        """
        Cierra un período para una entrada.
        """
        stats = entrada.get_estadisticas()
        if not stats or stats['cantidad_mediciones'] == 0:
            if self.debug:
                logger.debug(f"No hay estadísticas para cerrar período de entrada {entrada.dmac}")
            return
            
        if self.debug:
            logger.debug(
                f"Cerrando período para entrada {entrada.dmac}:\n"
                f"  Período: {stats['periodo_inicio']:%Y-%m-%d %H:%M} - {stats['periodo_fin']:%Y-%m-%d %H:%M}\n"
                f"  Mediciones: {stats['cantidad_mediciones']}\n"
                f"  Promedio: {stats['promedio']}°C\n"
                f"  Min/Max: {stats['minimo']}°C / {stats['maximo']}°C"
            )

        cursor = None
        try:
            cursor = self.conn.cursor()
            
            # Aquí está el cambio clave: usamos periodo_inicio para el registro
            registro = (
                0,  # idSistema
                entrada.id_entrada,
                entrada.periodo_fin,  # Usamos período_inicio directamente
                stats['promedio'],
                stats['maximo'],
                stats['minimo'],
                stats['desvio']
            )

            query = """
            INSERT IGNORE INTO registros 
                (idSistema, idEntrada, fecha, medio, maximo, minimo, desvio)
            VALUES 
                (%s, %s, %s, %s, %s, %s, %s)
            """
            
            if self.debug:
                logger.debug(f"Ejecutando query para período {entrada.periodo_inicio:%Y-%m-%d %H:%M}")
                
            cursor.execute(query, registro)
            self.conn.commit()
            
            if self.debug:
                logger.debug(f"Registro insertado exitosamente para entrada {entrada.dmac}")
                
        except Error as e:
            self.conn.rollback()
            logger.error(f"Error insertando registro para entrada {entrada.dmac}: {e}")
            if self.debug:
                logger.debug("Error detallado:", exc_info=True)
        finally:
            if cursor:
                cursor.close()
                
                
    def stop(self):
        """
        Detiene el procesador y limpia recursos.
        """
        logger.info("Deteniendo el procesador...")
        self.running = False
        
        if self.conn:
            try:
                if self.conn.is_connected():
                    try:
                        # Asegurarnos de que no hay resultados pendientes
                        cursor = self.conn.cursor(buffered=True)
                        cursor.close()
                    except:
                        pass
                        
                    try:
                        self.conn.commit()  # Commit cualquier transacción pendiente
                    except:
                        pass
                        
                    try:
                        self.conn.close()
                        logger.info("Conexión a base de datos cerrada")
                    except:
                        pass
            except Exception as e:
                logger.error(f"Error cerrando conexión: {e}")
            finally:
                self.conn = None