import json
import os
import re
import subprocess
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
# ---------------------------------------------------------------------------
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()

# REDES (cabeceras SNMP por red, credenciales Felix) vive en config.py,
# en el mismo directorio. Se importa DESPUES de _cargar_dotenv() porque
# felix_auth_header lee las variables FELIX_AUTH_* del .env via
# os.environ.get en el momento de construir la lista.
from config import REDES

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

# --- OIDs Huawei MA56xx/MA58xx (rama privada .1.3.6.1.4.1.2011.6.128) ---
OID_ONT_SN            = ".1.3.6.1.4.1.2011.6.128.1.1.2.43.1.3"   # string
OID_ONT_DESCRIPCION   = ".1.3.6.1.4.1.2011.6.128.1.1.2.43.1.9"   # string
OID_ONT_STATUS        = ".1.3.6.1.4.1.2011.6.128.1.1.2.46.1.15"  # int: 1 online, 2 offline
OID_ONT_LAST_DOWN_CAUSE = ".1.3.6.1.4.1.2011.6.128.1.1.2.46.1.24"  # int
OID_ONT_LAST_DOWN_DATE  = ".1.3.6.1.4.1.2011.6.128.1.1.2.46.1.23"  # hex-string
# Potencia optica que recibe la ONT (Rx), en centesimas: valor/100-100 = dBm
# (conversion sacada de la tabla de OIDs de la app view antigua).
OID_ONT_RX_POWER = ".1.3.6.1.4.1.2011.6.128.1.1.2.51.1.4"
# Estado del puerto PON completo (1 online, 2 offline), indexado por ifIndex.
OID_PON_STATUS = ".1.3.6.1.4.1.2011.6.128.1.1.2.21.1.10"
# Estado administrativo del puerto (IF-MIB, 1 up / 2 down): un puerto
# APAGADO a proposito desde la OLT tiene admin down y no debe avisar.
OID_IF_ADMIN_STATUS = ".1.3.6.1.2.1.2.2.1.7"

# OJO: el motivo 1 NO es LOS. La app view antigua lo excluia
# explicitamente (motivos_excluir = [-1, 1]): es la ONT desactivada
# administrativamente en la OLT, que es lo que ocurre con un impago o
# una baja ordenada desde Felix. Solo 2-6 son perdida de senal real.
MOTIVOS_CAIDA = {
    1: "Desactivada en OLT (impago/baja)",
    2: "Fallo LOS (perdida de senal)",
    3: "Fallo LOS (perdida de senal)",
    4: "Fallo LOS (perdida de senal)",
    5: "Fallo LOS (perdida de senal)",
    6: "Fallo LOS (perdida de senal)",
    9: "Reinicio de la ONT",
    13: "Corte de luz / ONT apagada",
    33: "Restablecida a valores de fabrica",
}

FELIX_ENDPOINT_CLIENTE = "/customer_technical_info_report"
FELIX_ENDPOINT_PAQUETES = "/product_package_group"
FELIX_TIMEOUT = 15  # segundos

# Referencia de caja dentro de external_plant de Felix (mismo regex que
# usaba el bot antiguo y actualizar_onts.php).
RE_CAJA = re.compile(r"\d{1,2}-\D{1,2}\d{0,2}-\D\d{0,2}-\d{2,4}-\d{1,2}")
SNMP_TIMEOUT = 8    # segundos por comando snmpwalk/snmpget
SNMP_RETRIES = 2

SYSTEM_CA_BUNDLE = "/etc/ssl/certs/ca-certificates.crt"

_felix_session = requests.Session()
_telegram_session = requests.Session()

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

BASE_DIR = os.path.dirname(os.path.abspath(__file__))


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


LOG_FILE = os.path.join(os.path.dirname(os.path.abspath(__file__)), "los_bot.log")

OFFSET_FILE = os.path.join(BASE_DIR, "telegram_offset.txt")
LOCK_FILE = os.path.join(BASE_DIR, "bot_felix.lock")
_lock_fd = None


def adquirir_lock_ejecucion():
    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
    _lock_fd.write(str(os.getpid()))
    _lock_fd.flush()
    return True


CHATS_PERMITIDOS = {c for c in (TELEGRAM_CHAT_ID, ADMIN_CHAT_ID) if c}

SLEEP_SECONDS = 300

# Numero de ciclos consecutivos offline antes de notificar (evita
# falsos positivos por glitches puntuales de la OLT).
LOS_CICLOS_CONFIRMACION = 2

# Una caida mas antigua que esto ya no se considera "LOS activo" en
# /los y /estado (ONTs retiradas, clientes que se fueron sin baja...).
# 7 dias: una averia real pendiente de visita sigue saliendo en /los
# aunque tarde varios dias en arreglarse; los registros fantasma de
# meses quedan fuera. (La app view antigua usaba 48h.)
LOS_MAX_HORAS = 168

# --- Deteccion de CAJA/PON caida ---
# Si en un mismo puerto PON (misma tarjeta+puerto de la OLT) cae por
# LOS el CAJA_CAIDA_PORCENTAJE de las ONTs que estaban online, y son
# al menos CAJA_CAIDA_MIN_ONTS, se manda UN aviso de caja caida con
# los clientes afectados (y se suprimen los avisos individuales de
# esas ONTs para no inundar el chat). Solo cuentan caidas recientes
# (ultima hora) para no arrastrar fantasmas viejos del mismo puerto.
CAJA_CAIDA_PORCENTAJE = 0.20
CAJA_CAIDA_MIN_ONTS = 2
CAJA_CAIDA_VENTANA_HORAS = 1
CAJA_CAIDA_VENTANA_MINUTOS = 10  # ventana para el aviso por PORCENTAJE (caja caida);
                                  # la de HORAS sigue usandose solo en PUERTO PON CAIDO
CAJA_CAIDA_MAX_LISTADO = 20  # clientes listados en el aviso como mucho

# Confirmacion por ciclos para PUERTO PON CAIDO (igual filosofia que
# LOS_CICLOS_CONFIRMACION a nivel de ONT): un corte de luz/reinicio de
# cabecera puede tumbar muchos puertos a la vez en un ciclo pero se
# recupera solo en el siguiente; un fallo real de puerto persiste. Si
# el puerto sigue caido tras confirmar, se avisa SIEMPRE, coincida o
# no con otros puertos cayendo a la vez (un puerto con problema real
# no debe silenciarse solo porque otros tuvieron un blip pasajero).
# Incidente 2026-07-28: PRADO DEL REY tuvo 23 puertos "caidos" a la vez
# en el mismo ciclo (corte de luz/reset de cabecera); con confirmacion
# de 2 ciclos, los que se recuperan solos ya no llegan a avisar.
PUERTO_CICLOS_CONFIRMACION = 2

# (cabecera_ip, ifindex_pon) -> ts del aviso enviado; se limpia cuando
# el puerto vuelve a estar por debajo del umbral.
_cajas_caidas_avisadas = {}

# Estado online/offline del puerto PON completo en el ciclo anterior,
# para avisar solo en la transicion online->offline (un puerto que ya
# estaba caido al arrancar el bot no dispara aviso).
_puertos_online_antes = {}
_ciclos_offline = {}
_puertos_ciclos_offline = {}

# Cache del mapa de clientes de Felix. El ciclo de monitorizacion lo
# refresca a la fuerza en cada pasada (~5 min), asi que este TTL es
# solo el respaldo para que un comando de Telegram no se ponga a
# descargar Felix en linea salvo que el ciclo lleve mucho parado.
CACHE_CLIENTE_TTL = 900
_cache_cliente = {}
_cache_cliente_lock = threading.Lock()

# Cache en memoria del ultimo estado SNMP conocido por red (slug ->
# {"red":..., "estado":{cpe_id: info}, "ts": epoch}). Los comandos de
# Telegram (/estado, /los, /buscar) leen de aqui en vez de relanzar el
# walk SNMP completo (que tarda minutos): responden con el dato del
# ultimo ciclo de monitorizacion, no en vivo.
_cache_estado_lock = threading.Lock()
_cache_estado = {}


def _actualizar_cache_estado(red, estado):
    with _cache_estado_lock:
        _cache_estado[red["slug"]] = {"red": red, "estado": estado, "ts": time.time()}


def _antiguedad_legible(ts):
    segundos = time.time() - ts
    if segundos < 60:
        return "hace {:.0f}s".format(segundos)
    if segundos < 3600:
        return "hace {:.0f}min".format(segundos / 60)
    return "hace {:.1f}h".format(segundos / 3600)

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)

_handler_fichero = logging.handlers.RotatingFileHandler(
    LOG_FILE, maxBytes=1024 * 1024, backupCount=5
)
_handler_fichero.setFormatter(_formato)
log.addHandler(_handler_fichero)


# ---------------------------------------------------------------------------
# ESTADO LOCAL (escritura atomica)
# ---------------------------------------------------------------------------

def _escribir_texto_atomico(path, texto):
    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:
        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):
        return {}
    try:
        with open(state_file, "r") as f:
            return json.load(f)
    except (ValueError, IOError) as e:
        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))


# ---------------------------------------------------------------------------
# SNMP (via snmpwalk / snmpget del paquete net-snmp)
# ---------------------------------------------------------------------------

class SnmpError(Exception):
    pass


_RE_SNMP_LINEA = re.compile(
    r'^(\.\d+(?:\.\d+)*)\s*=\s*([A-Za-z_-]+(?:-[A-Za-z]+)?)\s*:\s*(.*)$'
)


