# -*- coding: utf-8 -*- """telegram.py — Fuente Telegram del monitor (RF-1, RF-6). Recoge los mensajes del **dia en curso** de las 5 fuentes configuradas usando Telethon con una sesion DEDICADA (`~/.telegram-agent/sessions/monitor`), distinta de la del watcher (que mantiene la suya bloqueada). La primera vez hay que autorizar esa sesion manualmente: python scripts/monitor_estado.py login-telegram El nucleo (`normalizar_mensajes`) es puro y testeable sin conexion. """ from __future__ import annotations import asyncio import os import subprocess import sys import tempfile from datetime import date, datetime from pathlib import Path from .. import config as config_mod # ------------------------------------------------------------------ credenciales def _credenciales() -> tuple[int, str]: """Devuelve (api_id, api_hash) desde el entorno de usuario (igual que el watcher).""" api_id = int(os.environ.get("TELEGRAM_API_ID", "0") or 0) api_hash = os.environ.get("TELEGRAM_API_HASH", "") if not api_id or not api_hash: try: r1 = subprocess.run( ["powershell", "-Command", '[Environment]::GetEnvironmentVariable("TELEGRAM_API_ID","User")'], capture_output=True, text=True, ) r2 = subprocess.run( ["powershell", "-Command", '[Environment]::GetEnvironmentVariable("TELEGRAM_API_HASH","User")'], capture_output=True, text=True, ) api_id = int(r1.stdout.strip()) if r1.stdout.strip() else 0 api_hash = r2.stdout.strip() except Exception: pass return api_id, api_hash def _sesion_path(cfg: dict) -> str: """Ruta de la sesion Telethon (sin extension) expandiendo `~`.""" ruta = cfg.get("telegram", {}).get("sesion", "~/.telegram-agent/sessions/monitor") return str(Path(ruta).expanduser()) def construir_cliente(cfg: dict): """Crea un TelegramClient con la sesion dedicada del monitor.""" from telethon import TelegramClient api_id, api_hash = _credenciales() if not api_id or not api_hash: raise RuntimeError("Faltan TELEGRAM_API_ID / TELEGRAM_API_HASH en el entorno del usuario") return TelegramClient(_sesion_path(cfg), api_id, api_hash) # ------------------------------------------------------------------ nucleo puro def _dia_local(fecha: datetime) -> date: """Fecha local de un datetime (Telethon entrega fechas con zona).""" return fecha.astimezone().date() if fecha.tzinfo else fecha.date() def normalizar_mensajes(raw_mensajes, dia: date, nombre_fuente: str) -> list[dict]: """Convierte mensajes crudos en registros del dia, ordenados cronologicamente. Cada registro: {id, fecha (ISO local), quien, texto, fuente}. Solo incluye mensajes con texto y del dia indicado. No hace llamadas de red. """ salida: list[dict] = [] for m in raw_mensajes: fecha = getattr(m, "date", None) texto = getattr(m, "message", None) if not fecha or not texto: continue if _dia_local(fecha) != dia: continue quien = getattr(m, "_sender_name", None) or str(getattr(m, "sender_id", "") or "?") salida.append({ "id": getattr(m, "id", None), "fecha": fecha.astimezone().isoformat() if fecha.tzinfo else fecha.isoformat(), "quien": quien, "texto": texto, "fuente": nombre_fuente, }) salida.sort(key=lambda r: r["fecha"]) return salida def render_mensaje(m: dict) -> str: """Render Markdown de un mensaje para el fichero del dia.""" try: hora = datetime.fromisoformat(m["fecha"]).strftime("%H:%M") except Exception: hora = "--:--" texto = (m.get("texto") or "").replace("\r", "").rstrip() lineas = texto.split("\n") cuerpo = "\n".join(f" {ln}" for ln in lineas[1:]) cabecera = f"- **[{hora}] {m.get('quien', '?')}:** {lineas[0] if lineas else ''}" return f"{cabecera}\n{cuerpo}" if cuerpo else cabecera # ------------------------------------------------------------------ red (Telethon) async def _nombre_sender(mensaje) -> str: """Nombre legible del emisor de un mensaje.""" try: s = await mensaje.get_sender() if s is not None: return getattr(s, "first_name", None) or getattr(s, "title", None) or getattr(s, "username", None) or str(mensaje.sender_id) except Exception: pass return str(getattr(mensaje, "sender_id", "") or "?") async def recoger_fuente(client, nombre: str, chat_id, dia: date) -> list[dict]: """Recoge los mensajes de hoy de una fuente concreta (una llamada de red).""" entidad = await client.get_entity(chat_id) crudos = [] async for m in client.iter_messages(entidad, limit=None): fecha = getattr(m, "date", None) if not fecha: continue if _dia_local(fecha) < dia: break # ya salimos del dia en curso (vienen en orden descendente) m._sender_name = await _nombre_sender(m) crudos.append(m) return normalizar_mensajes(crudos, dia, nombre) async def _recoger_todas_async(cfg: dict, dia: date, client_factory=None) -> tuple[dict, list]: """Conecta una vez y recoge todas las fuentes, aislando errores por fuente (RF-6).""" client = (client_factory or construir_cliente)(cfg) await client.connect() try: if not await client.is_user_authorized(): raise RuntimeError("Sesion de Telegram no autorizada: ejecuta 'login-telegram'") # Cachea las entidades de los dialogos para poder resolver los IDs de las # fuentes por numero (Telethon necesita el access_hash de cada peer). await client.get_dialogs() resultado: dict[str, list] = {} errores: list[tuple[str, str]] = [] for f in config_mod.fuentes_telegram(cfg): nombre = f.get("nombre") try: resultado[nombre] = await recoger_fuente(client, nombre, f.get("id"), dia) except Exception as e: # una fuente caida no aborta el ciclo errores.append((nombre, str(e))) return resultado, errores finally: await client.disconnect() def recoger_todas(cfg: dict | None = None, dia: date | None = None, client_factory=None) -> tuple[dict, list]: """Version sincrona: {fuente: [mensajes de hoy]} y lista de errores por fuente.""" cfg = cfg or config_mod.cargar() dia = dia or date.today() return asyncio.run(_recoger_todas_async(cfg, dia, client_factory)) def login(cfg: dict | None = None, force_sms: bool = False) -> None: """Autoriza la sesion dedicada de forma interactiva (una sola vez). `force_sms=True` obliga a que el codigo llegue por SMS (util si no aparece en la app de Telegram). """ import getpass from telethon.errors import SessionPasswordNeededError cfg = cfg or config_mod.cargar() client = construir_cliente(cfg) async def _run(): await client.connect() if await client.is_user_authorized(): me = await client.get_me() print(f"[monitor_estado] Sesion ya autorizada: {me.first_name} (@{me.username or 'sin usuario'})") await client.disconnect() return telefono = cfg.get("telegram", {}).get("telefono") sent = await client.send_code_request(telefono, force_sms=force_sms) print(f"[login] Codigo solicitado (tipo: {type(sent.type).__name__}, force_sms={force_sms}).") print("[login] Introduce el codigo (llega a la app de Telegram o por SMS).") codigo = input("Codigo: ").strip() try: await client.sign_in(phone=telefono, code=codigo, phone_code_hash=sent.phone_code_hash) except SessionPasswordNeededError: password = getpass.getpass("Contrasena 2FA: ") await client.sign_in(password=password) me = await client.get_me() print(f"[monitor_estado] Sesion autorizada: {me.first_name} (@{me.username or 'sin usuario'})") await client.disconnect() asyncio.run(_run()) def login_qr(cfg: dict | None = None, timeout: int = 180) -> None: """Autoriza la sesion dedicada via codigo QR (escaneo desde la app de Telegram). Abre una imagen QR en pantalla; se escanea desde Telegram > Ajustes > Dispositivos > Vincular dispositivo. """ cfg = cfg or config_mod.cargar() client = construir_cliente(cfg) async def _run(): await client.connect() if await client.is_user_authorized(): me = await client.get_me() print(f"[monitor_estado] Sesion ya autorizada: {me.first_name} (@{me.username or 'sin usuario'})") await client.disconnect() return qr = await client.qr_login() try: import qrcode img = qrcode.make(qr.url) ruta = Path(tempfile.gettempdir()) / "monitor_tg_qr.png" img.save(ruta) print(f"[login-qr] QR generado en: {ruta}") print("[login-qr] Telegram > Ajustes > Dispositivos > Vincular dispositivo.") try: os.startfile(str(ruta)) # abre la imagen en el visor por defecto except Exception: pass except Exception as e: print(f"[login-qr] No se pudo generar la imagen ({e}). Usa la URL:\n{qr.url}") try: await qr.wait(timeout=timeout) except Exception as e: print(f"[login-qr] QR expirado o cancelado: {e}") await client.disconnect() return me = await client.get_me() print(f"[monitor_estado] Sesion autorizada: {me.first_name} (@{me.username or 'sin usuario'})") await client.disconnect() asyncio.run(_run()) if __name__ == "__main__": if len(sys.argv) > 1 and sys.argv[1] == "login": login() elif len(sys.argv) > 1 and sys.argv[1] == "login-qr": login_qr() else: mensajes, errs = recoger_todas() for nombre, msgs in mensajes.items(): print(f"{nombre}: {len(msgs)} mensajes hoy") for nombre, err in errs: print(f"ERROR {nombre}: {err}", file=sys.stderr)