#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
Bot de notificacion de LOS (Loss of Signal) - Multi-red
====================================================================
Estrategia: Opcion A (API Felix)

Por que esta opcion y no la BD (onts_estado):
- Revisando actualizar_onts.php se confirma que la tabla onts_estado
  NUNCA registra cambios del campo online: $onts_online se fija a 0 al
  insertar y no se vuelve a tocar en el UPDATE. Es decir, la BD local
  hoy no contiene informacion fiable de online/offline.
- customer_technical_info_report de Felix SI trae "online" por cpe_id,
  ademas de olt_id/frame/slot/port (para armar el PON) y full_name
  (cliente), asi que no hace falta tocar tu PHP ni el OLT via SNMP.

Requiere: Python 3.5+, requests
    pip3 install requests

IMPORTANTE - revisar antes de usar en produccion:
1) TELEGRAM_TOKEN / TELEGRAM_CHAT_ID: rellena con los tuyos.
2) La primera ejecucion de CADA RED no envia notificaciones: solo
   guarda el estado base de todas las ONTs de esa red, para no
   disparar una alarma masiva con ONTs que ya estaban offline desde
   antes de arrancar el bot. Cada red tiene su propio fichero de
   estado, asi que si anades una red nueva mas adelante, solo esa red
   pasara por "primera ejecucion".
3) Verificacion SSL: se deja activada (verify=True) por defecto contra
   el endpoint de Felix. Si el certificado da problemas, es mejor
   arreglarlo en el servidor que desactivar la verificacion.
4) Multi-red: REDES es una lista de dicts {nombre, base_url,
   auth_header}. El ciclo principal repite obtener_estado_actual() /
   comparar_y_notificar() por cada red, con su propio fichero de
   estado (uno por red) para no mezclar cpe_id de una red con otra.
   Si falla la consulta de una red, se avisa al admin y se continua
   con las demas (un fallo de una red no bloquea al resto).
5) Logging de fallos:
   - Todo se registra ademas en LOG_FILE (rotativo, no crece sin limite).
   - Si algo falla (Felix no responde, error inesperado, etc.) se manda
     un aviso a ADMIN_CHAT_ID por Telegram, separado del chat de LOS,
     para que un fallo del bot no se mezcle con los avisos de clientes.
   - Si Telegram mismo esta caido, el aviso de admin fallara tambien:
     en ese caso solo queda registrado en LOG_FILE (por eso conviene
     revisar el log de vez en cuando o vigilarlo con otra herramienta).