# Timeout del PROCESO snmpwalk completo, no de cada paquete individual
# (eso ya lo controla -t/-r). Una OLT con miles de ONTs puede tardar
# bastante mas que unos pocos segundos en volcar toda la tabla via
# GETNEXT/GETBULK, aunque cada paquete individual responda rapido.
SNMP_WALK_TIMEOUT = 180  # segundos


def _snmp_walk_con_tipo(ip, community, oid_base):
    """
    Devuelve dict {indice: (tipo, valor)} donde 'indice' es la parte
    del OID que sobra tras 'oid_base' (p.ej. '0.12345' para el ppg 0,
    ont 12345), 'tipo' es el tipo SNMP tal cual lo da net-snmp (STRING,
    Hex-STRING, INTEGER...) y 'valor' el valor ya limpio (sin comillas).
    Lanza SnmpError si el comando falla o no hay respuesta.
    """
    # snmpbulkwalk (GETBULK) en vez de snmpwalk (GETNEXT uno a uno):
    # trae muchas filas por paquete, muchisimo mas rapido en tablas
    # grandes (miles de ONTs). -Cr25 = hasta 25 filas por respuesta.
    cmd = [
        "snmpbulkwalk", "-v2c", "-c", community,
        "-t", str(SNMP_TIMEOUT), "-r", str(SNMP_RETRIES), "-Cr25",
        "-On", "-Oe", ip, oid_base,
    ]
    try:
        resultado = subprocess.run(
            cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
            universal_newlines=True, timeout=SNMP_WALK_TIMEOUT
        )
    except subprocess.TimeoutExpired:
        raise SnmpError("Timeout ({}s) ejecutando snmpwalk contra {} ({})".format(
            SNMP_WALK_TIMEOUT, ip, oid_base))
    except FileNotFoundError:
        raise SnmpError("Comando 'snmpwalk' no encontrado. Instala el paquete 'snmp' (apt install snmp).")

    salida = resultado.stdout.strip()
    if resultado.returncode != 0 or not salida:
        stderr = resultado.stderr.strip()
        if "No Such" in stderr or "No Such" in salida:
            return {}
        raise SnmpError(
            "snmpwalk fallo contra {} ({}): rc={} stderr={}".format(
                ip, oid_base, resultado.returncode, stderr
            )
        )
    if "Timeout" in salida or "No Response" in salida:
        raise SnmpError("Sin respuesta SNMP de {} ({})".format(ip, oid_base))

    valores = {}
    base_sin_punto_inicial = oid_base.lstrip(".")
    for linea in salida.splitlines():
        m = _RE_SNMP_LINEA.match(linea.strip())
        if not m:
            continue
        oid_completo, tipo, valor = m.groups()
        oid_completo = oid_completo.lstrip(".")
        if not oid_completo.startswith(base_sin_punto_inicial):
            continue
        indice = oid_completo[len(base_sin_punto_inicial):].lstrip(".")
        valor = valor.strip().strip('"')
        valores[indice] = (tipo, valor)
    return valores


def _snmp_walk(ip, community, oid_base):
    """Como _snmp_walk_con_tipo pero descarta el tipo (dict indice -> valor)."""
    return {indice: valor for indice, (_tipo, valor) in _snmp_walk_con_tipo(ip, community, oid_base).items()}


def _decodificar_sn(tipo, valor):
    """
    Los SN de ONT Huawei llegan por SNMP como Hex-STRING de 8 bytes:
    los 4 primeros son el vendor id en ASCII (p.ej. 'HWTC') y los 4
    ultimos son el numero de serie en binario. El SN "de verdad" (el
    que usa Felix como cpe_id) es vendor + hex en mayusculas, p.ej.
    '48 57 54 43 09 8E 14 AE' -> 'HWTC098E14AE'.
    Si no viene como Hex-STRING de 8 bytes, se devuelve el valor tal cual.
    """
    if tipo.lower() != "hex-string":
        return valor
    octetos = valor.split()
    if len(octetos) != 8:
        return valor
    try:
        vendor = "".join(chr(int(o, 16)) for o in octetos[:4])
        serie = "".join(o.upper() for o in octetos[4:])
    except ValueError:
        return valor
    if not vendor.isalnum():
        return valor
    return vendor + serie


def _parsear_fecha_caida(tipo, valor):
    """
    La fecha de ultima caida llega como Hex-STRING de 11 bytes:
    '07 EA 05 05 0C 1C 31 ...' = 0x07EA(2026) 05 05 12:28:49 (+tz).
    Devuelve epoch (time.time()) o None si no se puede parsear o
    viene a ceros (ONT que nunca ha caido).
    """
    if tipo.lower() != "hex-string":
        return None
    octetos = valor.split()
    if len(octetos) < 7:
        return None
    try:
        anio = int(octetos[0], 16) * 256 + int(octetos[1], 16)
        mes, dia = int(octetos[2], 16), int(octetos[3], 16)
        hora, minuto, segundo = int(octetos[4], 16), int(octetos[5], 16), int(octetos[6], 16)
    except ValueError:
        return None
    if anio < 2000 or not (1 <= mes <= 12) or not (1 <= dia <= 31):
        return None
    try:
        return time.mktime((anio, mes, dia, hora, minuto, segundo, 0, 0, -1))
    except (ValueError, OverflowError):
        return None


def _sn_formato_felix(sn):
    """
    Convierte el SN legible ('HWTC47FAA5A8') al formato que usa Felix
    como cpe_id: los 8 bytes completos en hex minuscula, con el vendor
    tambien en hex ('4857544347faa5a8'). Es el mismo formato que la app
    view vieja pasaba a Felix tal cual salia del walk SNMP.
    Si el SN no tiene la pinta esperada (4 chars vendor + 8 hex), se
    devuelve tal cual en minusculas.
    """
    sn = (sn or "").strip()
    if len(sn) == 12 and sn[:4].isalpha():
        vendor_hex = "".join("{:02x}".format(ord(c)) for c in sn[:4])
        return vendor_hex + sn[4:].lower()
    return sn.lower()


def _snmp_get(ip, community, oid):
    """snmpget de un unico valor. Devuelve el valor (string) o None."""
    cmd = [
        "snmpget", "-v2c", "-c", community,
        "-t", "3", "-r", "1", "-On", "-Oe", ip, oid,
    ]
    try:
        resultado = subprocess.run(
            cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
            universal_newlines=True, timeout=15
        )
    except (subprocess.TimeoutExpired, FileNotFoundError):
        return None
    m = _RE_SNMP_LINEA.match(resultado.stdout.strip())
    if not m:
        return None
    return m.group(3).strip().strip('"')


def _texto_potencia_desde_crudo(valor):
    """
    Convierte el valor crudo SNMP de OID_ONT_RX_POWER a texto tipo
    '-21.4 dBm', o None si no da un valor plausible. Segun firmware, la
    OLT devuelve la potencia en centesimas SIN signo con offset 100
    (7950 -> -20.5 dBm, formula de la app view antigua) o directamente
    CON signo (-2050 -> -20.5). Se prueba primero el formato del view y
    si no da un valor plausible, el formato con signo. Rango plausible:
    -45..+10.
    """
    if valor is None:
        return None
    try:
        crudo = int(valor)
    except ValueError:
        return None
    dbm = crudo / 100.0 - 100
    if not (-45 <= dbm <= 10):
        dbm = crudo / 100.0
    if not (-45 <= dbm <= 10):
        return None
    return "{:.1f} dBm".format(dbm)


def _leer_potencia_rx(red, info, reintentos=3, espera=10):
    """
    Lee en vivo la potencia optica Rx de una ONT (solo tiene sentido
    con la ONT online, p.ej. al notificar una recuperacion o un alta).
    Una ONT que acaba de levantar puede tardar un poco en dar lectura
    DDM valida, asi que se reintenta varias veces con espera entre
    intentos. Devuelve texto tipo '-21.4 dBm' o None si no hay lectura.
    """
    ip = info.get("cabecera")
    indice = info.get("indice")
    if not ip or not indice:
        return None
    community = None
    cabecera_cfg = None
    for cabecera in red.get("cabeceras", []):
        if cabecera["ip"] == ip:
            community = cabecera["community"]
            cabecera_cfg = cabecera
            break
    if not community:
        return None
    # Algunas cabeceras (p.ej. BORNOS) tardan mas en refrescar la
    # lectura DDM tras un reenganche; permite alargar los reintentos
    # por cabecera via "potencia_reintentos" en REDES.
    reintentos = (cabecera_cfg or {}).get("potencia_reintentos", reintentos)

    for intento in range(reintentos):
        if intento > 0:
            time.sleep(espera)
        valor = _snmp_get(ip, community, OID_ONT_RX_POWER + "." + indice)
        texto = _texto_potencia_desde_crudo(valor)
        if texto is None:
            continue
        return texto
    log.info("Sin lectura DDM valida para indice %s en %s tras %d intentos.",
             indice, ip, reintentos)
    return None


def _pon_legible(indice, ip):
    """
    Traduce el indice SNMP 'ifIndex.ont' a puerto fisico legible.
    En OLTs Huawei el ifIndex de un puerto GPON es
    4194304000 + slot*8192 + puerto*256, asi que se puede decodificar
    a frame/slot/puerto (verificado contra los datos de Felix:
    4194330880 -> 0/3/9, 4194321152 -> 0/2/3).
    """
    partes = indice.split(".")
    if len(partes) != 2:
        return indice
    try:
        ifidx = int(partes[0])
    except ValueError:
        return indice
    delta = ifidx - 4194304000
    if delta < 0:
        return indice
    slot = delta // 8192
    puerto = (delta % 8192) // 256
    return "0/{}/{} ONT {} ({})".format(slot, puerto, partes[1], ip)


