264 lines
10 KiB
Python
264 lines
10 KiB
Python
# -*- 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)
|