"""

import json
import os
import re
import sys
import tempfile
import time
try:
    import fcntl
    _FCNTL_AVAILABLE = True
except ImportError:
    _FCNTL_AVAILABLE = False
import logging
import logging.handlers
import threading
import traceback

import requests

# ---------------------------------------------------------------------------
# CARGA DE .env (antes de CONFIGURACION para que os.environ.get funcione)
# ---------------------------------------------------------------------------
# Si existe un fichero .env junto al script, carga sus variables en el
# entorno. Formato: KEY=value (una por linea, # para comentarios).
# Las variables ya presentes en el entorno tienen prioridad sobre el .env.
def _cargar_dotenv():
    _env_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".env")
    if not os.path.exists(_env_path):
        return
    with open(_env_path) as _f:
        for _linea in _f:
            _linea = _linea.strip()
            if not _linea or _linea.startswith("#") or "=" not in _linea:
                continue
            _clave, _, _valor = _linea.partition("=")
            _clave = _clave.strip()
            _valor = _valor.strip().strip('"').strip("'")
            if _clave and _clave not in os.environ:
                os.environ[_clave] = _valor

_cargar_dotenv()

# ---------------------------------------------------------------------------
# CONFIGURACION
# ---------------------------------------------------------------------------

# Cada red tiene su propia base_url y su propio header de autorizacion
# Basic de Felix. "slug" se usa para nombrar el fichero de estado de
# cada red (last_felix_state_<slug>.json), asi que si lo cambias en
# una red ya en produccion, esa red volvera a pasar por "primera
# ejecucion" (perdera la memoria de estado anterior).
REDES = [
    {
        "nombre": "Moron de la frontera",
        "slug": "moron",
        "base_url": "https://fiberpuebla-flx1.alea-soluciones.com:9199/fiberpuebla",
        "auth_header": os.environ.get("FELIX_AUTH_MORON", "Basic "),
    },
    {
        "nombre": "Fiberalg",
        "slug": "fiberalg",
        # OJO: la URL que te paso aqui tiene "/fiberalg" al final aunque tu
        # solo me diste el dominio suelto. Se lo he anadido porque asi es
        # como esta montada la de Moron ("/fiberpuebla" al final) y como
        # salian el resto de filas en la tabla de la base de datos (cada
        # una terminaba en "/" + su propio slug). Si al probar da 404,
        # revisa este trozo primero.
        "base_url": "https://fiberalg-flx1.alea-soluciones.com:9199/fiberalg",
        "auth_header": os.environ.get("FELIX_AUTH_FIBERALG", "Basic "),
    },
    {
        "nombre": "FiberPlus",
        "slug": "fiberplus",
        # Mismo razonamiento que en Fiberalg: anadido "/fiberplus" al final.
        "base_url": "https://fiberplus-flx1.alea-soluciones.com:9199/fiberplus",
        "auth_header": os.environ.get("FELIX_AUTH_FIBERPLUS", "Basic "),
    },
]

CUSTOMER_ENDPOINT = "/customer_technical_info_report"
PACKAGE_ENDPOINT = "/product_package_group"

FELIX_TIMEOUT = 30  # segundos

# Python (via certifi) trae su propio paquete de certificados raiz,
# separado del sistema operativo, y puede estar desactualizado en
# instalaciones viejas (como Python 3.5 en Debian Stretch). Usamos el
# bundle del sistema, que es el que se actualiza con
# update-ca-certificates, para evitar depender de dos fuentes distintas.
SYSTEM_CA_BUNDLE = "/etc/ssl/certs/ca-certificates.crt"

# Sesiones HTTP reutilizables: evitan abrir una conexion TCP/TLS nueva
# en cada llamada (a Felix o a Telegram), reutilizando el pool de
# conexiones de requests. Mas relevante cuando REDES crece o el bot
# lleva mucho tiempo corriendo. requests.Session() es seguro de usar
# desde varios hilos para peticiones simples como las de este bot
# (monitorizacion y comandos usan hilos distintos).
_felix_session = requests.Session()
_telegram_session = requests.Session()

TELEGRAM_TOKEN = os.environ.get("TELEGRAM_TOKEN", "")
TELEGRAM_CHAT_ID = os.environ.get("TELEGRAM_CHAT_ID", "")

# Chat donde llegan los AVISOS DE FALLO del propio bot (Felix caido,
# excepcion no controlada, etc.), separado del chat de LOS de clientes.
# Puede ser el mismo TELEGRAM_CHAT_ID si prefieres tenerlo todo junto,
# o el chat_id de un admin/grupo tecnico distinto.
ADMIN_CHAT_ID = os.environ.get("ADMIN_CHAT_ID", "")

# Carpeta donde guardamos el ultimo estado online/offline conocido por
# cpe_id (SN de la ONT), un fichero por red, para detectar transiciones
# entre ejecuciones sin mezclar cpe_id de una red con otra.
BASE_DIR = os.path.dirname(os.path.abspath(__file__))


def state_file_de_red(slug):
    return os.path.join(BASE_DIR, "last_felix_state_{}.json".format(slug))

# Log de fichero en el servidor, con rotacion para que no crezca sin
# limite (5 ficheros de 1 MB cada uno, se van rotando).
LOG_FILE = os.path.join(
    os.path.dirname(os.path.abspath(__file__)), "los_bot.log"
)

# Regex para extraer la referencia de caja desde external_plant, igual
# que hace felixsncaj() en actualizar_onts.php. Si no matchea, se
# muestra el external_plant tal cual.
RE_CAJA = re.compile(r"\d{1,2}-\D{1,2}\d{0,2}-\D\d{0,2}-\d{2,4}-\d{1,2}")

# Fichero donde se guarda el ultimo update_id de Telegram procesado,
# para el modo de escucha de comandos (long polling con getUpdates).
OFFSET_FILE = os.path.join(BASE_DIR, "telegram_offset.txt")

# Lock de ejecucion unica: evita que se lancen dos instancias del bot a
# la vez (p.ej. si se relanza con nohup sin matar la anterior), que es
# lo que provoca el "409 Conflict" de Telegram en getUpdates y una
# carrera al escribir los ficheros de estado. Se usa fcntl.flock, un
# lock a nivel de kernel atado al proceso: si el proceso muere por
# cualquier motivo (incluido kill -9), Linux lo libera solo, asi que no
# hay riesgo de que quede un "lock fantasma" bloqueando arranques
# futuros (a diferencia de comprobar un PID guardado en un fichero).
LOCK_FILE = os.path.join(BASE_DIR, "bot_felix.lock")
_lock_fd = None  # hay que mantener el file descriptor abierto: si se
                  # cierra (o lo recoge el garbage collector), el lock
                  # se libera aunque el proceso siga vivo.


def adquirir_lock_ejecucion():
    """
    Intenta tomar el lock exclusivo de ejecucion. Devuelve True si lo
    consigue (no hay otra instancia corriendo) o False si ya hay otra
    instancia con el lock tomado.
    """
    global _lock_fd
    if not _FCNTL_AVAILABLE:
        log.warning("fcntl no disponible (Windows): lock de ejecucion unica desactivado.")
        return True
    _lock_fd = open(LOCK_FILE, "w")
    try:
        fcntl.flock(_lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
    except (IOError, OSError):
        _lock_fd.close()
        _lock_fd = None
        return False
    # Dejamos el PID escrito solo a efectos informativos (para que un
    # humano mirando el fichero sepa quien tiene el lock); el lock en
    # si no depende de este contenido, depende del flock del fd.
    _lock_fd.write(str(os.getpid()))
    _lock_fd.flush()
    return True

# Solo se responde a comandos que lleguen desde estos chats (el de
# clientes y el de admin). Cualquier otro chat se ignora, para que no
# cualquiera pueda consultar datos de clientes escribiendole al bot.
CHATS_PERMITIDOS = {c for c in (TELEGRAM_CHAT_ID, ADMIN_CHAT_ID) if c}

# --- MODO DE EJECUCION ---
# Por defecto (sin argumentos) el script se queda corriendo para
# siempre en segundo plano, haciendo dos cosas a la vez:
#   1) Monitorizacion de LOS cada SLEEP_SECONDS (hilo principal).
#   2) Escucha de comandos de Telegram /estado /offline /buscar (hilo aparte).
# Ya no hace falta cron ni systemd para tenerlo funcionando: basta con
# arrancarlo una vez con nohup (ver el bloque final del fichero) y se
# mantiene vivo el solo, reintentando si algo falla puntualmente.
SLEEP_SECONDS = 300

# Numero de ciclos consecutivos que una ONT debe estar offline antes de
# notificar. Con 1 se notifica en el primer ciclo (puede dar falsos positivos
# por glitches puntuales de Felix). Con 2 se espera un ciclo extra de
# confirmacion (~SLEEP_SECONDS segundos mas) antes de avisar.
LOS_CICLOS_CONFIRMACION = 2

# Contadores en memoria de ciclos consecutivos offline por ONT.
# Clave: (slug_red, cpe_id) -> numero de ciclos seguidos offline.
# Se resetea al reiniciar el bot (correcto: volvemos a exigir N ciclos).
_ciclos_offline = {}

log = logging.getLogger("los_bot")
log.setLevel(logging.INFO)

_formato = logging.Formatter("%(asctime)s [%(levelname)s] %(message)s")

_handler_consola = logging.StreamHandler()
_handler_consola.setFormatter(_formato)
log.addHandler(_handler_consola)

# maxBytes=1MB, quedandonos con los ultimos 5 ficheros rotados
_handler_fichero = logging.handlers.RotatingFileHandler(
    LOG_FILE, maxBytes=1024 * 1024, backupCount=5
)
_handler_fichero.setFormatter(_formato)
log.addHandler(_handler_fichero)


# ---------------------------------------------------------------------------
# ESTADO LOCAL (ultimo online/offline conocido por cpe_id)
# ---------------------------------------------------------------------------

def _escribir_texto_atomico(path, texto):
    """
    Escribe 'texto' en 'path' de forma atomica: escribe primero a un
    fichero temporal en el MISMO directorio (para que este en el mismo
    filesystem) y luego hace os.replace(), que en POSIX es una
    operacion atomica a nivel de sistema de ficheros. Asi evitamos que
    un kill -9 / corte de luz / OOM a mitad de escritura deje el
    fichero truncado o a medias: o se ve la version vieja completa, o
    la nueva completa, nunca algo intermedio.
    fsync() antes del rename asegura que los datos ya estan en disco
    (y no solo en el buffer del SO) antes de que el rename los deje
    visibles.
    """
    directorio = os.path.dirname(path) or "."
    fd, tmp_path = tempfile.mkstemp(dir=directorio, prefix=".tmp_", suffix=".swap")
    try:
        with os.fdopen(fd, "w") as f:
            f.write(texto)
            f.flush()
            os.fsync(f.fileno())
        os.replace(tmp_path, path)
    except Exception:
        # Si algo falla a medio camino, no dejamos el temporal huerfano.
        try:
            os.remove(tmp_path)
        except OSError:
            pass
        raise


def cargar_estado_anterior(slug):
    state_file = state_file_de_red(slug)
    if not os.path.exists(state_file):
        # Caso normal: primera ejecucion de esta red, todavia no hay
        # fichero de estado. No es un fallo, no se loguea como tal.
        return {}
    try:
        with open(state_file, "r") as f:
            return json.load(f)
    except (ValueError, IOError) as e:
        # El fichero SI existe pero no se puede leer o parsear (JSON
        # corrupto, permisos, etc.). A diferencia del caso anterior,
        # esto es una anomalia real: si lo tratamos igual y devolvemos
        # {} en silencio, el bot lo interpreta como "primera ejecucion"
        # y no notificara las transiciones de este ciclo. Avisamos al
        # admin para que se sepa que el estado guardado se ha perdido.
        enviar_alerta_admin(
            "No se pudo leer el fichero de estado '{}' (posible JSON "
            "corrupto): {}\nSe trata como si no hubiera estado previo "
            "para esta red (no se notificaran transiciones en este "
            "ciclo).".format(state_file, e),
            clave="estado_corrupto_{}".format(slug),
        )
        return {}


def guardar_estado(slug, estado):
    _escribir_texto_atomico(state_file_de_red(slug), json.dumps(estado))


# ---------------------------------------------------------------------------
# CONSULTAS A FELIX
# ---------------------------------------------------------------------------

def consultar_felix(red, endpoint):
    url = red["base_url"] + endpoint
    headers = {
        "authorization": red["auth_header"],
        "cache-control": "no-cache",
    }
    verify = SYSTEM_CA_BUNDLE if os.path.exists(SYSTEM_CA_BUNDLE) else True
    r = _felix_session.get(url, headers=headers, timeout=FELIX_TIMEOUT, verify=verify)
    r.raise_for_status()
    return r.json()


def construir_mapa_cajas(red):
    """cpe_id -> referencia de caja (extraida de external_plant)."""
    data = consultar_felix(red, PACKAGE_ENDPOINT)
    mapa = {}
    for grupo in data:
        external_plant = (grupo.get("external_plant") or "").strip().upper()
        match = RE_CAJA.search(external_plant)
        referencia = match.group(0) if match else (external_plant or "SIN REFERENCIA")
        for paquete in grupo.get("product_packages", []):
            cpe_id = paquete.get("cpe_id")
            if cpe_id:
                mapa[cpe_id] = referencia
    return mapa


def _parsear_online(valor):
    """
    Interpreta el campo "online" de Felix de forma robusta. Si viene
    como bool o int (0/1) nativo, bool(valor) ya es correcto. Pero si
    alguna vez llega como string (p.ej. "0" o "false"), bool("0") da
    True por ser un string no vacio, que es el resultado contrario al
    que se quiere. No hemos visto ese caso en respuestas reales de
    Felix, pero esto evita que sea un fallo silencioso si ocurriera.
    Si el campo no viene (None), se interpreta como offline.
    """
    if isinstance(valor, str):
        return valor.strip().lower() not in ("", "0", "false", "no", "off")
    return bool(valor)


def construir_pon(ppg):
    olt = ppg.get("olt_id") or ""
    if not olt:
        # Felix borra olt_id cuando la ONT cae pero deja frame/slot/port
        # en 0 en vez de None, asi que sin este atajo devolveríamos
        # "N/D 0/0/0". Devolvemos solo "N/D" para que quede claro que
        # el dato no está disponible, no que la ONT esté en el PON 0/0/0.
        return "N/D"
    frame = ppg.get("frame")
    slot = ppg.get("slot")
    port = ppg.get("port")
    if frame is None or slot is None or port is None:
        return olt
    return "{} {}/{}/{}".format(olt, frame, slot, port)


def obtener_estado_actual(red, mapa_cajas):
    """
    cpe_id -> {online: bool, activa: bool, cliente, caja, pon}
    Incluye ppgs activos E inactivos (siempre que tengan cpe_id).
    Las ONTs con active=False se marcan con activa=False para que los
    mensajes y comandos puedan distinguirlas visualmente.
    """
    data = consultar_felix(red, CUSTOMER_ENDPOINT)
    estado = {}
    for cliente in data:
        nombre_cliente = cliente.get("full_name") or "N/D"
        for ppg in cliente.get("ppgs", []):
            cpe_id = ppg.get("cpe_id")
            if not cpe_id:
                continue
            activa = bool(ppg.get("active", True))
            estado[cpe_id] = {
                "online": _parsear_online(ppg.get("online")),
                "activa": activa,
                "cliente": nombre_cliente,
                "caja": mapa_cajas.get(cpe_id, "N/D"),
                "pon": construir_pon(ppg),
            }
    return estado


def obtener_estado_todas_redes():
    """
    Consulta Felix EN VIVO (no usa el estado guardado en disco) para
    todas las redes configuradas. Devuelve una tupla:
      (datos, redes_fallidas)
    donde 'datos' es una lista de tuplas (nombre_red, cpe_id, info)
    y 'redes_fallidas' es una lista de nombres de redes que no
    respondieron. Se usa desde los comandos de Telegram, que necesitan
    el estado actual real, no el ultimo guardado por cron.
    """
    resultado = []
    redes_fallidas = []
    for red in REDES:
        try:
            mapa_cajas = construir_mapa_cajas(red)
            estado = obtener_estado_actual(red, mapa_cajas)
        except requests.RequestException as e:
            enviar_alerta_admin(
                "Error consultando la API de Felix para la red '{}' "
                "al atender un comando de Telegram.\nDetalle: {}".format(
                    red["nombre"], e
                ),
                clave="felix_error_comando_{}".format(red["slug"]),
            )
            redes_fallidas.append(red["nombre"])
            continue
        except Exception as e:
            enviar_alerta_admin(
                "Error inesperado consultando/procesando Felix para la "
                "red '{}' al atender un comando de Telegram.\n"
                "Detalle: {}\n```\n{}\n```".format(
                    red["nombre"], e, traceback.format_exc()
                ),
                clave="felix_error_comando_inesperado_{}".format(red["slug"]),
            )
            redes_fallidas.append(red["nombre"])
            continue
        for cpe_id, info in estado.items():
            resultado.append((red["nombre"], cpe_id, info))
    return resultado, redes_fallidas


# ---------------------------------------------------------------------------
# MENSAJE Y ENVIO
# ---------------------------------------------------------------------------

class _RateLimited(Exception):
    """Señal interna: Telegram respondio 429, con el retry_after indicado."""
    def __init__(self, retry_after):
        super(_RateLimited, self).__init__("429 rate limited")
        self.retry_after = retry_after

_RE_MARKDOWN_ESPECIALES = re.compile(r"([_*\[\]()~`#+\-=|{}.!\\>])")


def _escapar_markdown(texto):
    """
    Escapa los caracteres especiales de MarkdownV2 de Telegram para que
    contenido dinamico (nombres de cliente, referencias de caja, etc.)
    no rompa el formato del mensaje. Se aplica solo al construir el texto
    a mostrar, no a los valores guardados en estado_actual (para no
    romper la busqueda de /buscar, que compara sobre el valor crudo).
    """
    if texto is None:
        return texto
    return _RE_MARKDOWN_ESPECIALES.sub(r"\\\1", str(texto))


def construir_mensaje(red, cpe_id, info, tipo):
    if tipo == "offline":
        icono, titulo = "\U0001F534", "LOS detectado"
    else:
        icono, titulo = "\U0001F7E2", "Recuperado"

    fecha_str = time.strftime("%Y-%m-%d %H:%M:%S")

    aviso_inactiva = (
        "\n\u26A0\uFE0F _Marcada como inactiva en Felix_"
        if not info.get("activa", True) else ""
    )

    return (
        "{icono} *{titulo}*\n"
        "*Red:* {red}\n"
        "*Cliente:* {cliente}\n"
        "*SN ONT:* `{sn}`\n"
        "*Caja:* {caja}\n"
        "*PON:* {pon}\n"
        "*Hora:* {fecha}"
        "{aviso}"
    ).format(
        icono=icono,
        titulo=titulo,
        red=_escapar_markdown(red["nombre"]),
        cliente=_escapar_markdown(info["cliente"]),
        sn=cpe_id,
        caja=_escapar_markdown(info["caja"]),
        pon=_escapar_markdown(info["pon"]),
        fecha=_escapar_markdown(fecha_str),
        aviso=aviso_inactiva,
    )


TELEGRAM_MAX_CHARS = 4000  # limite real de Telegram es 4096, dejamos margen

# Telegram limita aprox. 20 msg/min en grupos (~1 cada 3s) y ~30 msg/s
# en total contra la API. Con este intervalo minimo entre envios nos
# quedamos comodos por debajo de ambos limites incluso en una caida
# masiva con decenas de notificaciones seguidas.
TELEGRAM_MIN_INTERVALO = 3.5  # segundos entre mensajes salientes
_ultimo_envio_lock = threading.Lock()
_ultimo_envio_ts = [0.0]

# Si aun asi Telegram responde 429, cuantos reintentos hacemos
# respetando el "retry_after" que el propio Telegram indica.
TELEGRAM_MAX_REINTENTOS_429 = 3


def _esperar_turno_envio():
    """Throttle global: fuerza un hueco minimo entre envios salientes,
    compartido por todos los hilos/llamadas (monitorizacion y comandos)."""
    with _ultimo_envio_lock:
        ahora = time.time()
        espera = _ultimo_envio_ts[0] + TELEGRAM_MIN_INTERVALO - ahora
        if espera > 0:
            time.sleep(espera)
        _ultimo_envio_ts[0] = time.time()


def enviar_telegram(texto, chat_id=None):
    destino = chat_id or TELEGRAM_CHAT_ID
    url = "https://api.telegram.org/bot{}/sendMessage".format(TELEGRAM_TOKEN)

    # Mensajes demasiado largos (p.ej. un traceback completo) Telegram
    # los rechaza con otro error distinto ("message is too long").
    # Truncamos para evitar ese segundo motivo de fallo.
    if len(texto) > TELEGRAM_MAX_CHARS:
        texto = texto[:TELEGRAM_MAX_CHARS] + "\n... (mensaje truncado, ver LOG_FILE)"

    def _post(payload):
        _esperar_turno_envio()
        r = _telegram_session.post(url, json=payload, timeout=10)
        if r.status_code == 429:
            # Telegram indica cuantos segundos esperar en el body:
            # {"ok": false, "error_code": 429, "parameters": {"retry_after": N}}
            try:
                retry_after = r.json().get("parameters", {}).get("retry_after", 3)
            except ValueError:
                retry_after = 3
            raise _RateLimited(retry_after)
        r.raise_for_status()

    payload = {
        "chat_id": destino,
        "text": texto,
        "parse_mode": "MarkdownV2",
    }

    def _post_con_reintentos_429(payload_a_enviar):
        for intento in range(TELEGRAM_MAX_REINTENTOS_429 + 1):
            try:
                _post(payload_a_enviar)
                return
            except _RateLimited as e:
                if intento == TELEGRAM_MAX_REINTENTOS_429:
                    raise
                log.warning(
                    "Rate limit (429) de Telegram (chat_id=%s), esperando %ss "
                    "y reintentando (intento %d/%d).",
                    destino, e.retry_after, intento + 1, TELEGRAM_MAX_REINTENTOS_429,
                )
                time.sleep(e.retry_after)

    try:
        _post_con_reintentos_429(payload)
        return
    except _RateLimited:
        log.error(
            "Se agotaron los reintentos por rate limit (429) de Telegram "
            "(chat_id=%s), se descarta este mensaje.", destino,
        )
        return
    except requests.RequestException as e:
        cuerpo = getattr(e.response, "text", "")
        # Fallo tipico: el texto tiene "_", "*" o "`" sueltos (por
        # ejemplo un traceback con nombres de fichero) y Telegram no
        # puede parsear el Markdown. En vez de perder el aviso,
        # reintentamos como texto plano.
        if getattr(e.response, "status_code", None) == 400 and "parse" in cuerpo.lower():
            log.warning(
                "Fallo el envio con Markdown (chat_id=%s), reintentando en texto plano. "
                "Motivo Telegram: %s", destino, cuerpo
            )
            try:
                _post_con_reintentos_429({"chat_id": destino, "text": texto})
                log.info("Reintento en texto plano OK (chat_id=%s).", destino)
                return
            except _RateLimited:
                log.error(
                    "Se agotaron los reintentos por rate limit (429) de Telegram "
                    "en texto plano (chat_id=%s), se descarta este mensaje.", destino,
                )
                return
            except requests.RequestException as e2:
                cuerpo2 = getattr(e2.response, "text", "")
                log.error(
                    "Error enviando a Telegram en texto plano (chat_id=%s): %s | %s",
                    destino, e2, cuerpo2,
                )
                return

        log.error(
            "Error enviando a Telegram (chat_id=%s): %s | %s",
            destino, e, cuerpo,
        )


# Si el mismo tipo de fallo se repite (p.ej. Felix caido durante varios
# ciclos seguidos), no queremos mandar un aviso identico cada
# SLEEP_SECONDS: es ruido. Agrupamos los avisos por "clave" (el
# llamador indica de que tipo de fallo se trata) y solo dejamos pasar
# uno por clave cada ALERTA_ADMIN_COOLDOWN segundos; el resto se sigue
# logueando en LOG_FILE, solo se suprime el envio a Telegram.
ALERTA_ADMIN_COOLDOWN = 1800  # 30 min
_alertas_admin_lock = threading.Lock()
_ultimas_alertas_admin = {}  # clave -> timestamp del ultimo envio a Telegram


def enviar_alerta_admin(mensaje, clave=None):
    """
    Avisa al chat de admin de que algo ha fallado. Siempre queda
    registrado en LOG_FILE independientemente de si Telegram responde
    o si el envio se suprime por repeticion.

    'clave' agrupa avisos del mismo tipo de fallo (p.ej. el slug de la
    red, o "excepcion_no_controlada"). Si no se indica, se usa el
    propio mensaje como clave (solo deduplica avisos identicos).
    """
    log.error(mensaje)

    clave_agrupacion = clave or mensaje
    ahora = time.time()
    with _alertas_admin_lock:
        ultimo_envio = _ultimas_alertas_admin.get(clave_agrupacion)
        if ultimo_envio is not None and (ahora - ultimo_envio) < ALERTA_ADMIN_COOLDOWN:
            log.info(
                "Aviso a admin suprimido por repeticion (clave='%s', "
                "ultimo envio hace %.0fs).", clave_agrupacion, ahora - ultimo_envio,
            )
            return
        _ultimas_alertas_admin[clave_agrupacion] = ahora

    texto = "\u26A0\uFE0F *Fallo en telegram\\_los\\_bot*\n{}".format(mensaje)
    enviar_telegram(texto, chat_id=ADMIN_CHAT_ID)


# ---------------------------------------------------------------------------
# COMANDOS DE TELEGRAM (consulta en vivo, requiere modo "comandos")
# ---------------------------------------------------------------------------

MAX_LINEAS_LISTADO = 30  # limite para no mandar mensajes gigantes


def comando_estado():
    datos, fallidas = obtener_estado_todas_redes()
    if not datos:
        return "No se ha podido consultar Felix ahora mismo\. Revisa el log\."

    # Contadores por red, separando LOS reales de ONTs inactivas en Felix
    # (una ONT inactiva aparece como online=False pero no es una caida real)
    stats_por_red = {}
    for nombre_red, _, info in datos:
        if nombre_red not in stats_por_red:
            stats_por_red[nombre_red] = {"online": 0, "los": 0, "inactivas_los": 0}
        if info["online"]:
            stats_por_red[nombre_red]["online"] += 1
        elif not info.get("activa", True):
            stats_por_red[nombre_red]["inactivas_los"] += 1
        else:
            stats_por_red[nombre_red]["los"] += 1

    total_online    = sum(s["online"]        for s in stats_por_red.values())
    total_los       = sum(s["los"]           for s in stats_por_red.values())
    total_inactivas = sum(s["inactivas_los"] for s in stats_por_red.values())

    lineas = ["*Estado actual*"]
    for nombre_red, s in stats_por_red.items():
        partes = "\U0001F7E2 {}".format(s["online"])
        if s["los"] > 0:
            partes += "  \U0001F534 {} LOS".format(s["los"])
        if s["inactivas_los"] > 0:
            partes += "  \u26A0\uFE0F {} inactiva{}".format(
                s["inactivas_los"], "s" if s["inactivas_los"] > 1 else ""
            )
        lineas.append("{}: {}".format(_escapar_markdown(nombre_red), partes))

    los_str = str(total_los)
    if total_inactivas > 0:
        los_str += " \\(\\+{} inactiva{} en Felix\\)".format(
            total_inactivas, "s" if total_inactivas > 1 else ""
        )
    lineas.append(
        "\nTotal: {} ONTs  \U0001F7E2 {}  \U0001F534 {}".format(
            len(datos), total_online, los_str
        )
    )

    if fallidas:
        lineas.append(
            "\n\u26A0\uFE0F Datos incompletos, sin respuesta de: {}".format(
                ", ".join(_escapar_markdown(r) for r in fallidas)
            )
        )

    return "\n".join(lineas)


def comando_offline():
    datos, fallidas = obtener_estado_todas_redes()
    if not datos:
        return "No se ha podido consultar Felix ahora mismo\. Revisa el log\."

    caidas = [(red, cpe_id, info) for red, cpe_id, info in datos if not info["online"]]
    if not caidas:
        aviso_fallidas = (
            "\n\u26A0\uFE0F Sin datos de: {}".format(", ".join(_escapar_markdown(r) for r in fallidas))
            if fallidas else ""
        )
        return "\u2705 No hay ninguna ONT en LOS ahora mismo\." + aviso_fallidas

    lineas = ["*ONTs en LOS ahora mismo \\({}\\):*".format(len(caidas))]
    for red, cpe_id, info in caidas[:MAX_LINEAS_LISTADO]:
        aviso = " \u26A0\uFE0F" if not info.get("activa", True) else ""
        lineas.append(
            "\\- `{sn}`{aviso} \\| {cliente} \\| caja {caja} \\| {red}".format(
                sn=cpe_id,
                aviso=aviso,
                cliente=_escapar_markdown(info["cliente"]),
                caja=_escapar_markdown(info["caja"]),
                red=_escapar_markdown(red),
            )
        )
    if len(caidas) > MAX_LINEAS_LISTADO:
        lineas.append("\\.\\.\\. y {} mas".format(len(caidas) - MAX_LINEAS_LISTADO))
    if fallidas:
        lineas.append(
            "\n\u26A0\uFE0F Sin datos de: {}".format(", ".join(_escapar_markdown(r) for r in fallidas))
        )
    return "\n".join(lineas)


def comando_buscar(texto_busqueda):
    if not texto_busqueda:
        return "Uso: /buscar <SN o nombre de cliente\\>"

    datos, fallidas = obtener_estado_todas_redes()
    if not datos:
        return "No se ha podido consultar Felix ahora mismo\. Revisa el log\."

    filtro = texto_busqueda.lower()
    coincidencias = [
        (red, cpe_id, info) for red, cpe_id, info in datos
        if filtro in cpe_id.lower() or filtro in info["cliente"].lower()
    ]
    if not coincidencias:
        aviso_fallidas = (
            " \\(sin datos de: {}\\)".format(", ".join(_escapar_markdown(r) for r in fallidas)) if fallidas else ""
        )
        return "No se encontraron ONTs que coincidan con '{}'\.{}".format(
            _escapar_markdown(texto_busqueda), aviso_fallidas
        )

    lineas = []
    for red, cpe_id, info in coincidencias[:MAX_LINEAS_LISTADO]:
        icono = "\U0001F7E2" if info["online"] else "\U0001F534"
        aviso = " \u26A0\uFE0F" if not info.get("activa", True) else ""
        lineas.append(
            "{icono} `{sn}`{aviso} \\| {cliente} \\| caja {caja} \\| {pon} \\| {red}".format(
                icono=icono, sn=cpe_id, aviso=aviso,
                cliente=_escapar_markdown(info["cliente"]),
                caja=_escapar_markdown(info["caja"]),
                pon=_escapar_markdown(info["pon"]),
                red=_escapar_markdown(red),
            )
        )
    if len(coincidencias) > MAX_LINEAS_LISTADO:
        lineas.append("\\.\\.\\. y {} mas \\(afina la busqueda\\)".format(len(coincidencias) - MAX_LINEAS_LISTADO))
    if fallidas:
        lineas.append(
            "\n\u26A0\uFE0F Sin datos de: {}".format(", ".join(_escapar_markdown(r) for r in fallidas))
        )
    return "\n".join(lineas)


def procesar_comando(texto, chat_id):
    partes = texto.strip().split(None, 1)
    comando = partes[0].lower().split("@")[0]  # quita @NombreDelBot si viene incluido
    argumento = partes[1].strip() if len(partes) > 1 else ""

    if comando == "/estado":
        respuesta = comando_estado()
    elif comando == "/offline":
        respuesta = comando_offline()
    elif comando == "/buscar":
        respuesta = comando_buscar(argumento)
    elif comando in ("/start", "/ayuda", "/help"):
        respuesta = (
            "Comandos disponibles:\n"
            "/estado \\- resumen de ONTs online/LOS\n"
            "/offline \\- lista de ONTs caidas ahora mismo\n"
            "/buscar <SN o cliente\\> \\- estado de una ONT concreta"
        )
    else:
        respuesta = "Comando no reconocido\\. Prueba /ayuda\\."

    enviar_telegram(respuesta, chat_id=chat_id)


# ---------------------------------------------------------------------------
# CICLO PRINCIPAL
# ---------------------------------------------------------------------------

def procesar_red(red, silencioso=False):
    """Procesa una unica red: consulta Felix, compara y notifica.
    Con silencioso=True solo guarda el estado actual sin enviar nada."""
    try:
        mapa_cajas = construir_mapa_cajas(red)
        estado_actual = obtener_estado_actual(red, mapa_cajas)
    except requests.RequestException as e:
        enviar_alerta_admin(
            "Error consultando la API de Felix para la red '{}', se aborta "
            "esta red en este ciclo.\nDetalle: {}".format(red["nombre"], e),
            clave="felix_error_{}".format(red["slug"]),
        )
        return
    except Exception as e:
        # Cualquier otro fallo inesperado (JSON con forma distinta a la
        # esperada -> KeyError/AttributeError, etc.): lo tratamos igual
        # que un fallo de red en vez de dejarlo subir sin avisar a
        # nadie, para no perder esta red en silencio.
        enviar_alerta_admin(
            "Error inesperado consultando/procesando Felix para la red "
            "'{}', se aborta esta red en este ciclo.\nDetalle: {}\n"
            "```\n{}\n```".format(red["nombre"], e, traceback.format_exc()),
            clave="felix_error_inesperado_{}".format(red["slug"]),
        )
        return

    estado_anterior = cargar_estado_anterior(red["slug"])
    primera_ejecucion = not estado_anterior or silencioso

    notificados = 0
    nuevo_estado = {}
    for cpe_id, info in estado_actual.items():
        online_ahora = info["online"]

        # Compatibilidad con el formato antiguo de estado (bool puro) y
        # con el nuevo formato enriquecido (dict con online/pon/caja/cliente).
        anterior_raw = estado_anterior.get(cpe_id)
        if isinstance(anterior_raw, dict):
            online_antes  = anterior_raw.get("online")
            pon_guardado  = anterior_raw.get("pon",  "N/D")
            caja_guardada = anterior_raw.get("caja", "N/D")
        elif anterior_raw is not None:
            online_antes  = bool(anterior_raw)  # formato antiguo: solo bool
            pon_guardado  = "N/D"
            caja_guardada = "N/D"
        else:
            online_antes  = None
            pon_guardado  = "N/D"
            caja_guardada = "N/D"

        # Felix borra olt_id/frame/slot/port cuando una ONT cae, asi que
        # info["pon"] llega como "N/D" justo en el momento del LOS.
        # Recuperamos el ultimo valor bueno guardado en el ciclo anterior
        # para que el aviso muestre la ubicacion real de la caja.
        pon_efectivo  = info["pon"]  if info["pon"]  != "N/D" else pon_guardado
        caja_efectiva = info["caja"] if info["caja"] != "N/D" else caja_guardada

        # Guardar estado enriquecido. Usamos pon_efectivo/caja_efectiva
        # (el mejor valor conocido) para que el PON/caja se conserve
        # aunque la ONT lleve varios ciclos seguidos en LOS y Felix
        # siga sin devolver los campos en cada consulta.
        nuevo_estado[cpe_id] = {
            "online":  online_ahora,
            "pon":     pon_efectivo,
            "caja":    caja_efectiva,
            "cliente": info["cliente"],
        }

        if primera_ejecucion or online_antes is None:
            continue

        clave = (red["slug"], cpe_id)
        activa = info.get("activa", True)

        if not activa:
            # Inactiva en Felix: ignorar transiciones, limpiar contador.
            if online_antes != online_ahora:
                tipo = "online" if online_ahora else "offline"
                log.info("[%s] Transicion ignorada (inactiva) sn=%s -> %s",
                         red["nombre"], cpe_id, tipo)
            _ciclos_offline.pop(clave, None)
            continue

        if not online_ahora:
            if online_antes:
                # Transicion online → offline detectada en este ciclo.
                _ciclos_offline[clave] = 1
            elif clave in _ciclos_offline:
                # Sigue offline y ya lo rastreabamos: acumular ciclos.
                _ciclos_offline[clave] += 1
            # else: ya estaba offline al arrancar el bot → ignorar.

            ciclos = _ciclos_offline.get(clave, 0)
            if ciclos == LOS_CICLOS_CONFIRMACION:
                info_msg = dict(info)
                info_msg["pon"]  = pon_efectivo
                info_msg["caja"] = caja_efectiva
                texto = construir_mensaje(red, cpe_id, info_msg, "offline")
                enviar_telegram(texto)
                notificados += 1
                log.info("[%s] LOS notificado sn=%s (ciclo %d/%d).",
                         red["nombre"], cpe_id, ciclos, LOS_CICLOS_CONFIRMACION)
            elif 0 < ciclos < LOS_CICLOS_CONFIRMACION:
                log.info("[%s] LOS tentativo sn=%s (ciclo %d/%d), esperando confirmacion.",
                         red["nombre"], cpe_id, ciclos, LOS_CICLOS_CONFIRMACION)
        else:
            # ONT online ahora
            if not online_antes:
                # Recuperacion
                ciclos_fue = _ciclos_offline.pop(clave, 0)
                info_msg = dict(info)
                info_msg["pon"]  = pon_efectivo
                info_msg["caja"] = caja_efectiva
                if ciclos_fue >= LOS_CICLOS_CONFIRMACION:
                    # El LOS habia sido notificado: avisar de la recuperacion.
                    texto = construir_mensaje(red, cpe_id, info_msg, "online")
                    enviar_telegram(texto)
                    notificados += 1
                    log.info("[%s] Recuperacion notificada sn=%s.", red["nombre"], cpe_id)
                else:
                    # LOS no habia llegado al umbral: recuperacion silenciosa.
                    log.info("[%s] Recuperacion silenciosa sn=%s (LOS no confirmado, %d/%d ciclos).",
                             red["nombre"], cpe_id, ciclos_fue, LOS_CICLOS_CONFIRMACION)

    guardar_estado(red["slug"], nuevo_estado)

    if primera_ejecucion:
        motivo = "Arranque silencioso" if silencioso else "Primera ejecucion"
        log.info(
            "[%s] %s: estado base guardado (%d ONTs), sin notificaciones.",
            red["nombre"], motivo, len(nuevo_estado),
        )
    else:
        log.info(
            "[%s] Ciclo completado: %d ONTs revisadas, %d notificaciones enviadas.",
            red["nombre"], len(nuevo_estado), notificados,
        )


def procesar_una_vez():
    """Recorre todas las redes configuradas en REDES."""
    for red in REDES:
        procesar_red(red)


# ---------------------------------------------------------------------------
# MODO COMANDOS (long polling, proceso aparte del cron de monitorizacion)
# ---------------------------------------------------------------------------

def cargar_offset():
    if not os.path.exists(OFFSET_FILE):
        return None
    try:
        with open(OFFSET_FILE, "r") as f:
            return int(f.read().strip())
    except (ValueError, IOError):
        return None


def guardar_offset(offset):
    _escribir_texto_atomico(OFFSET_FILE, str(offset))


def escuchar_comandos():
    """
    Bucle infinito de long polling contra getUpdates de Telegram.
    Pensado para correr como proceso aparte (systemd, screen, nohup...),
    NO por cron, ya que necesita quedarse escuchando de forma continua.
    """
    log.info("Escuchando comandos de Telegram (long polling)...")
    offset = cargar_offset()
    url = "https://api.telegram.org/bot{}/getUpdates".format(TELEGRAM_TOKEN)

    while True:
        try:
            params = {"timeout": 30}
            if offset is not None:
                params["offset"] = offset
            r = _telegram_session.get(url, params=params, timeout=40)
            r.raise_for_status()
            data = r.json()

            for update in data.get("result", []):
                offset = update["update_id"] + 1
                guardar_offset(offset)

                mensaje = update.get("message") or update.get("edited_message")
                if not mensaje:
                    continue
                texto = mensaje.get("text", "")
                chat_id = str(mensaje.get("chat", {}).get("id", ""))

                if not texto.startswith("/"):
                    continue
                if chat_id not in CHATS_PERMITIDOS:
                    log.info("Comando ignorado de chat no autorizado: %s", chat_id)
                    continue

                try:
                    procesar_comando(texto, chat_id)
                except Exception:
                    log.error(
                        "Error procesando comando '%s':\n%s",
                        texto, traceback.format_exc(),
                    )
        except requests.RequestException as e:
            log.error("Error en long polling de Telegram: %s", e)
            time.sleep(5)
        except Exception:
            log.error(
                "Excepcion no controlada en escuchar_comandos:\n%s",
                traceback.format_exc(),
            )
            time.sleep(5)


def bucle_monitorizacion():
    """Hilo de monitorizacion de LOS: se repite cada SLEEP_SECONDS para siempre."""
    log.info("Iniciando monitorizacion en bucle, intervalo=%ss", SLEEP_SECONDS)
    # Pasada inicial silenciosa: actualiza el estado base sin enviar
    # notificaciones, para que un reinicio del bot no dispare alertas
    # masivas por ONTs que ya estaban offline antes de arrancar.
    log.info("Pasada inicial silenciosa: actualizando estado base sin notificar...")
    for red in REDES:
        try:
            procesar_red(red, silencioso=True)
        except Exception:
            log.error(
                "Error en pasada inicial silenciosa de '%s':\n%s",
                red["nombre"], traceback.format_exc(),
            )
    log.info("Pasada inicial completada. Iniciando monitorizacion normal.")
    while True:
        try:
            procesar_una_vez()
        except Exception:
            enviar_alerta_admin(
                "Excepcion no controlada en el ciclo de monitorizacion:\n"
                "```\n{}\n```".format(traceback.format_exc()),
                clave="excepcion_no_controlada_monitorizacion",
            )
        time.sleep(SLEEP_SECONDS)


if __name__ == "__main__":
    # --- USO NORMAL: sin argumentos ---
    # El script se queda corriendo para siempre en primer plano,
    # haciendo monitorizacion de LOS (hilo principal) y escuchando
    # comandos de Telegram (hilo aparte) a la vez. Para tenerlo
    # funcionando en segundo plano en el servidor:
    #
    #   nohup python3 bot_felix.py > bot_felix.log 2>&1 &
    #
    # Eso lo deja corriendo aunque cierres la sesion SSH. Para pararlo:
    #   ps aux | grep telegram_los_bot   (busca el PID)
    #   kill <PID>
    #
    # --- USO OPCIONAL (depuracion / testing) ---
    #   python3 telegram_los_bot.py una-vez     -> una sola pasada de LOS y termina
    #   python3 telegram_los_bot.py comandos    -> solo escucha comandos, sin monitorizar
    # Validacion de configuracion critica antes de arrancar
    _config_ok = True
    for _var, _val in [("TELEGRAM_TOKEN", TELEGRAM_TOKEN),
                       ("TELEGRAM_CHAT_ID", TELEGRAM_CHAT_ID),
                       ("ADMIN_CHAT_ID", ADMIN_CHAT_ID)]:
        if not _val:
            log.error("CONFIGURACION INCOMPLETA: '%s' esta vacio. "
                      "Defínelo en el fichero .env o como variable de entorno.", _var)
            _config_ok = False
    if not _config_ok:
        log.error("El bot no puede arrancar con la configuracion incompleta. "
                  "Crea un fichero .env junto al script con el contenido:\n"
                  "TELEGRAM_TOKEN=tu_token\n"
                  "TELEGRAM_CHAT_ID=tu_chat_id\n"
                  "ADMIN_CHAT_ID=tu_admin_chat_id\n"
                  "FELIX_AUTH_MORON=Basic base64...\n"
                  "FELIX_AUTH_FIBERALG=Basic base64...\n"
                  "FELIX_AUTH_FIBERPLUS=Basic base64...")
        sys.exit(1)

    if len(sys.argv) > 1 and sys.argv[1] == "una-vez":
        # "una-vez" tambien toma el lock de ejecucion unica: si el bot
        # principal ya esta corriendo (modo por defecto), lanzar
        # "una-vez" a mano a la vez podria hacer que ambos escriban el
        # fichero de estado de la misma red casi simultaneamente y uno
        # pise la transicion detectada por el otro. No corrompe el
        # fichero (la escritura es atomica), pero puede perderse una
        # notificacion, asi que preferimos que "una-vez" se aparte si
        # ya hay otra instancia con el lock tomado.
        if not adquirir_lock_ejecucion():
            log.error(
                "Ya hay otra instancia de este bot corriendo (lock '%s' en "
                "uso). 'una-vez' se cancela para no pisar el fichero de "
                "estado a la vez que la instancia principal.",
                LOCK_FILE,
            )
            sys.exit(1)
        try:
            procesar_una_vez()
        except Exception:
            enviar_alerta_admin(
                "Excepcion no controlada al ejecutar el bot (modo una-vez):\n"
                "```\n{}\n```".format(traceback.format_exc()),
                clave="excepcion_no_controlada_una_vez",
            )
            sys.exit(1)
        sys.exit(0)

    # Los modos "comandos" y por defecto hacen long polling contra
    # getUpdates, que Telegram no permite en paralelo desde dos
    # procesos con el mismo token (responde 409 Conflict). Comprobamos
    # el lock de ejecucion unica ANTES de arrancar nada: si ya hay otra
    # instancia corriendo, avisamos y salimos limpio en vez de dejar
    # dos procesos peleando por el mismo polling y escribiendo a la vez
    # los ficheros de estado.
    if not adquirir_lock_ejecucion():
        log.error(
            "Ya hay otra instancia de este bot corriendo (lock '%s' en uso). "
            "Esta instancia se cierra sin hacer nada. Si crees que no deberia "
            "haber ninguna corriendo, revisa 'ps aux | grep %s' y limpia el "
            "proceso viejo antes de relanzar.",
            LOCK_FILE, os.path.basename(__file__),
        )
        sys.exit(1)

    if len(sys.argv) > 1 and sys.argv[1] == "comandos":
        escuchar_comandos()
        sys.exit(0)

    # Modo por defecto: monitorizacion + comandos a la vez.
    hilo_comandos = threading.Thread(target=escuchar_comandos, daemon=True)
    hilo_comandos.start()
    bucle_monitorizacion()