def _indice_par(indice):
    """El indice de OLT Huawei para (ppg, ont) son los DOS ultimos
    numeros del sufijo del OID (el resto, si existe, es prefijo de
    tabla/version segun modelo)."""
    partes = indice.split(".")
    if len(partes) < 2:
        return None
    return partes[-2] + "." + partes[-1]


def _consultar_estado_cabecera(ip, community, nombre=""):
    """SNMP walk de una unica cabecera/OLT. Devuelve tupla:
    (dict cpe_id(SN) -> info, dict ifindex_puerto -> online_bool)."""
    sns_crudos    = _snmp_walk_con_tipo(ip, community, OID_ONT_SN)
    estados       = _snmp_walk(ip, community, OID_ONT_STATUS)
    descripciones = _snmp_walk(ip, community, OID_ONT_DESCRIPCION)
    motivos       = _snmp_walk(ip, community, OID_ONT_LAST_DOWN_CAUSE)
    fechas_caida  = _snmp_walk_con_tipo(ip, community, OID_ONT_LAST_DOWN_DATE)
    try:
        potencias = _snmp_walk(ip, community, OID_ONT_RX_POWER)
    except SnmpError:
        # Complementario: si falla no se pierde el ciclo, solo se deja
        # de refrescar la ultima potencia conocida de esta cabecera.
        potencias = {}

    puertos = {}
    try:
        for indice_puerto, valor in _snmp_walk(ip, community, OID_PON_STATUS).items():
            puertos[indice_puerto.split(".")[-1]] = str(valor).strip() == "1"
    except SnmpError:
        # El estado de puertos es complementario: si este walk falla no
        # se pierde el ciclo de la cabecera, solo la deteccion de
        # puerto caido de esta pasada.
        puertos = {}

    estado = {}
    for indice, (tipo, valor) in sns_crudos.items():
        sn = _decodificar_sn(tipo, valor)
        if not sn:
            continue
        par = _indice_par(indice)
        estado_valor = estados.get(indice)
        if estado_valor is None:
            continue
        online = str(estado_valor).strip() == "1"

        motivo_texto = ""
        caida_ts = None
        if not online:
            motivo_raw = motivos.get(indice)
            if motivo_raw is not None:
                try:
                    motivo_codigo = int(motivo_raw)
                    motivo_texto = MOTIVOS_CAIDA.get(motivo_codigo, "Desconocido ({})".format(motivo_codigo))
                except ValueError:
                    pass
            fecha_raw = fechas_caida.get(indice)
            if fecha_raw is not None:
                caida_ts = _parsear_fecha_caida(fecha_raw[0], fecha_raw[1])

        # Potencia solo se toma del walk en vivo si esta online (offline
        # no da lectura DDM valida). Cuando esta offline se deja en
        # None aqui: el llamador (ciclo principal) la rellena con la
        # ultima conocida persistida en disco de cuando SI estaba online.
        potencia_texto = _texto_potencia_desde_crudo(potencias.get(indice)) if online else None

        estado[sn] = {
            "online": online,
            "indice": par or indice,
            "pon_legible": _pon_legible(par or indice, nombre or ip),
            "descripcion": descripciones.get(indice, ""),
            "motivo": motivo_texto,
            "caida_ts": caida_ts,
            "cabecera": ip,
            "cabecera_nombre": nombre,
            "potencia": potencia_texto,
        }
    return estado, puertos


def _es_los_activo(info):
    """
    LOS "de verdad": offline por perdida de senal Y caida reciente
    (menos de LOS_MAX_HORAS). Caidas de hace semanas/meses (clientes
    idos, ONTs retiradas) no cuentan como LOS activo, igual que en la
    app view antigua. Si la OLT no da fecha, se asume reciente.
    """
    if info["online"]:
        return False
    if not (info.get("motivo") or "").startswith("Fallo LOS"):
        return False
    caida_ts = info.get("caida_ts")
    if caida_ts is None:
        return True
    return (time.time() - caida_ts) <= LOS_MAX_HORAS * 3600


def _mejor_registro(a, b):
    """
    El mismo SN puede aparecer registrado en varias OLTs/puertos (mudanzas
    de puerto, registros fantasma antiguos). Se queda con el mejor:
    online gana siempre; si ambos offline, gana el de caida mas reciente.
    """
    if a["online"] != b["online"]:
        return a if a["online"] else b
    ts_a = a.get("caida_ts") or 0
    ts_b = b.get("caida_ts") or 0
    return a if ts_a >= ts_b else b


def consultar_estado_olt(red):
    """
    Vuelca por SNMP el estado de todas las ONTs de TODAS las cabeceras
    (OLTs) de una red, y junta los resultados en un unico dict
    cpe_id(SN) -> {online, indice, descripcion, motivo}.

    Si una cabecera concreta falla pero hay otras, se seguimos con las
    demas y se loguea el fallo parcial (no se aborta toda la red por
    una sola OLT caida). Solo se lanza SnmpError si TODAS las
    cabeceras configuradas fallan, o si la red no tiene ninguna
    cabecera dada de alta todavia.
    """
    cabeceras = red.get("cabeceras") or []
    if not cabeceras:
        log.warning("[%s] Red sin cabeceras SNMP configuradas todavia, se omite.", red["nombre"])
        return {}, []

    # Cada cabecera se consulta en su propio hilo: los walks SNMP son
    # lentos (varios minutos en serie) y asi el tiempo total de la red
    # es el de la cabecera mas lenta, no la suma de todas.
    resultados = {}  # ip -> dict parcial o SnmpError

    def _consultar_con_reintento(cabecera):
        # Reintento propio por cabecera (aparte del -r de snmpbulkwalk):
        # algunas OLTs fallan con "No Response" de forma intermitente.
        error_final = None
        for intento in range(2):
            try:
                resultados[cabecera["ip"]] = _consultar_estado_cabecera(
                    cabecera["ip"], cabecera["community"], cabecera.get("nombre", ""))
                return
            except SnmpError as e:
                error_final = e
            except Exception as e:  # no dejar morir el hilo en silencio
                error_final = SnmpError("Error inesperado: {}".format(e))
        resultados[cabecera["ip"]] = error_final

    hilos = []
    for cabecera in cabeceras:
        hilo = threading.Thread(target=_consultar_con_reintento, args=(cabecera,))
        hilo.start()
        hilos.append(hilo)
    for hilo in hilos:
        hilo.join()

    estado = {}
    puertos_todos = []  # lista de (ip, nombre, ifindex, online_bool)
    errores = []
    for cabecera in cabeceras:
        resultado = resultados.get(cabecera["ip"])
        if isinstance(resultado, SnmpError) or resultado is None:
            errores.append(str(resultado))
            log.warning("[%s] Fallo SNMP en cabecera %s (tras reintento): %s",
                        red["nombre"], cabecera["ip"], resultado)
            continue
        parcial, puertos = resultado
        for sn, registro in parcial.items():
            if sn in estado:
                estado[sn] = _mejor_registro(estado[sn], registro)
            else:
                estado[sn] = registro
        for ifindex, online in puertos.items():
            puertos_todos.append(
                (cabecera["ip"], cabecera.get("nombre", ""), ifindex, online)
            )

    if errores and len(errores) == len(cabeceras):
        raise SnmpError(
            "Todas las cabeceras de la red '{}' fallaron: {}".format(
                red["nombre"], "; ".join(errores)
            )
        )
    return estado, puertos_todos


# ---------------------------------------------------------------------------
# FELIX: solo para SN -> cliente / activo (bajo demanda, con cache)
# ---------------------------------------------------------------------------

def _ultimos8_hex(sn):
    """
    Los 8 caracteres finales del SN son la parte de numero de serie en
    hex (los 8 primeros son el vendor id, sea como sea que cada lado
    -bot via SNMP, Felix via su propio origen- lo represente: ASCII
    'HWTC', hex '48575443', mayus/minuscula...). Comparando solo esta
    cola evitamos tener que adivinar el formato exacto que usa Felix.
    """
    sn = (sn or "").strip()
    return sn[-8:].upper() if len(sn) >= 8 else sn.upper()


def _consultar_felix_clientes_red(red):
    """Descarga TODOS los clientes/ppgs de una red desde Felix (sin
    filtrar por cpe_id) -mismo endpoint que ya se usaba, sin parametro-."""
    url = red["felix_base_url"] + FELIX_ENDPOINT_CLIENTE
    headers = {
        "authorization": red["felix_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 _consultar_felix_paquetes_red(red):
    """Descarga los product_package_group de la red (para sacar la
    referencia de caja de external_plant, como hacia el bot antiguo)."""
    url = red["felix_base_url"] + FELIX_ENDPOINT_PAQUETES
    headers = {
        "authorization": red["felix_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 _obtener_mapa_clientes_felix(red, forzar=False):
    """
    Dict {ultimos8_hex_del_sn: {'cliente':str, 'activa':bool, 'caja':str}}
    para toda la red, refrescado como mucho cada CACHE_CLIENTE_TTL
    segundos (1 descarga completa por red, no 1 consulta por SN).
    Con forzar=True refresca siempre (lo usa el ciclo de monitorizacion
    para que los comandos de Telegram nunca tengan que descargar Felix).
    """
    slug = red["slug"]
    ahora = time.time()
    with _cache_cliente_lock:
        cache_hit = _cache_cliente.get(slug)
        if not forzar and cache_hit and (ahora - cache_hit["ts"]) < CACHE_CLIENTE_TTL:
            return cache_hit["mapa"]

    mapa = {}
    try:
        data = _consultar_felix_clientes_red(red)
        for cliente in data or []:
            nombre = cliente.get("full_name") or "N/D"
            for ppg in cliente.get("ppgs", []):
                cpe_id_felix = ppg.get("cpe_id")
                if not cpe_id_felix:
                    continue
                # Impago: Felix lo marca metiendo "Impago" en la lista
                # de services del ppg (no en el campo "active").
                services = [str(s).strip().lower() for s in (ppg.get("services") or [])]
                impago = "impago" in services

                # Ubicacion en la OLT segun Felix: olt_id + frame/slot/port.
                olt_id = (ppg.get("olt_id") or "").strip()
                frame, slot, port = ppg.get("frame"), ppg.get("slot"), ppg.get("port")
                if olt_id and frame is not None and slot is not None and port is not None:
                    olt_pon = "{} {}/{}/{}".format(olt_id, frame, slot, port)
                elif olt_id:
                    olt_pon = olt_id
                else:
                    olt_pon = "N/D"

                mapa[_ultimos8_hex(cpe_id_felix)] = {
                    "cliente": nombre,
                    "activa": bool(ppg.get("active", True)) and not impago,
                    "impago": impago,
                    "caja": "N/D",
                    "olt_pon": olt_pon,
                }
    except requests.RequestException as e:
        log.warning("No se pudo refrescar el mapa de clientes de Felix para '%s': %s",
                    red["nombre"], e)
        if cache_hit:
            # Felix caido pero habia cache previa: mejor usar la
            # vieja (aunque desactualizada) que quedarse sin nada.
            return cache_hit["mapa"]

    # Referencia de caja (external_plant) por SN, si el endpoint responde.
    try:
        for grupo in _consultar_felix_paquetes_red(red) or []:
            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 "N/D")
            for paquete in grupo.get("product_packages", []):
                cpe_id_felix = paquete.get("cpe_id")
                if not cpe_id_felix:
                    continue
                clave = _ultimos8_hex(cpe_id_felix)
                if clave in mapa:
                    mapa[clave]["caja"] = referencia
    except requests.RequestException as e:
        log.warning("No se pudo refrescar las referencias de caja de Felix para '%s': %s",
                    red["nombre"], e)

    with _cache_cliente_lock:
        _cache_cliente[slug] = {"mapa": mapa, "ts": ahora}
    return mapa


def obtener_info_cliente(red, cpe_id):
    """
    Devuelve {'cliente','activa','caja'} para un SN, casando por los
    ultimos 8 caracteres hex contra el mapa de Felix de esa red. Si no
    hay match o Felix falla, devuelve valores por defecto sin bloquear
    el aviso de LOS.
    """
    info = obtener_info_cliente_o_none(red, cpe_id)
    if info is not None:
        return info
    return {"cliente": "N/D", "activa": True, "caja": "N/D", "olt_pon": "N/D"}


def obtener_info_cliente_o_none(red, cpe_id):
    """Como obtener_info_cliente pero devuelve None si el SN no esta
    asignado a ningun cliente en Felix (sin match)."""
    mapa = _obtener_mapa_clientes_felix(red)
    return mapa.get(_ultimos8_hex(cpe_id))


# ---------------------------------------------------------------------------
# MENSAJE Y ENVIO A TELEGRAM
# ---------------------------------------------------------------------------

class _RateLimited(Exception):
    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):
    if texto is None:
        return texto
    return _RE_MARKDOWN_ESPECIALES.sub(r"\\\1", str(texto))


def _olt_con_nombre(olt_pon, nombre_cabecera):
    """
    Inserta el nombre de la cabecera en la ubicacion OLT de Felix:
    'fiberplus:gpon2 0/3/2' + 'MONTELLANO' ->
    'fiberplus:gpon2 (MONTELLANO) 0/3/2'.
    """
    olt_pon = olt_pon or "N/D"
    if not nombre_cabecera:
        return olt_pon
    partes = olt_pon.split(" ", 1)
    if len(partes) == 2:
        return "{} ({}) {}".format(partes[0], nombre_cabecera, partes[1])
    return "{} ({})".format(olt_pon, nombre_cabecera)


def _icono_titulo_offline(motivo):
    """Titulo del aviso segun el motivo real de caida (no todo es LOS)."""
    motivo = motivo or ""
    if motivo.startswith("Fallo LOS"):
        return "\U0001F534", "LOS detectado"
    if motivo.startswith("Corte de luz"):
        return "\U000026A1", "Corte de luz detectado"
    if motivo.startswith("Reinicio"):
        return "\U0001F501", "Reinicio de ONT detectado"
    if motivo.startswith("Restablecida"):
        return "\U00002699", "Valores de fabrica restablecidos"
    if motivo.startswith("Desactivada"):
        return "\U0001F4A4", "ONT desactivada (impago/baja)"
    if motivo:
        return "\U00002753", "Caida detectada ({})".format(motivo)
    return "\U0001F534", "Offline detectado"


def construir_mensaje(red, cpe_id, info_snmp, info_cliente, tipo):
    if tipo == "offline":
        icono, titulo = _icono_titulo_offline(info_snmp.get("motivo"))
    elif tipo == "alta":
        icono, titulo = "\U00002705", "Alta realizada"
    else:
        icono, titulo = "\U0001F7E2", "Recuperado"

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

    motivo_linea = (
        "\n⚠️ *Motivo:* {}".format(_escapar_markdown(info_snmp.get("motivo")))
        if tipo == "offline" and info_snmp.get("motivo") else ""
    )
    if tipo == "offline" and info_snmp.get("potencia"):
        potencia_linea = "\n\U0001F4F6 *Ultima potencia registrada:* {}".format(
            _escapar_markdown(info_snmp.get("potencia")))
    elif tipo in ("online", "alta") and info_snmp.get("potencia"):
        potencia_linea = "\n\U0001F4F6 *Potencia Rx:* {}".format(
            _escapar_markdown(info_snmp.get("potencia")))
    else:
        potencia_linea = ""
    aviso_inactiva = (
        "\n\U0001F6AB _Marcada como inactiva en Felix_"
        if not info_cliente.get("activa", True) else ""
    )

    return (
        "{icono} *{titulo}*\n"
        "━━━━━━━━━━━━━━\n"
        "\U0001F4E1 *Red:* {red}\n"
        "\U0001F464 *Cliente:* {cliente}\n"
        "\U0001F50C *SN:* `{sn}`\n"
        "\U0001F4E6 *Caja:* {caja}\n"
        "\U0001F5A5 *OLT:* {olt_pon}\n"
        "\U0001F4CD *PON:* {indice}\n"
        "\U0001F550 *Hora:* {fecha}"
        "{motivo}"
        "{potencia}"
        "{aviso}"
    ).format(
        icono=icono,
        titulo=_escapar_markdown(titulo),
        red=_escapar_markdown(red["nombre"]),
        cliente=_escapar_markdown(info_cliente["cliente"]),
        sn=_sn_formato_felix(cpe_id),
        caja=_escapar_markdown(info_cliente.get("caja", "N/D")),
        olt_pon=_escapar_markdown(_olt_con_nombre(
            info_cliente.get("olt_pon"), info_snmp.get("cabecera_nombre", ""))),
        indice=_escapar_markdown(info_snmp.get("pon_legible") or info_snmp.get("indice", "N/D")),
        fecha=_escapar_markdown(fecha_str),
        motivo=motivo_linea,
        potencia=potencia_linea,
        aviso=aviso_inactiva,
    )


TELEGRAM_MAX_CHARS = 4000
TELEGRAM_MIN_INTERVALO = 3.5
_ultimo_envio_lock = threading.Lock()
_ultimo_envio_ts = [0.0]
TELEGRAM_MAX_REINTENTOS_429 = 3


def _esperar_turno_envio():
    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)

    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:
            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", "")
        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,
        )


ALERTA_ADMIN_COOLDOWN = 1800
_alertas_admin_lock = threading.Lock()
_ultimas_alertas_admin = {}


def enviar_alerta_admin(mensaje, clave=None):
    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 = "⚠️ *Fallo en telegram\\_los\\_bot*\n{}".format(mensaje)
    enviar_telegram(texto, chat_id=ADMIN_CHAT_ID)


# ---------------------------------------------------------------------------
# COMANDOS DE TELEGRAM (leen del cache, actualizado cada ciclo de
# monitorizacion, en vez de relanzar el walk SNMP completo -que tarda
# minutos- en cada comando)
# ---------------------------------------------------------------------------

MAX_LINEAS_LISTADO = 30


def obtener_estado_todas_redes():
    """
    Devuelve (resultado, redes_sin_cache) leyendo del cache en memoria
    que va rellenando procesar_red() en cada ciclo. 'redes_sin_cache'
    son redes que todavia no han completado ni un ciclo (bot recien
    arrancado) o que no tienen cabeceras configuradas.
    """
    resultado = []
    redes_sin_cache = []
    with _cache_estado_lock:
        snapshot = dict(_cache_estado)
    for red in REDES:
        entrada = snapshot.get(red["slug"])
        if entrada is None:
            redes_sin_cache.append(red["nombre"])
            continue
        for cpe_id, info in entrada["estado"].items():
            resultado.append((entrada["red"], cpe_id, info))
    return resultado, redes_sin_cache


def _antiguedad_cache_mas_vieja():
    """Antiguedad (texto legible) del dato en cache mas desactualizado
    entre todas las redes con cache. None si no hay ninguna."""
    with _cache_estado_lock:
        timestamps = [v["ts"] for v in _cache_estado.values()]
    if not timestamps:
        return None
    return _antiguedad_legible(min(timestamps))


def comando_estado():
    datos, fallidas = obtener_estado_todas_redes()
    if not datos:
        return "Todavia no hay datos en cache\\. Prueba en unos minutos\\."

    # Tres categorias: online / LOS / offline (resto de caidas). El
    # contador de LOS usa EXACTAMENTE el mismo criterio que /los
    # (senal reciente + cliente activo en Felix) para que ambos
    # comandos siempre cuadren entre si.
    stats_por_red = {}
    for red, cpe_id, info in datos:
        nombre_red = red["nombre"]
        if nombre_red not in stats_por_red:
            stats_por_red[nombre_red] = {"online": 0, "los": 0, "offline": 0}
        if info["online"]:
            stats_por_red[nombre_red]["online"] += 1
        elif _es_los_activo(info):
            cliente = obtener_info_cliente_o_none(red, cpe_id)
            if cliente is not None and cliente.get("activa", True):
                stats_por_red[nombre_red]["los"] += 1
            else:
                stats_por_red[nombre_red]["offline"] += 1
        else:
            stats_por_red[nombre_red]["offline"] += 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_offline = sum(s["offline"] for s in stats_por_red.values())

    lineas = ["\U0001F4CA *Estado actual*"]
    for nombre_red, s in stats_por_red.items():
        lineas.append(
            "\n*{}*\n"
            "  \U0001F7E2 Online: {}\n"
            "  \U0001F534 LOS: {}\n"
            "  \U000026AB Offline: {}".format(
                _escapar_markdown(nombre_red), s["online"], s["los"], s["offline"]
            )
        )

    lineas.append(
        "\n*Total: {} ONTs*\n"
        "  \U0001F7E2 Online: {}\n"
        "  \U0001F534 LOS: {}\n"
        "  \U000026AB Offline: {}".format(
            len(datos), total_online, total_los, total_offline
        )
    )

    antiguedad = _antiguedad_cache_mas_vieja()
    if antiguedad:
        lineas.append("\n\U0001F551 Dato de {}".format(_escapar_markdown(antiguedad)))

    if fallidas:
        lineas.append(
            "\n⚠️ Sin cache todavia de: {}".format(
                ", ".join(_escapar_markdown(r) for r in fallidas)
            )
        )

    return "\n".join(lineas)


def _trocear_mensajes(bloques, cabecera):
    """
    Junta 'bloques' (lista de strings, uno por ONT) en el minimo numero
    de mensajes de Telegram sin pasarse de TELEGRAM_MAX_CHARS. El
    primer mensaje lleva la cabecera. Devuelve lista de mensajes.
    """
    mensajes = []
    actual = cabecera
    for bloque in bloques:
        candidato = actual + "\n" + bloque if actual else bloque
        if len(candidato) > TELEGRAM_MAX_CHARS - 100:
            mensajes.append(actual)
            actual = bloque.lstrip("\n")
        else:
            actual = candidato
    if actual:
        mensajes.append(actual)
    return mensajes


def comando_offline():
    """
    Solo LOS real (perdida de senal). Solo se muestran las ONTs que
    tienen cliente asignado en Felix (sin match -> se omiten, suelen
    ser ONTs dadas de baja o de pruebas). Si el listado no cabe en un
    mensaje de Telegram, se parte en varios.
    """
    datos, fallidas = obtener_estado_todas_redes()
    if not datos:
        return "Todavia no hay datos en cache \\(el bot no ha completado su primera pasada, o ninguna red tiene cabeceras configuradas\\)\\. Prueba en unos minutos\\."

    caidas = []
    sin_match = 0
    inactivas = 0
    antiguas = 0
    for red, cpe_id, info in datos:
        if info["online"] or not (info.get("motivo") or "").startswith("Fallo LOS"):
            continue
        if not _es_los_activo(info):
            # Caida por senal pero de hace mas de LOS_MAX_HORAS: ONT
            # retirada/cliente ido, no averia activa.
            antiguas += 1
            continue
        cliente = obtener_info_cliente_o_none(red, cpe_id)
        if cliente is None:
            sin_match += 1
            continue
        if not cliente.get("activa", True):
            # Inactiva en Felix (impago/baja): cortada a proposito, no
            # es una averia -> fuera del listado de LOS.
            inactivas += 1
            continue
        caidas.append((red, cpe_id, info, cliente))

    omitidas = []
    if sin_match:
        omitidas.append("{} sin cliente en Felix".format(sin_match))
    if inactivas:
        omitidas.append("{} inactivas/impago".format(inactivas))
    if antiguas:
        omitidas.append("{} caidas antiguas de mas de {}h".format(antiguas, LOS_MAX_HORAS))
    nota_omitidas = ", ".join(omitidas)

    aviso_fallidas = (
        "\n⚠️ Sin cache todavia de: {}".format(", ".join(_escapar_markdown(r) for r in fallidas))
        if fallidas else ""
    )
    if not caidas:
        extra = " \\({} omitidas\\)".format(_escapar_markdown(nota_omitidas)) if nota_omitidas else ""
        return "✅ No hay ninguna ONT en LOS ahora mismo\\.{}".format(extra) + aviso_fallidas

    cabecera = "\U0001F534 *ONTs en LOS \\({}\\)*".format(len(caidas))
    if nota_omitidas:
        cabecera += "\n_\\({} omitidas\\)_".format(_escapar_markdown(nota_omitidas))

    bloques = []
    for red, cpe_id, info, cliente in caidas:
        bloques.append(
            "\n\U0001F464 *{cliente}*\n"
            "\U0001F50C SN: `{sn_felix}`\n"
            "\U0001F4E6 Caja: {caja}\n"
            "\U0001F5A5 OLT: {olt_pon}\n"
            "\U0001F4CD PON {indice} \\| {red}".format(
                cliente=_escapar_markdown(cliente["cliente"]),
                sn_felix=_sn_formato_felix(cpe_id),
                caja=_escapar_markdown(cliente.get("caja", "N/D")),
                olt_pon=_escapar_markdown(_olt_con_nombre(
                    cliente.get("olt_pon"), info.get("cabecera_nombre", ""))),
                indice=_escapar_markdown(info.get("pon_legible") or info.get("indice", "N/D")),
                red=_escapar_markdown(red["nombre"]),
            )
        )
    if aviso_fallidas:
        bloques.append(aviso_fallidas)
    return _trocear_mensajes(bloques, cabecera)


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

    datos, fallidas = obtener_estado_todas_redes()
    if not datos:
        return "Todavia no hay datos en cache\\. Prueba en unos minutos\\."

    # Acepta el SN en cualquiera de los dos formatos: 'HWTC47FAA5A8'
    # (como lo da la OLT) o '4857544347faa5a8' (como lo usa Felix).
    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 _sn_formato_felix(cpe_id)
    ]
    if not coincidencias:
        aviso_fallidas = (
            " \\(sin cache todavia 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"
        cliente = obtener_info_cliente(red, cpe_id)
        lineas.append(
            "{icono} *{cliente}*\n"
            "\U0001F50C SN: `{sn_felix}`\n"
            "\U0001F4E6 Caja: {caja}\n"
            "\U0001F5A5 OLT: {olt_pon}\n"
            "\U0001F4CD PON {indice} \\| {red}".format(
                icono=icono,
                cliente=_escapar_markdown(cliente["cliente"]),
                sn_felix=_sn_formato_felix(cpe_id),
                caja=_escapar_markdown(cliente.get("caja", "N/D")),
                olt_pon=_escapar_markdown(_olt_con_nombre(
                    cliente.get("olt_pon"), info.get("cabecera_nombre", ""))),
                indice=_escapar_markdown(info.get("pon_legible") or info.get("indice", "N/D")),
                red=_escapar_markdown(red["nombre"]),
            )
        )
    if len(coincidencias) > MAX_LINEAS_LISTADO:
        lineas.append("\\.\\.\\. y {} mas \\(afina la busqueda\\)".format(len(coincidencias) - MAX_LINEAS_LISTADO))
    if fallidas:
        lineas.append(
            "\n⚠️ Sin cache todavia 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]
    argumento = partes[1].strip() if len(partes) > 1 else ""

    if comando == "/estado":
        respuesta = comando_estado()
    elif comando == "/los":
        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"
            "/los \\- lista de ONTs en LOS ahora mismo\n"
            "/buscar <SN\\> \\- estado de una ONT concreta"
        )
    else:
        respuesta = "Comando no reconocido\\. Prueba /ayuda\\."

    # Los comandos pueden devolver un mensaje (str) o varios (list)
    # si el listado no cabe en un solo mensaje de Telegram.
    if isinstance(respuesta, list):
        for mensaje in respuesta:
            enviar_telegram(mensaje, chat_id=chat_id)
    else:
        enviar_telegram(respuesta, chat_id=chat_id)


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

def _detectar_puertos_caidos(red, puertos, estado_actual):
    """
    Detecta la caida de un puerto PON COMPLETO (OID de estado del
    puerto, no calculo por ONTs) y avisa con los clientes afectados.
    Devuelve el conjunto de cpe_id de las ONTs de puertos caidos (para
    suprimir sus avisos individuales).

    Exige PUERTO_CICLOS_CONFIRMACION ciclos seguidos offline antes de
    avisar (misma logica que LOS_CICLOS_CONFIRMACION a nivel de ONT):
    un corte de luz/reinicio de cabecera puede tumbar muchos puertos a
    la vez mostrando offline en UN ciclo, pero se recupera solo antes
    de confirmar; un fallo real de puerto persiste y sigue avisando
    aunque coincida con otros puertos cayendo a la vez.
    """
    afectadas = set()
    for ip, nombre, ifindex, online in puertos:
        clave = (ip, ifindex)
        antes = _puertos_online_antes.get(clave)
        _puertos_online_antes[clave] = online

        # ONTs que cuelgan de este puerto en esta cabecera
        onts_puerto = [
            (cpe_id, info) for cpe_id, info in estado_actual.items()
            if info.get("cabecera") == ip
            and (info.get("indice") or "").split(".")[0] == ifindex
        ]

        if not online:
            for cpe_id, _info in onts_puerto:
                afectadas.add(cpe_id)

        if not online:
            if antes:
                _puertos_ciclos_offline[clave] = 1
            elif clave in _puertos_ciclos_offline:
                _puertos_ciclos_offline[clave] += 1
            ciclos = _puertos_ciclos_offline.get(clave, 0)

            if not onts_puerto or ciclos != PUERTO_CICLOS_CONFIRMACION:
                if 0 < ciclos < PUERTO_CICLOS_CONFIRMACION:
                    log.info("[%s] Puerto PON %s de %s caido, ciclo %d/%d, "
                             "esperando confirmacion.", red["nombre"], ifindex,
                             nombre or ip, ciclos, PUERTO_CICLOS_CONFIRMACION)
                continue
            # Distinguir APAGADO (admin down: alguien cerro el puerto
            # a proposito en la OLT) de caida real. Solo avisa la real.
            if _puerto_apagado_admin(red, ip, ifindex):
                log.info("[%s] Puerto PON %s de %s pasa a offline por APAGADO "
                         "administrativo, no se avisa.", red["nombre"], ifindex, nombre or ip)
                continue
            # Puerto caido de verdad = NINGUNA ONT del puerto sigue
            # online. Si alguna aparece online (desfase entre walks),
            # el puerto no esta caido y no se avisa.
            if any(info["online"] for _c, info in onts_puerto):
                log.info("[%s] Puerto PON %s de %s marcado offline pero con ONTs "
                         "online, no se avisa (desfase).", red["nombre"], ifindex, nombre or ip)
                continue
            # Confirmar que los clientes del puerto estan de verdad en
            # LOS (perdida de senal) Y que la caida es RECIENTE (de
            # esta misma incidencia, no un cliente que ya llevaba dias
            # caido de antes). Sin este filtro de fecha, un puerto que
            # se cae de verdad arrastraba en el listado a clientes
            # viejos que no tienen nada que ver con el corte de ahora.
            ahora = time.time()
            ventana = CAJA_CAIDA_VENTANA_HORAS * 3600
            onts_en_los = [
                (cpe_id, info) for cpe_id, info in onts_puerto
                if not info["online"]
                and (info.get("motivo") or "").startswith("Fallo LOS")
                and info.get("caida_ts") is not None
                and (ahora - info["caida_ts"]) <= ventana
            ]
            if not onts_en_los:
                log.info("[%s] Puerto PON %s de %s offline pero sin ONTs en LOS "
                         "recientes, no se avisa.", red["nombre"], ifindex, nombre or ip)
                continue
            # Marca tambien la deteccion por porcentaje para que no
            # duplique aviso sobre el mismo puerto.
            if clave not in _cajas_caidas_avisadas:
                _cajas_caidas_avisadas[clave] = time.time()
                _avisar_puerto_caido(red, ip, nombre, ifindex, onts_en_los)
        else:
            _puertos_ciclos_offline.pop(clave, None)
            if antes is False:
                log.info("[%s] Puerto PON %s de %s recuperado.",
                         red["nombre"], _pon_legible(ifindex + ".0", nombre or ip).rsplit(" ONT", 1)[0], ip)
    return afectadas


def _puerto_apagado_admin(red, ip, ifindex):
    """True si el puerto esta administrativamente apagado (ifAdminStatus=2).
    Si no se puede leer, se asume que NO esta apagado (mejor avisar de
    mas que callar una caida real)."""
    community = None
    for cabecera in red.get("cabeceras", []):
        if cabecera["ip"] == ip:
            community = cabecera["community"]
            break
    if not community:
        return False
    valor = _snmp_get(ip, community, OID_IF_ADMIN_STATUS + "." + ifindex)
    return valor is not None and valor.strip() == "2"


def _avisar_puerto_caido(red, ip, nombre, ifindex, onts_puerto):
    clientes_listado = []
    olt_pon = "N/D"
    for cpe_id, _info in onts_puerto:
        cliente = obtener_info_cliente_o_none(red, cpe_id)
        if cliente is None or not cliente.get("activa", True):
            continue
        if olt_pon == "N/D" and cliente.get("olt_pon", "N/D") != "N/D":
            olt_pon = cliente["olt_pon"]
        clientes_listado.append(
            "\\- {} \\(`{}`\\)".format(
                _escapar_markdown(cliente["cliente"]), _sn_formato_felix(cpe_id)
            )
        )

    if not clientes_listado:
        # Puerto sin ningun cliente activo en Felix (ONTs de pruebas,
        # retiradas, bajas...): caida irrelevante, no se avisa.
        log.info("[%s] Puerto PON caido en %s (ifindex %s) sin clientes en Felix, no se avisa.",
                 red["nombre"], nombre or ip, ifindex)
        return

    extra = ""
    if len(clientes_listado) > CAJA_CAIDA_MAX_LISTADO:
        extra = "\n\\.\\.\\. y {} mas".format(len(clientes_listado) - CAJA_CAIDA_MAX_LISTADO)
        clientes_listado = clientes_listado[:CAJA_CAIDA_MAX_LISTADO]

    puerto_str = _pon_legible(ifindex + ".0", nombre or ip).rsplit(" ONT", 1)[0]
    texto = (
        "\U0001F6A8 *PUERTO PON CAIDO*\n"
        "━━━━━━━━━━━━━━\n"
        "\U0001F4E1 *Red:* {red}\n"
        "\U0001F3E0 *Cabecera:* {cabecera}\n"
        "\U0001F5A5 *OLT:* {olt_pon}\n"
        "\U0001F4CD *Puerto:* {puerto}\n"
        "\U0001F465 *ONTs afectadas:* {n_onts}\n"
        "\U0001F550 *Hora:* {fecha}\n"
        "\n*Clientes afectados:*\n{listado}{extra}"
    ).format(
        red=_escapar_markdown(red["nombre"]),
        cabecera=_escapar_markdown(nombre or ip),
        olt_pon=_escapar_markdown(_olt_con_nombre(olt_pon, nombre)),
        puerto=_escapar_markdown(puerto_str),
        n_onts=len(onts_puerto),
        fecha=_escapar_markdown(time.strftime("%Y-%m-%d %H:%M:%S")),
        listado="\n".join(clientes_listado) if clientes_listado else "\\(sin clientes con match en Felix\\)",
        extra=extra,
    )
    enviar_telegram(texto)
    log.info("[%s] PUERTO PON CAIDO avisado: %s en %s, %d ONTs.",
             red["nombre"], puerto_str, ip, len(onts_puerto))


def _detectar_cajas_caidas(red, estado_actual):
    """
    Agrupa las ONTs por puerto PON fisico (cabecera + primer numero del
    indice SNMP = misma tarjeta y mismo puerto de la OLT) y detecta
    puertos donde ha caido por LOS al menos CAJA_CAIDA_PORCENTAJE de
    las ONTs activas en Felix, en los ultimos CAJA_CAIDA_VENTANA_MINUTOS
    minutos. Manda un unico aviso de CAJA CAIDA con
    los clientes afectados y devuelve el conjunto de cpe_id cubiertos
    (para suprimir sus avisos individuales).
    """
    ahora = time.time()
    ventana = CAJA_CAIDA_VENTANA_MINUTOS * 60

    grupos = {}
    for cpe_id, info in estado_actual.items():
        pon_fisico = (info.get("indice") or "").split(".")[0]
        if not pon_fisico:
            continue
        clave_pon = (info.get("cabecera", ""), pon_fisico)
        grupos.setdefault(clave_pon, []).append((cpe_id, info))

    afectadas = set()
    for clave_pon, onts in grupos.items():
        activas = [
            o for o in onts
            if (obtener_info_cliente_o_none(red, o[0]) or {}).get("activa", False)
        ]
        total_relevante = len(activas)
        if total_relevante == 0:
            continue
        caidas_los = [
            o for o in activas
            if not o[1]["online"]
            and (o[1].get("motivo") or "").startswith("Fallo LOS")
            and o[1].get("caida_ts") is not None
            and (ahora - o[1]["caida_ts"]) <= ventana
        ]
        porcentaje = len(caidas_los) / float(total_relevante)

        if len(caidas_los) >= CAJA_CAIDA_MIN_ONTS and porcentaje >= CAJA_CAIDA_PORCENTAJE:
            for cpe_id, _info in caidas_los:
                afectadas.add(cpe_id)
            if clave_pon in _cajas_caidas_avisadas:
                continue  # ya avisada, no repetir cada ciclo
            _cajas_caidas_avisadas[clave_pon] = ahora
            _avisar_caja_caida(red, caidas_los, total_relevante, porcentaje)
        else:
            # Puerto por debajo del umbral otra vez: se rearma para
            # poder avisar en una incidencia futura.
            _cajas_caidas_avisadas.pop(clave_pon, None)

    return afectadas


def _avisar_caja_caida(red, caidas_los, total_relevante, porcentaje):
    clientes_listado = []
    olt_pon = "N/D"
    nombre_cabecera = caidas_los[0][1].get("cabecera_nombre", "") if caidas_los else ""
    for cpe_id, _info in caidas_los:
        cliente = obtener_info_cliente_o_none(red, cpe_id)
        if cliente is None or not cliente.get("activa", True):
            continue
        if olt_pon == "N/D" and cliente.get("olt_pon", "N/D") != "N/D":
            olt_pon = cliente["olt_pon"]
        clientes_listado.append(
            "\\- {} \\(`{}`\\)".format(
                _escapar_markdown(cliente["cliente"]), _sn_formato_felix(cpe_id)
            )
        )

    extra = ""
    if len(clientes_listado) > CAJA_CAIDA_MAX_LISTADO:
        extra = "\n\\.\\.\\. y {} mas".format(len(clientes_listado) - CAJA_CAIDA_MAX_LISTADO)
        clientes_listado = clientes_listado[:CAJA_CAIDA_MAX_LISTADO]

    texto = (
        "\U0001F6A8 *CAJA CAIDA*\n"
        "━━━━━━━━━━━━━━\n"
        "\U0001F4E1 *Red:* {red}\n"
        "\U0001F5A5 *OLT:* {olt_pon}\n"
        "\U0001F4C9 *{caidas} de {total} ONTs del puerto en LOS \\({pct}%\\)*\n"
        "\U0001F550 *Hora:* {fecha}\n"
        "\n*Clientes afectados:*\n{listado}{extra}"
    ).format(
        red=_escapar_markdown(red["nombre"]),
        olt_pon=_escapar_markdown(_olt_con_nombre(olt_pon, nombre_cabecera)),
        caidas=len(caidas_los),
        total=total_relevante,
        pct=int(porcentaje * 100),
        fecha=_escapar_markdown(time.strftime("%Y-%m-%d %H:%M:%S")),
        listado="\n".join(clientes_listado) if clientes_listado else "\\(sin clientes con match en Felix\\)",
        extra=extra,
    )
    enviar_telegram(texto)
    log.info("[%s] CAJA CAIDA avisada: %d/%d ONTs en LOS (%.0f%%).",
             red["nombre"], len(caidas_los), total_relevante, porcentaje * 100)


def procesar_red(red, silencioso=False):
    try:
        estado_actual, puertos_pon = consultar_estado_olt(red)
    except SnmpError as e:
        enviar_alerta_admin(
            "Error SNMP consultando la OLT de la red '{}', se aborta esta "
            "red en este ciclo.\nDetalle: {}".format(red["nombre"], e),
            clave="snmp_error_{}".format(red["slug"]),
        )
        return
    except Exception as e:
        enviar_alerta_admin(
            "Error inesperado consultando/procesando SNMP para la red "
            "'{}', se aborta esta red en este ciclo.\nDetalle: {}\n"
            "```\n{}\n```".format(red["nombre"], e, traceback.format_exc()),
            clave="snmp_error_inesperado_{}".format(red["slug"]),
        )
        return

    _actualizar_cache_estado(red, estado_actual)

    # Precarga/refresco del mapa de clientes de Felix DENTRO del ciclo:
    # asi los comandos de Telegram (/estado, /los, /buscar) siempre
    # encuentran el cache caliente y responden al instante, sin tener
    # que descargar Felix en medio del comando.
    try:
        _obtener_mapa_clientes_felix(red, forzar=True)
    except Exception:
        log.warning("[%s] No se pudo precargar el mapa de Felix:\n%s",
                    red["nombre"], traceback.format_exc())

    # Deteccion de puerto PON caido (OID de estado del puerto, aviso
    # inmediato) + caja caida por porcentaje de ONTs. Las ONTs
    # cubiertas por cualquiera de los dos no generan avisos
    # individuales de LOS.
    afectadas_caja = set()
    if not silencioso:
        try:
            afectadas_caja |= _detectar_puertos_caidos(red, puertos_pon, estado_actual)
        except Exception:
            log.error("[%s] Error en deteccion de puertos caidos:\n%s",
                      red["nombre"], traceback.format_exc())
        try:
            afectadas_caja |= _detectar_cajas_caidas(red, estado_actual)
        except Exception:
            log.error("[%s] Error en deteccion de cajas caidas:\n%s",
                      red["nombre"], traceback.format_exc())
    else:
        # En la pasada inicial solo se guarda el estado base de los
        # puertos, sin avisar.
        for ip, _nombre, ifindex, online in puertos_pon:
            _puertos_online_antes[(ip, ifindex)] = online

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

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

        anterior = estado_anterior.get(cpe_id)
        online_antes = anterior.get("online") if isinstance(anterior, dict) else None
        # "notificado" persiste en disco: si este LOS ya se aviso por
        # Telegram, el flag sobrevive a reinicios del bot para que la
        # recuperacion se notifique aunque pasen dias y reinicios.
        notificado_antes = bool(anterior.get("notificado")) if isinstance(anterior, dict) else False
        motivo_antes = (anterior.get("motivo") or "") if isinstance(anterior, dict) else ""
        caida_ts_antes = anterior.get("caida_ts") if isinstance(anterior, dict) else None
        # La potencia solo se puede leer en vivo con la ONT online, asi
        # que mientras este offline se conserva la ultima lectura buena
        # (persistida en disco) en vez de perderla. Se usa para mostrar
        # "ultima potencia registrada" en el aviso de LOS.
        potencia_antes = anterior.get("potencia") if isinstance(anterior, dict) else None
        # Cliente Felix asignado al SN la ultima vez que se guardo estado
        # (nombre, no id: Felix no da un id estable aparte del nombre).
        # Se usa para detectar ONTs reutilizadas: mismo SN, cliente nuevo.
        cliente_antes = anterior.get("cliente") if isinstance(anterior, dict) else None
        cliente_info_actual = obtener_info_cliente_o_none(red, cpe_id)
        nombre_cliente_actual = cliente_info_actual.get("cliente") if cliente_info_actual else None
        nuevo_estado[cpe_id] = {
            "online": online_ahora,
            "indice": info.get("indice", "N/D"),
            "notificado": notificado_antes and not online_ahora,
            # Motivo y fecha de la caida guardados en disco: permiten
            # avisar la recuperacion de un LOS que nunca se aviso
            # (p.ej. cayo estando el bot parado).
            "motivo": info.get("motivo") or "",
            "caida_ts": info.get("caida_ts"),
            "potencia": info.get("potencia") or potencia_antes,
            # Si Felix no da match este ciclo (fallo puntual), se
            # conserva el ultimo cliente conocido en vez de perderlo.
            "cliente": nombre_cliente_actual if nombre_cliente_actual is not None else cliente_antes,
        }

        # --- ALTA REALIZADA (candidata) ---
        # SN totalmente NUEVO (no existia en la pasada anterior, ni
        # online ni offline) que aparece ONLINE y con cliente ACTIVO en
        # Felix (impagos/bajas excluidos): instalacion recien hecha.
        # Una ONT ya conocida que estaba offline y levanta NO es alta
        # (eso es recuperacion o silencio), EXCEPTO si el cliente
        # asignado en Felix cambio respecto al que tenia antes: ONT
        # reutilizada (el tecnico la recupera de un cliente de baja y
        # la instala en uno nuevo) -> es alta para el cliente nuevo,
        # no recuperacion del servicio del cliente anterior. Ver mas
        # abajo, rama "not online_antes", donde se comprueba esto.
        # Se acumulan y se envian al final del ciclo, con tope de
        # seguridad (muchas "altas" de golpe = perdida de estado, no
        # instalaciones reales).
        if (not primera_ejecucion and online_ahora and anterior is None):
            if cliente_info_actual is not None and cliente_info_actual.get("activa", True):
                altas_pendientes.append((cpe_id, info, cliente_info_actual))
            _ciclos_offline.pop((red["slug"], cpe_id), None)
            continue

        if primera_ejecucion or online_antes is None:
            continue

        clave = (red["slug"], cpe_id)

        if not online_ahora:
            if online_antes:
                _ciclos_offline[clave] = 1
            elif clave in _ciclos_offline:
                _ciclos_offline[clave] += 1

            ciclos = _ciclos_offline.get(clave, 0)
            if ciclos == LOS_CICLOS_CONFIRMACION:
                motivo = info.get("motivo") or ""
                if motivo.startswith("Desactivada") or motivo.startswith("Corte de luz"):
                    # Desactivacion administrativa (impago/baja) o corte
                    # de luz en casa del cliente: no es averia de fibra,
                    # no se avisa. Se saca del contador para que tampoco
                    # genere aviso de "Recuperado" al volver.
                    _ciclos_offline.pop(clave, None)
                    log.info("[%s] Caida ignorada (%s) sn=%s.",
                             red["nombre"], motivo, cpe_id)
                elif cpe_id in afectadas_caja:
                    # Cubierta por un aviso de CAJA CAIDA: no se manda
                    # aviso individual, pero se marca notificado para
                    # que su recuperacion si se avise.
                    nuevo_estado[cpe_id]["notificado"] = True
                    log.info("[%s] LOS cubierto por aviso de caja caida sn=%s.",
                             red["nombre"], cpe_id)
                else:
                    info_cliente = obtener_info_cliente(red, cpe_id)
                    if info_cliente.get("activa", True):
                        # La ONT ya esta offline: no se puede leer
                        # potencia en vivo. Se muestra la ultima
                        # registrada (persistida cuando aun estaba
                        # online), carried-over en nuevo_estado.
                        info_msg = dict(info)
                        info_msg["potencia"] = nuevo_estado[cpe_id].get("potencia")
                        texto = construir_mensaje(red, cpe_id, info_msg, info_cliente, "offline")
                        enviar_telegram(texto)
                        notificados += 1
                        nuevo_estado[cpe_id]["notificado"] = True
                        log.info("[%s] LOS notificado sn=%s (ciclo %d/%d).",
                                 red["nombre"], cpe_id, ciclos, LOS_CICLOS_CONFIRMACION)
                    else:
                        log.info("[%s] LOS ignorado (inactiva en Felix) sn=%s.",
                                 red["nombre"], cpe_id)
            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:
            if not online_antes:
                # ONT REASIGNADA: mismo SN, pero Felix ahora lo tiene
                # asociado a un cliente ACTIVO distinto del que tenia
                # registrado la ultima vez (y antes SI tenia uno
                # registrado, para no disparar esto en el primer ciclo
                # tras desplegar este fix, cuando "cliente_antes" aun
                # no existe para ningun SN). Es una instalacion nueva
                # para ese cliente, no la recuperacion del servicio del
                # cliente anterior -> Alta, no Recuperado.
                cliente_reasignado = (
                    cliente_antes is not None
                    and nombre_cliente_actual is not None
                    and nombre_cliente_actual != cliente_antes
                    and cliente_info_actual.get("activa", True)
                )
                if cliente_reasignado:
                    info_msg = dict(info)
                    try:
                        potencia_relectura = _leer_potencia_rx(red, info)
                    except Exception:
                        potencia_relectura = None
                    info_msg["potencia"] = potencia_relectura or info.get("potencia")
                    texto = construir_mensaje(red, cpe_id, info_msg, cliente_info_actual, "alta")
                    enviar_telegram(texto)
                    notificados += 1
                    log.info(
                        "[%s] Alta realizada notificada (ONT reasignada, cliente "
                        "antes='%s' ahora='%s') sn=%s.",
                        red["nombre"], cliente_antes, nombre_cliente_actual, cpe_id,
                    )
                    _ciclos_offline.pop(clave, None)
                    continue

                ciclos_fue = _ciclos_offline.pop(clave, 0)
                # Se avisa la recuperacion si:
                #  a) el LOS se notifico en su momento (en esta sesion o
                #     en una anterior, via flag "notificado" en disco), o
                #  b) estaba caida por LOS reciente (< LOS_MAX_HORAS)
                #     aunque nunca se avisara (p.ej. cayo con el bot
                #     parado). Impagos, cortes de luz y fantasmas de
                #     meses siguen recuperandose en silencio.
                era_los_reciente = (
                    motivo_antes.startswith("Fallo LOS")
                    and (caida_ts_antes is None
                         or (time.time() - caida_ts_antes) <= LOS_MAX_HORAS * 3600)
                )
                if (ciclos_fue >= LOS_CICLOS_CONFIRMACION or notificado_antes
                        or era_los_reciente):
                    info_cliente = obtener_info_cliente(red, cpe_id)
                    info_msg = dict(info)
                    # El walk ya trae potencia valida (ONT esta online
                    # en este ciclo). Se intenta una relectura en vivo
                    # por si acaba de reengancharse y la del walk quedo
                    # desfasada, pero si la relectura falla (DDM aun no
                    # estable) NO se pisa el valor bueno del walk con
                    # None: se conserva ese como fallback.
                    try:
                        potencia_relectura = _leer_potencia_rx(red, info)
                    except Exception:
                        potencia_relectura = None
                    info_msg["potencia"] = potencia_relectura or info.get("potencia")
                    texto = construir_mensaje(red, cpe_id, info_msg, info_cliente, "online")
                    enviar_telegram(texto)
                    notificados += 1
                    log.info("[%s] Recuperacion notificada sn=%s.", red["nombre"], cpe_id)
                else:
                    log.info("[%s] Recuperacion silenciosa sn=%s (LOS no confirmado, %d/%d ciclos).",
                             red["nombre"], cpe_id, ciclos_fue, LOS_CICLOS_CONFIRMACION)
            _ciclos_offline.pop(clave, None)

    # Envio de altas con tope de seguridad: mas de unas pocas en el
    # mismo ciclo no son instalaciones reales sino perdida de estado
    # (fichero borrado, cabecera que reaparece...). En ese caso se
    # suprime el envio masivo y se avisa al admin una sola vez.
    ALTAS_MAX_POR_CICLO = 5
    if len(altas_pendientes) > ALTAS_MAX_POR_CICLO:
        enviar_alerta_admin(
            "Se han detectado {} 'altas' de golpe en la red '{}' — "
            "imposible que sean instalaciones reales (posible perdida "
            "de estado). Avisos de alta suprimidos este ciclo.".format(
                len(altas_pendientes), red["nombre"]),
            clave="altas_masivas_{}".format(red["slug"]),
        )
    else:
        for cpe_id_alta, info_alta, cliente_alta in altas_pendientes:
            info_msg = dict(info_alta)
            # Mismo fallback que en la recuperacion: si la relectura en
            # vivo falla, se conserva la potencia que ya trajo el walk
            # (la ONT esta online en este ciclo, asi que suele venir
            # rellena) en vez de perderla y dejar el mensaje sin dato.
            try:
                potencia_relectura = _leer_potencia_rx(red, info_alta)
            except Exception:
                potencia_relectura = None
            info_msg["potencia"] = potencia_relectura or info_alta.get("potencia")
            texto = construir_mensaje(red, cpe_id_alta, info_msg, cliente_alta, "alta")
            enviar_telegram(texto)
            notificados += 1
            log.info("[%s] Alta realizada notificada sn=%s.", red["nombre"], cpe_id_alta)

    # Las ONTs que NO han aparecido en el walk de este ciclo (p.ej.
    # porque su cabecera fallo el SNMP esta pasada) conservan su estado
    # anterior en vez de borrarse. Sin esto, al volver la cabecera
    # todas sus ONTs parecerian "nuevas" (nunca vistas online) y
    # dispararian una cascada de "Alta realizada" falsas.
    for cpe_id_ant, registro_ant in estado_anterior.items():
        if cpe_id_ant not in nuevo_estado and isinstance(registro_ant, dict):
            nuevo_estado[cpe_id_ant] = registro_ant

    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(silencioso=False):
    """Procesa todas las redes EN PARALELO (una por hilo): el ciclo
    dura lo que la red mas lenta, no la suma de las tres."""
    hilos = []
    for red in REDES:
        hilo = threading.Thread(target=_procesar_red_protegido, args=(red, silencioso))
        hilo.start()
        hilos.append(hilo)
    for hilo in hilos:
        hilo.join()


def _procesar_red_protegido(red, silencioso):
    try:
        procesar_red(red, silencioso=silencioso)
    except Exception:
        log.error("Error no controlado procesando la red '%s':\n%s",
                  red["nombre"], traceback.format_exc())


# ---------------------------------------------------------------------------
# MODO COMANDOS (long polling)
# ---------------------------------------------------------------------------

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():
    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():
    log.info("Iniciando monitorizacion en bucle, intervalo=%ss", SLEEP_SECONDS)
    log.info("Pasada inicial silenciosa: actualizando estado base sin notificar...")
    procesar_una_vez(silencioso=True)
    log.info("Pasada inicial completada. Iniciando monitorizacion normal.")
    while True:
        inicio_ciclo = time.time()
        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",
            )
        # Se descuenta lo que ha tardado el ciclo (los walks SNMP de
        # todas las cabeceras pueden llevar varios minutos) para que la
        # cadencia real entre pasadas sea ~SLEEP_SECONDS, no
        # SLEEP_SECONDS + duracion del ciclo. Minimo 30s de respiro.
        duracion = time.time() - inicio_ciclo
        espera = max(30, SLEEP_SECONDS - duracion)
        log.info("Ciclo de %.0fs, durmiendo %.0fs hasta el siguiente.", duracion, espera)
        time.sleep(espera)


if __name__ == "__main__":
    # Uso normal (sin argumentos): monitorizacion + comandos a la vez,
    # en primer plano. Para dejarlo corriendo tras cerrar la SSH:
    #   nohup python3 telegram_los_bot_snmp.py > bot_felix.log 2>&1 &
    #   ps aux | grep telegram_los_bot_snmp   (ver PID)
    #   kill <PID>                            (parar)
    #
    # Modos opcionales:
    #   python3 telegram_los_bot_snmp.py una-vez     -> una pasada y termina
    #   python3 telegram_los_bot_snmp.py comandos    -> solo escucha comandos
    _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
    for _red in REDES:
        # Red sin cabeceras SNMP dadas de alta todavia: no bloquea el
        # arranque (procesar_red la omite sola), solo se avisa.
        if not _red.get("cabeceras"):
            log.warning("Red '%s' sin cabeceras SNMP configuradas, se omitira "
                        "hasta que se anadan a REDES.", _red["nombre"])
    if not _config_ok:
        log.error("El bot no puede arrancar con la configuracion incompleta. "
                  "Crea/edita el fichero .env junto al script con, p.ej.:\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":
        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)

    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)

    hilo_comandos = threading.Thread(target=escuchar_comandos, daemon=True)
    hilo_comandos.start()
    bucle_monitorizacion()
